From 61d2b08a05fdb4062734e843c9969e2c66fdd09a Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 00:24:24 +0800 Subject: [PATCH 1/4] feat(chat): bind attached sessions to GoalRef Signed-off-by: duanjialing.777 --- ...nstance-identity-and-orphan-recovery-v0.md | 14 + ...e-identity-and-orphan-recovery-v0.zh-CN.md | 12 + loopx/attached_session.py | 589 +++++++++++++++++- loopx/attached_session_api.py | 12 +- loopx/chat_runtime.py | 70 ++- loopx/chat_store.py | 85 ++- loopx/cli_commands/worker_bridge.py | 7 +- .../control_plane/effect_runtime_handlers.ts | 2 + .../goals/chat_session_lifecycle.ts | 278 +++++++++ .../goal_instance_binding_inventory_v1.json | 16 +- .../project_registry_io_manifest_v1.json | 48 +- .../test_goal_instance_binding_inventory.py | 4 + .../test_source_session_registry_denial.py | 7 + .../chat_session_lifecycle.test.ts | 167 +++++ .../effect_runtime_handlers.test.ts | 13 + tests/test_attached_session_goal_instance.py | 570 +++++++++++++++++ tsconfig.control-plane.json | 2 + 17 files changed, 1841 insertions(+), 55 deletions(-) create mode 100644 loopx/control_plane/goals/chat_session_lifecycle.ts create mode 100644 tests/control_plane_ts/chat_session_lifecycle.test.ts create mode 100644 tests/test_attached_session_goal_instance.py diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md index 12b9af5f4e..de06ec43ed 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.md @@ -769,6 +769,20 @@ promotion retain their own acceptance. No new paid cohort or soak is authorized. qualify the remaining effect owners before existing-project activation or global routing can open. +### 2026-09-26: M3 attached-host Chat candidate + +- **Baseline:** `9849366c6`. +- **Proposed:** Bind attached Chat sessions to the current exact GoalRef and + recheck it for enqueue, resume, and new claims under the M2 lifetime guard. + A delayed result may complete only when its persisted claim admission still + matches the historical session. +- **Compatibility:** Non-source profiles keep the existing writers, lookup, + broker payloads, and serialized bytes. Clients do not submit + `goal_instance_id`. +- **Remaining hold:** This qualifies only `attached_host_chat_session`. + Managed provider startup and downstream host effects remain blocked by the + separate `first_party_host_runtime` row and the overall M3 activation hold. + ### 2026-09-27: first-party Host runtime partial enforcement - **Baseline:** `fd96e5e2574272262b9ea604a96581a0d20e94d1` diff --git a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md index 68c47ed23c..a5c831e63a 100644 --- a/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md +++ b/docs/architecture/rfcs/goal-instance-identity-and-orphan-recovery-v0.zh-CN.md @@ -698,6 +698,18 @@ service adoption、D1–D3 provider promotion 保留各自验收。不授权付 - **剩余 hold:** 所有结果均为 `execution_authority: false`。M3 必须先完成其余 effect owner 资格化,才能开放既有项目 activation 或 global routing。 +### 2026-09-26:M3 attached-host Chat 候选 + +- **基线:** `9849366c6`。 +- **候选实现:** attached Chat Session 绑定当前精确 GoalRef;enqueue、resume + 与新 claim 在 M2 lifetime guard 内重新校验。迟到结果只有在持久化 claim + admission 仍与历史 Session 一致时才能完成。 +- **兼容性:** 非 source profile 保留既有 writer、lookup、broker payload 与 + 序列化字节;客户端不提交 `goal_instance_id`。 +- **剩余 hold:** 本切片只资格化 `attached_host_chat_session`。Managed + provider 启动和下游 host effect 仍受独立 `first_party_host_runtime` 行与 + M3 总 activation hold 阻断。 + ### 2026-09-27:第一方 Host runtime 部分 enforcement - **基线:** `fd96e5e2574272262b9ea604a96581a0d20e94d1` diff --git a/loopx/attached_session.py b/loopx/attached_session.py index cc77c6f776..6c1835776b 100644 --- a/loopx/attached_session.py +++ b/loopx/attached_session.py @@ -5,14 +5,24 @@ import hashlib import math import time -from collections.abc import Mapping +from collections.abc import Iterator, Mapping +from contextlib import contextmanager from pathlib import Path from typing import Any from .chat import normalize_agent_response from .chat_store import CHAT_SESSION_MODE_ATTACHED, ChatSessionStore +from .control_plane.effect_runtime import effect_runtime_result +from .control_plane.goals.source_session_registry_state import ( + current_goal_ref, + guard_path, +) +from .control_plane.projects.registry_codec import ( + SOURCE_SESSION_PROFILE_ID, + load_project_registry, +) from .control_plane.todos.contract import normalize_todo_claimed_by -from .file_lock import exclusive_file_lock +from .file_lock import exclusive_cross_runtime_file_lock, exclusive_file_lock from .registry import find_registry_goal from .thread_agent_binding import resolve_thread_agent_binding @@ -69,10 +79,248 @@ def _require_bound_host( ) +def _session_fact(session: Mapping[str, Any]) -> dict[str, Any]: + return { + "session_id": str(session.get("session_id") or ""), + "goal_id": str(session.get("goal_id") or ""), + "goal_instance_id": session.get("goal_instance_id"), + "updated_at": str(session.get("updated_at") or ""), + } + + +def _turn_fact( + session: Mapping[str, Any], + turn: Mapping[str, Any] | None, +) -> dict[str, Any] | None: + if turn is None: + return None + return { + "goal_id": str(session.get("goal_id") or ""), + "goal_instance_id": turn.get("goal_instance_id"), + "admitted_goal_instance_id": turn.get("admitted_goal_instance_id"), + } + + +def _lifecycle_decision( + *, + operation: str, + registry: Mapping[str, Any], + current_ref: Mapping[str, str], + session: Mapping[str, Any] | None = None, + turn: Mapping[str, Any] | None = None, + candidates: list[dict[str, Any]] | None = None, +) -> dict[str, Any]: + facts: dict[str, Any] = { + "operation": operation, + "profile_id": registry.get("profile_id"), + "current_goal_ref": dict(current_ref), + } + if operation == "select": + facts["candidates"] = [ + _session_fact(candidate) for candidate in candidates or [] + ] + else: + if session is None: + raise RuntimeError("Chat lifecycle admission requires a Session") + facts["session"] = _session_fact(session) + facts["turn"] = _turn_fact(session, turn) + decision = effect_runtime_result( + "goal.chat_session.lifecycle.decide", + facts, + ) + if not isinstance(decision, dict): + raise RuntimeError("Chat session lifecycle decision must be an object") + if decision.get("kind") == "reject": + raise ValueError(f"attached Chat session rejected: {decision.get('code')}") + return decision + + +def _strict_registry( + registry_path: Path | None, +) -> dict[str, Any] | None: + if registry_path is None or not registry_path.exists(): + return None + registry = load_project_registry(registry_path) + return registry if registry.get("profile_id") == SOURCE_SESSION_PROFILE_ID else None + + +def load_attached_session_registry(registry_path: Path) -> dict[str, Any]: + """Load the source-aware registry for this qualified owner only.""" + + return load_project_registry(registry_path) + + +def _select_current_session( + *, + store: ChatSessionStore, + registry: Mapping[str, Any], + current_ref: Mapping[str, str], + goal_id: str, + agent_id: str, + channel_id: str, +) -> dict[str, Any] | None: + candidates = store.resumable_session_candidates( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + decision = _lifecycle_decision( + operation="select", + registry=registry, + current_ref=current_ref, + candidates=candidates, + ) + if decision.get("kind") == "create": + return None + if decision.get("kind") != "reuse": + raise RuntimeError("Chat session selection decision is unsupported") + session_id = str(decision.get("session_id") or "") + selected = next( + ( + candidate + for candidate in candidates + if candidate.get("session_id") == session_id + ), + None, + ) + if selected is None: + raise RuntimeError("Chat session selection omitted its selected Session") + return selected + + +def _bind_source_session( + *, + store: ChatSessionStore, + registry_path: Path, + goal_id: str, + agent_id: str, + host_surface: str, + host_session_id: str, + executor_endpoint_id: str, + channel_id: str | None, + execute: bool, +) -> dict[str, Any]: + guard = guard_path(registry_path, goal_id) + with exclusive_cross_runtime_file_lock( + guard, + operation="source_session_goal_lifetime", + ): + registry = load_project_registry(registry_path) + current_ref, goal = current_goal_ref(registry, goal_id=goal_id) + normalized_agent = _registered_agent(goal, agent_id) + _require_bound_host( + goal=goal, + agent_id=normalized_agent, + host_surface=host_surface, + host_session_id=host_session_id, + ) + selected_channel = channel_id or f"goal.{goal_id}" + lock_path = _binding_lock_path( + store, + goal_id=goal_id, + agent_id=normalized_agent, + channel_id=selected_channel, + ) + with exclusive_file_lock( + lock_path, + agent_id="loopx-chat", + operation="bind_attached_agent_session", + ): + latest = _select_current_session( + store=store, + registry=registry, + current_ref=current_ref, + goal_id=goal_id, + agent_id=normalized_agent, + channel_id=selected_channel, + ) + if latest is not None: + exact_match = ( + latest.get("session_mode") == CHAT_SESSION_MODE_ATTACHED + and latest.get("host_surface") == host_surface + and latest.get("upstream_thread_id") == host_session_id + and latest.get("executor_endpoint_id") == executor_endpoint_id + ) + if not exact_match: + raise ValueError( + "an active working Session already exists for this " + "Goal, Agent, and channel" + ) + return _bind_result( + store=store, + session=latest, + execute=execute, + changed=False, + created=False, + ) + if not execute: + return { + "ok": True, + "schema_version": ATTACHED_SESSION_BROKER_SCHEMA_VERSION, + "action": "bind", + "execute": False, + "changed": True, + "created": False, + "binding": { + "goal_id": goal_id, + "agent_id": normalized_agent, + "executor_endpoint_id": executor_endpoint_id, + "host_surface": host_surface, + "channel_id": selected_channel, + "session_mode": CHAT_SESSION_MODE_ATTACHED, + }, + } + session = store.create_session( + goal_id=goal_id, + goal_instance_id=current_ref["goal_instance_id"], + agent_id=normalized_agent, + executor_endpoint_id=executor_endpoint_id, + adapter_kind=ATTACHED_SESSION_ADAPTER_KIND, + upstream_thread_id=host_session_id, + upstream_mode=ATTACHED_SESSION_UPSTREAM_MODE, + channel_id=selected_channel, + session_mode=CHAT_SESSION_MODE_ATTACHED, + host_surface=host_surface, + attached_capabilities={ + "live_steering": False, + "session_queue": True, + "claim_wait": True, + "reply_readback": True, + }, + ) + return _bind_result( + store=store, + session=session, + execute=True, + changed=True, + created=True, + ) + + +def _bind_result( + *, + store: ChatSessionStore, + session: Mapping[str, Any], + execute: bool, + changed: bool, + created: bool, +) -> dict[str, Any]: + return { + "ok": True, + "schema_version": ATTACHED_SESSION_BROKER_SCHEMA_VERSION, + "action": "bind", + "execute": execute, + "changed": changed, + "created": created, + "session": store.public_session(dict(session)), + } + + def bind_attached_agent_session( *, store: ChatSessionStore, registry: dict[str, Any], + registry_path: Path | None = None, goal_id: str, agent_id: str, host_surface: str, @@ -83,6 +331,20 @@ def bind_attached_agent_session( ) -> dict[str, Any]: """Bind one existing host session without starting or resuming an adapter.""" + if registry.get("profile_id") == SOURCE_SESSION_PROFILE_ID: + if registry_path is None: + raise ValueError("source-session attached binding requires registry_path") + return _bind_source_session( + store=store, + registry_path=registry_path, + goal_id=goal_id, + agent_id=agent_id, + host_surface=host_surface, + host_session_id=host_session_id, + executor_endpoint_id=executor_endpoint_id, + channel_id=channel_id, + execute=execute, + ) goal = find_registry_goal(registry, goal_id) if goal is None: raise ValueError(f"goal_id not found in registry: {goal_id}") @@ -196,23 +458,284 @@ def _require_attached_host( return session -def claim_attached_agent_turn( +@contextmanager +def _attached_session_lifetime( *, store: ChatSessionStore, + registry_path: Path | None, + session_id: str, + allow_closed: bool = False, +) -> Iterator[ + tuple[dict[str, Any], dict[str, Any] | None, dict[str, str] | None] +]: + """Load authoritative Goal context before choosing strict or legacy behavior.""" + + session = store.load_session(session_id) + if session is None or (session.get("status") == "closed" and not allow_closed): + raise KeyError("attached Agent session was not found") + if session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED: + raise ValueError("the selected Session is not an attached host session") + exact_session = session.get("goal_instance_id") is not None + if registry_path is None: + if exact_session: + raise ValueError("source-session attached operation requires registry_path") + yield session, None, None + return + if not registry_path.exists(): + if exact_session: + raise FileNotFoundError(registry_path) + yield session, None, None + return + + goal_id = str(session.get("goal_id") or "") + with exclusive_cross_runtime_file_lock( + guard_path(registry_path, goal_id), + operation="source_session_goal_lifetime", + ): + registry = load_project_registry(registry_path) + session = store.load_session(session_id) + if session is None or ( + session.get("status") == "closed" and not allow_closed + ): + raise KeyError("attached Agent session was not found") + if session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED: + raise ValueError("the selected Session is not an attached host session") + if registry.get("profile_id") != SOURCE_SESSION_PROFILE_ID: + yield session, None, None + return + current_ref, _goal = current_goal_ref(registry, goal_id=goal_id) + yield session, registry, current_ref + + +def select_current_attached_session( + *, + store: ChatSessionStore, + registry_path: Path | None, + goal_id: str, + agent_id: str, + channel_id: str, +) -> tuple[dict[str, Any] | None, bool]: + """Select an attached Session and report whether strict identity is active.""" + + if registry_path is not None and not registry_path.exists(): + candidates = store.session_candidates( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + if any(candidate.get("goal_instance_id") is not None for candidate in candidates): + raise FileNotFoundError(registry_path) + registry = _strict_registry(registry_path) + if registry is None: + return ( + store.latest_session( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ), + False, + ) + assert registry_path is not None + guard = guard_path(registry_path, goal_id) + with exclusive_cross_runtime_file_lock( + guard, + operation="source_session_goal_lifetime", + ): + registry = load_project_registry(registry_path) + current_ref, _goal = current_goal_ref(registry, goal_id=goal_id) + session = _select_current_session( + store=store, + registry=registry, + current_ref=current_ref, + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + if ( + session is not None + and session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED + ): + raise ValueError( + "source-session managed Chat is not qualified for execution" + ) + return session, True + + +@contextmanager +def _current_attached_session_guard( + *, + store: ChatSessionStore, + registry_path: Path | None, + session_id: str, +) -> Iterator[dict[str, Any]]: + """Hold Goal lifetime authority while an exact attached Session is used.""" + + with _attached_session_lifetime( + store=store, + registry_path=registry_path, + session_id=session_id, + ) as (session, registry, current_ref): + if registry is None or current_ref is None: + yield session + return + _lifecycle_decision( + operation="admit", + registry=registry, + current_ref=current_ref, + session=session, + ) + yield session + + +def require_current_attached_session( + *, + store: ChatSessionStore, + registry_path: Path | None, + session_id: str, +) -> dict[str, Any]: + """Reject stale exact sessions before granting new authority.""" + + with _current_attached_session_guard( + store=store, + registry_path=registry_path, + session_id=session_id, + ) as session: + return session + + +def resume_attached_agent_session( + *, + store: ChatSessionStore, + registry_path: Path | None, + session_id: str, +) -> dict[str, Any]: + """Resume an attached Session while its Goal lifetime remains stable.""" + + with _current_attached_session_guard( + store=store, + registry_path=registry_path, + session_id=session_id, + ) as session: + if session.get("active_turn_id"): + return session + return store.update_session( + session_id, + status="ready", + active_turn_id=None, + last_error_code=None, + ) + + +def enqueue_attached_agent_turn( + *, + store: ChatSessionStore, + registry_path: Path | None, + session_id: str, + client_turn_id: str, + message: str, + origin: str, +) -> tuple[dict[str, Any], bool]: + """Enqueue work while the exact attached Session is current.""" + + with _attached_session_lifetime( + store=store, + registry_path=registry_path, + session_id=session_id, + ) as (session, registry, current_ref): + if registry is None or current_ref is None: + return store.create_queued_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + origin=origin, + ) + decision = _lifecycle_decision( + operation="admit", + registry=registry, + current_ref=current_ref, + session=session, + ) + goal_ref = decision.get("goal_ref") + if not isinstance(goal_ref, dict): + raise RuntimeError("Chat enqueue admission omitted its GoalRef") + return store.create_queued_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + goal_instance_id=str(goal_ref["goal_instance_id"]), + origin=origin, + ) + + +def _claim_attached_turn_once( + *, + store: ChatSessionStore, + registry_path: Path | None, session_id: str, host_surface: str, host_session_id: str, claim_id: str, - wait_seconds: float = 0.0, -) -> dict[str, Any]: - """Claim or bounded-wait for the oldest queued message for the exact host.""" - - _require_attached_host( +) -> dict[str, Any] | None: + session = _require_attached_host( store=store, session_id=session_id, host_surface=host_surface, host_session_id=host_session_id, ) + with _attached_session_lifetime( + store=store, + registry_path=registry_path, + session_id=session_id, + ) as (_session, registry, current_ref): + session = _require_attached_host( + store=store, + session_id=session_id, + host_surface=host_surface, + host_session_id=host_session_id, + ) + if registry is None or current_ref is None: + return store.claim_next_queued_turn( + session_id, + host_claim_id=claim_id, + ) + active_turn_id = str(session.get("active_turn_id") or "") + active_turn = ( + store.load_turn(session_id, active_turn_id) if active_turn_id else None + ) + replay = bool( + active_turn + and active_turn.get("host_claim_id") == claim_id + and active_turn.get("status") in {"starting", "running"} + ) + decision = _lifecycle_decision( + operation="replay_claim" if replay else "claim", + registry=registry, + current_ref=current_ref, + session=session, + turn=active_turn if replay else None, + ) + goal_ref = decision.get("goal_ref") + if not isinstance(goal_ref, dict): + raise RuntimeError("Chat claim admission omitted its GoalRef") + return store.claim_next_queued_turn( + session_id, + host_claim_id=claim_id, + admitted_goal_instance_id=str(goal_ref["goal_instance_id"]), + ) + + +def claim_attached_agent_turn( + *, + store: ChatSessionStore, + registry_path: Path | None = None, + session_id: str, + host_surface: str, + host_session_id: str, + claim_id: str, + wait_seconds: float = 0.0, +) -> dict[str, Any]: + """Claim or bounded-wait for the oldest queued message for the exact host.""" + normalized_wait = float(wait_seconds) if ( not math.isfinite(normalized_wait) @@ -226,7 +749,14 @@ def claim_attached_agent_turn( deadline = time.monotonic() + normalized_wait turn = None while turn is None: - turn = store.claim_next_queued_turn(session_id, host_claim_id=claim_id) + turn = _claim_attached_turn_once( + store=store, + registry_path=registry_path, + session_id=session_id, + host_surface=host_surface, + host_session_id=host_session_id, + claim_id=claim_id, + ) if turn is not None or normalized_wait == 0: break remaining = deadline - time.monotonic() @@ -259,6 +789,7 @@ def claim_attached_agent_turn( def complete_attached_agent_turn( *, store: ChatSessionStore, + registry_path: Path | None = None, session_id: str, turn_id: str, host_surface: str, @@ -269,7 +800,7 @@ def complete_attached_agent_turn( ) -> dict[str, Any]: """Write back one attached-host response with duplicate-safe receipts.""" - _require_attached_host( + session = _require_attached_host( store=store, session_id=session_id, host_surface=host_surface, @@ -285,14 +816,36 @@ def complete_attached_agent_turn( raise ValueError("response.message is required") if len(message) > 200_000: raise ValueError("response.message is too large") - _turn, created = store.complete_attached_turn( - session_id, - turn_id, - claim_id=claim_id, - completion_id=completion_id, - response=normalized_response, - agent_message=message, - ) + with _attached_session_lifetime( + store=store, + registry_path=registry_path, + session_id=session_id, + allow_closed=True, + ) as (_session, registry, current_ref): + session = _require_attached_host( + store=store, + session_id=session_id, + host_surface=host_surface, + host_session_id=host_session_id, + allow_closed=True, + ) + if registry is not None and current_ref is not None: + turn = store.load_turn(session_id, turn_id) + _lifecycle_decision( + operation="complete", + registry=registry, + current_ref=current_ref, + session=session, + turn=turn, + ) + _turn, created = store.complete_attached_turn( + session_id, + turn_id, + claim_id=claim_id, + completion_id=completion_id, + response=normalized_response, + agent_message=message, + ) return { "ok": True, "schema_version": ATTACHED_SESSION_BROKER_SCHEMA_VERSION, diff --git a/loopx/attached_session_api.py b/loopx/attached_session_api.py index e7ee6bd1f4..585290ab18 100644 --- a/loopx/attached_session_api.py +++ b/loopx/attached_session_api.py @@ -4,7 +4,10 @@ from typing import Any -from .attached_session import bind_attached_agent_session +from .attached_session import ( + bind_attached_agent_session, + load_attached_session_registry, +) def _compact_id(value: Any, *, limit: int = 160) -> str: @@ -60,10 +63,15 @@ def _attach_session(self) -> None: missing = [key for key, value in required.items() if not value] if missing: raise ValueError(f"missing attached session fields: {', '.join(missing)}") - registry, _goal = self._registry_and_goal(required["goal_id"]) + registry_path = getattr(self.server, "registry_path", None) + if registry_path is None: + registry, _goal = self._registry_and_goal(required["goal_id"]) + else: + registry = load_attached_session_registry(registry_path) packet = bind_attached_agent_session( store=self.server.chat_store, registry=registry, + registry_path=registry_path, goal_id=required["goal_id"], agent_id=required["agent_id"], host_surface=required["host_surface"], diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index fde39ee609..7b32d1e5bd 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -54,6 +54,12 @@ utc_now, ) from .chat_providers import ClaudeCodeAdapter, direct_model_from_environment +from .attached_session import ( + enqueue_attached_agent_turn, + require_current_attached_session, + resume_attached_agent_session, + select_current_attached_session, +) EventSink = Callable[[str, dict[str, Any]], None] @@ -565,6 +571,23 @@ def open_session( route_lock = self.session_open_locks.setdefault(route_key, threading.Lock()) with route_lock: latest = None + if not is_manager_channel(selected_channel): + exact_attached, strict_profile = select_current_attached_session( + store=self.store, + registry_path=self.registry_path, + goal_id=goal_id, + agent_id=agent_id, + channel_id=selected_channel, + ) + if strict_profile: + if mode == "resume_latest" and exact_attached is not None: + return exact_attached, True + raise CodexChatAgentError( + "Managed Chat execution is not qualified for the " + "source-session profile.", + error_code="source_session_managed_chat_unsupported", + gate=None, + ) if mode == "resume_latest": latest = self.store.latest_session( goal_id=None if is_manager_channel(selected_channel) else goal_id, @@ -906,8 +929,10 @@ def submit_turn( if session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: if attachments: raise ValueError("attached host session queue does not yet accept attachments") - return self.store.create_queued_turn( - session_id, + return enqueue_attached_agent_turn( + store=self.store, + registry_path=self.registry_path, + session_id=session_id, client_turn_id=client_turn_id, message=message, origin="web", @@ -1022,6 +1047,12 @@ def steer_active_turn( session = self.store.load_session(session_id) if session is None or session.get("status") == "closed": raise KeyError("chat session was not found") + if session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: + session = require_current_attached_session( + store=self.store, + registry_path=self.registry_path, + session_id=session_id, + ) receipt, created = self.store.create_ingress_receipt( session_id, client_ingress_id=client_ingress_id, @@ -1149,13 +1180,22 @@ def enqueue_turn( session = self.store.load_session(session_id) if session is None or session.get("status") == "closed": raise KeyError("chat session was not found") - turn, created = self.store.create_queued_turn( - session_id, - client_turn_id=client_turn_id, - message=message, - origin=origin, - ) - if session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED: + if session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: + turn, created = enqueue_attached_agent_turn( + store=self.store, + registry_path=self.registry_path, + session_id=session_id, + client_turn_id=client_turn_id, + message=message, + origin=origin, + ) + else: + turn, created = self.store.create_queued_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + origin=origin, + ) self.resume_session_queue( session_id=session_id, work_dir=work_dir, @@ -1753,15 +1793,11 @@ def resume_session(self, *, session_id: str, work_dir: Path, objective: str) -> raise KeyError("chat session was not found") self._check_codex_home(session) if session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: - if session.get("active_turn_id"): - return session - restored = self.store.update_session( - session_id, - status="ready", - active_turn_id=None, - last_error_code=None, + return resume_attached_agent_session( + store=self.store, + registry_path=self.registry_path, + session_id=session_id, ) - return restored with self.lock: current = self.adapters.get(session_id) adapter_healthy = current is not None and current.healthcheck() diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 4b5a650d5c..6cd8feb1dc 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -232,6 +232,7 @@ def create_session( self, *, goal_id: str, + goal_instance_id: str | None = None, agent_id: str, adapter_kind: str, upstream_thread_id: str, @@ -288,6 +289,16 @@ def create_session( "schema_version": CHAT_SESSION_SCHEMA_VERSION, "session_id": token, "goal_id": normalized_goal_id, + **( + { + "goal_instance_id": _opaque_id( + goal_instance_id, + field="goal_instance_id", + ) + } + if goal_instance_id is not None + else {} + ), "agent_id": normalized_agent_id, "executor_endpoint_id": normalized_executor_endpoint_id, "adapter_kind": normalized_adapter_kind, @@ -549,6 +560,41 @@ def latest_session( agent_id: str, channel_id: str | None = None, ) -> dict[str, Any] | None: + candidates = self.resumable_session_candidates( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + return max(candidates, key=lambda item: str(item.get("updated_at") or ""), default=None) + + def resumable_session_candidates( + self, + *, + goal_id: str | None, + agent_id: str, + channel_id: str | None = None, + ) -> list[dict[str, Any]]: + """Return matching storage facts without deciding Goal identity.""" + + return [ + candidate + for candidate in self.session_candidates( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + if candidate.get("status") in RESUMABLE_SESSION_STATES + ] + + def session_candidates( + self, + *, + goal_id: str | None, + agent_id: str, + channel_id: str | None = None, + ) -> list[dict[str, Any]]: + """Return matching Session records without lifecycle filtering.""" + if goal_id is None and channel_id is None: raise ValueError("channel_id is required when goal_id is omitted") selected_channel = channel_id or f"goal.{goal_id}" @@ -559,11 +605,10 @@ def latest_session( payload.get("schema_version") == CHAT_SESSION_SCHEMA_VERSION and (goal_id is None or payload.get("goal_id") == goal_id) and payload.get("agent_id") == agent_id - and payload.get("status") in RESUMABLE_SESSION_STATES and _session_channel(payload) == selected_channel ): candidates.append(payload) - return max(candidates, key=lambda item: str(item.get("updated_at") or ""), default=None) + return candidates def list_sessions( self, @@ -956,6 +1001,7 @@ def create_queued_turn( *, client_turn_id: str, message: str, + goal_instance_id: str | None = None, ttl_seconds: int = SESSION_QUEUE_TTL_SECONDS, origin: str = "external", ) -> tuple[dict[str, Any], bool]: @@ -997,6 +1043,16 @@ def create_queued_turn( "schema_version": CHAT_TURN_SCHEMA_VERSION, "turn_id": turn_id, "session_id": session_id, + **( + { + "goal_instance_id": _opaque_id( + goal_instance_id, + field="goal_instance_id", + ) + } + if goal_instance_id is not None + else {} + ), "client_turn_id": client_id, "status": "queued", "message": str(message), @@ -1109,6 +1165,7 @@ def claim_next_queued_turn( session_id: str, *, host_claim_id: str | None = None, + admitted_goal_instance_id: str | None = None, ) -> dict[str, Any] | None: """Atomically make the oldest live queued Turn active for its Session.""" @@ -1133,6 +1190,14 @@ def claim_next_queued_turn( and active.get("host_claim_id") == host_claim_id and active.get("status") in {"starting", "running"} ): + if admitted_goal_instance_id is not None and ( + active.get("goal_instance_id") != admitted_goal_instance_id + or active.get("admitted_goal_instance_id") + != admitted_goal_instance_id + ): + raise ValueError( + "active Turn Goal instance admission is invalid" + ) return active return None for turn in self._settle_expired_queued_turns( @@ -1140,6 +1205,12 @@ def claim_next_queued_turn( now=datetime.now(timezone.utc), ): turn_id = str(turn["turn_id"]) + if admitted_goal_instance_id is not None and ( + turn.get("goal_instance_id") != admitted_goal_instance_id + ): + raise ValueError( + "queued Turn Goal instance does not match its Session" + ) if host_claim_id: now_text = utc_now() turn.update( @@ -1149,6 +1220,16 @@ def claim_next_queued_turn( host_claim_id, field="host_claim_id", ), + **( + { + "admitted_goal_instance_id": _opaque_id( + admitted_goal_instance_id, + field="admitted_goal_instance_id", + ) + } + if admitted_goal_instance_id is not None + else {} + ), "started_at": turn.get("started_at") or now_text, "last_activity_at": now_text, } diff --git a/loopx/cli_commands/worker_bridge.py b/loopx/cli_commands/worker_bridge.py index 57ec986d22..48f1e87262 100644 --- a/loopx/cli_commands/worker_bridge.py +++ b/loopx/cli_commands/worker_bridge.py @@ -10,10 +10,10 @@ bind_attached_agent_session, claim_attached_agent_turn, complete_attached_agent_turn, + load_attached_session_registry, render_attached_session_broker_markdown, ) from ..chat_store import CHAT_SESSION_MODE_ATTACHED, ChatSessionStore -from ..history import load_registry from ..paths import resolve_runtime_root from ..worker_bridge import ( DEFAULT_ACTIVE_USER_CODEX_BIN, @@ -404,7 +404,7 @@ def handle_worker_bridge_command( try: if args.worker_bridge_command.startswith("attached-session-"): effective_registry_path = registry_path or Path(args.registry) - registry = load_registry(effective_registry_path) + registry = load_attached_session_registry(effective_registry_path) runtime_root = resolve_runtime_root( registry, args.runtime_root, @@ -416,6 +416,7 @@ def handle_worker_bridge_command( payload = bind_attached_agent_session( store=store, registry=registry, + registry_path=effective_registry_path, goal_id=args.goal_id, agent_id=args.agent_id, host_surface=args.host_surface, @@ -441,6 +442,7 @@ def handle_worker_bridge_command( elif args.worker_bridge_command == "attached-session-claim": payload = claim_attached_agent_turn( store=store, + registry_path=effective_registry_path, session_id=args.session_id, host_surface=args.host_surface, host_session_id=args.host_session_id, @@ -458,6 +460,7 @@ def handle_worker_bridge_command( raise ValueError("response JSON must be an object") payload = complete_attached_agent_turn( store=store, + registry_path=effective_registry_path, session_id=args.session_id, turn_id=args.turn_id, host_surface=args.host_surface, diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index ef844d79c6..e624e065bb 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -133,6 +133,7 @@ import { decideProjectSessionUnbind, } from "./goals/source_session_lifetime.ts"; import { decideFirstPartyHostRuntime } from "./goals/first_party_host_runtime.ts"; +import { decideChatSessionLifecycle } from "./goals/chat_session_lifecycle.ts"; import { evaluateDeliveryRoute, } from "./turn_driver/delivery_continuity.ts"; @@ -541,6 +542,7 @@ export function createEffectRuntimeHandlers( ["goal.source_session.unbind.decide", decideProjectSessionUnbind], ["goal.source_session.recreate.decide", decideGoalRecreation], ["goal.first_party_host_runtime.decide", decideFirstPartyHostRuntime], + ["goal.chat_session.lifecycle.decide", decideChatSessionLifecycle], ["goal.acceptance.inspect", inspectLocalGoalAcceptance], ["goal.acceptance.configure", commitLocalGoalAcceptance], ["goal.acceptance.verify.commit", commitLocalGoalAcceptanceVerification], diff --git a/loopx/control_plane/goals/chat_session_lifecycle.ts b/loopx/control_plane/goals/chat_session_lifecycle.ts new file mode 100644 index 0000000000..8d3bbea924 --- /dev/null +++ b/loopx/control_plane/goals/chat_session_lifecycle.ts @@ -0,0 +1,278 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { assertNever, jsonObject } from "../runtime_decode.ts"; +import { + parseExactGoalRef, + type ExactGoalRef, +} from "./goal_instance_identity.ts"; +import { SOURCE_SESSION_PROFILE_ID } from "./source_session_lifetime.ts"; + +type WireGoalRef = Readonly<{ + goal_id: string; + goal_instance_id: string; +}>; + +type SessionFact = + | Readonly<{ kind: "legacy"; sessionId: string; updatedAt: string }> + | Readonly<{ + kind: "exact"; + sessionId: string; + goalRef: ExactGoalRef; + updatedAt: string; + }>; + +type TurnFact = Readonly<{ + goalRef: ExactGoalRef | null; + admittedGoalRef: ExactGoalRef | null; +}>; + +type SelectionFacts = Readonly<{ + operation: "select"; + profileId: string | null; + currentGoalRef: ExactGoalRef | null; + candidates: readonly SessionFact[]; +}>; + +type AdmissionFacts = Readonly<{ + operation: "admit" | "claim" | "replay_claim" | "complete"; + profileId: string | null; + currentGoalRef: ExactGoalRef | null; + session: SessionFact; + turn: TurnFact | null; +}>; + +type LifecycleFacts = SelectionFacts | AdmissionFacts; + +export type ChatSessionLifecycleDecision = + | Readonly<{ kind: "legacy" }> + | Readonly<{ kind: "create"; goal_ref: WireGoalRef }> + | Readonly<{ + kind: "reuse"; + session_id: string; + goal_ref: WireGoalRef; + }> + | Readonly<{ kind: "allow_current"; goal_ref: WireGoalRef }> + | Readonly<{ kind: "allow_historical"; goal_ref: WireGoalRef }> + | Readonly<{ + kind: "reject"; + code: + | "goal_instance_id_missing" + | "stale_goal_instance" + | "turn_goal_instance_mismatch" + | "turn_not_admitted"; + }>; + +function requiredObject(value: unknown, label: string): JsonObject { + const result = jsonObject(value); + if (!result) { + throw new EffectRuntimeRequestError(`${label} must be an object`); + } + return result; +} + +function requiredString(value: unknown, label: string): string { + if (typeof value !== "string" || value.trim() === "") { + throw new EffectRuntimeRequestError(`${label} must be a non-empty string`); + } + return value; +} + +function optionalProfile(value: unknown): string | null { + if (value === null || value === undefined) return null; + return requiredString(value, "profile_id"); +} + +function exactGoalRef(value: unknown, label: string): ExactGoalRef { + const parsed = parseExactGoalRef(value); + if (parsed.kind === "invalid") { + throw new EffectRuntimeRequestError(`${label} ${parsed.issue}`); + } + return parsed.value; +} + +function optionalGoalRef(value: unknown, label: string): ExactGoalRef | null { + return value === null ? null : exactGoalRef(value, label); +} + +function wireGoalRef(value: ExactGoalRef): WireGoalRef { + return { + goal_id: value.goalId.value, + goal_instance_id: value.goalInstanceId.value, + }; +} + +function goalRefsEqual(left: ExactGoalRef, right: ExactGoalRef): boolean { + return left.goalId.value === right.goalId.value + && left.goalInstanceId.value === right.goalInstanceId.value; +} + +function sessionFact(value: unknown, label: string): SessionFact { + const raw = requiredObject(value, label); + const sessionId = requiredString(raw.session_id, `${label}.session_id`); + const updatedAt = requiredString(raw.updated_at, `${label}.updated_at`); + if (raw.goal_instance_id === null || raw.goal_instance_id === undefined) { + return { kind: "legacy", sessionId, updatedAt }; + } + return { + kind: "exact", + sessionId, + updatedAt, + goalRef: exactGoalRef( + { + goal_id: raw.goal_id, + goal_instance_id: raw.goal_instance_id, + }, + `${label}.goal_ref`, + ), + }; +} + +function turnFact(value: unknown): TurnFact | null { + if (value === null) return null; + const raw = requiredObject(value, "turn"); + const goalRef = raw.goal_instance_id === null + || raw.goal_instance_id === undefined + ? null + : exactGoalRef( + { + goal_id: raw.goal_id, + goal_instance_id: raw.goal_instance_id, + }, + "turn.goal_ref", + ); + const admittedGoalRef = raw.admitted_goal_instance_id === null + || raw.admitted_goal_instance_id === undefined + ? null + : exactGoalRef( + { + goal_id: raw.goal_id, + goal_instance_id: raw.admitted_goal_instance_id, + }, + "turn.admitted_goal_ref", + ); + return { goalRef, admittedGoalRef }; +} + +function decodeFacts(value: unknown): LifecycleFacts { + const raw = requiredObject(value, "Chat session lifecycle facts"); + const operation = requiredString(raw.operation, "operation"); + const profileId = optionalProfile(raw.profile_id); + const currentGoalRef = optionalGoalRef( + raw.current_goal_ref ?? null, + "current_goal_ref", + ); + if (operation === "select") { + if (!Array.isArray(raw.candidates)) { + throw new EffectRuntimeRequestError("candidates must be an array"); + } + return { + operation, + profileId, + currentGoalRef, + candidates: raw.candidates.map((candidate, index) => + sessionFact(candidate, `candidates[${index}]`) + ), + }; + } + if ( + operation === "admit" + || operation === "claim" + || operation === "replay_claim" + || operation === "complete" + ) { + return { + operation, + profileId, + currentGoalRef, + session: sessionFact(raw.session, "session"), + turn: turnFact(raw.turn ?? null), + }; + } + throw new EffectRuntimeRequestError("unsupported Chat session lifecycle operation"); +} + +function requireExactSession( + session: SessionFact, +): ExactGoalRef | ChatSessionLifecycleDecision { + return session.kind === "exact" + ? session.goalRef + : { kind: "reject", code: "goal_instance_id_missing" }; +} + +function decideSelection( + facts: SelectionFacts, +): ChatSessionLifecycleDecision { + if (facts.currentGoalRef === null) { + return { kind: "reject", code: "goal_instance_id_missing" }; + } + const currentGoalRef = facts.currentGoalRef; + const current = facts.candidates + .filter((candidate): candidate is Extract => + candidate.kind === "exact" + && goalRefsEqual(candidate.goalRef, currentGoalRef) + ) + .sort((left, right) => + right.updatedAt.localeCompare(left.updatedAt) + || right.sessionId.localeCompare(left.sessionId) + )[0]; + return current + ? { + kind: "reuse", + session_id: current.sessionId, + goal_ref: wireGoalRef(currentGoalRef), + } + : { kind: "create", goal_ref: wireGoalRef(currentGoalRef) }; +} + +function decideAdmission( + facts: AdmissionFacts, +): ChatSessionLifecycleDecision { + if (facts.currentGoalRef === null) { + return { kind: "reject", code: "goal_instance_id_missing" }; + } + const sessionGoalRef = requireExactSession(facts.session); + if (!("kind" in sessionGoalRef) || sessionGoalRef.kind !== "goal_ref") { + return sessionGoalRef; + } + if (facts.operation === "admit" || facts.operation === "claim") { + return goalRefsEqual(sessionGoalRef, facts.currentGoalRef) + ? { kind: "allow_current", goal_ref: wireGoalRef(sessionGoalRef) } + : { kind: "reject", code: "stale_goal_instance" }; + } + const turn = facts.turn; + if (turn === null || turn.goalRef === null) { + return { kind: "reject", code: "turn_goal_instance_mismatch" }; + } + if (!goalRefsEqual(turn.goalRef, sessionGoalRef)) { + return { kind: "reject", code: "turn_goal_instance_mismatch" }; + } + if ( + turn.admittedGoalRef === null + || !goalRefsEqual(turn.admittedGoalRef, sessionGoalRef) + ) { + return { kind: "reject", code: "turn_not_admitted" }; + } + return goalRefsEqual(sessionGoalRef, facts.currentGoalRef) + ? { kind: "allow_current", goal_ref: wireGoalRef(sessionGoalRef) } + : { kind: "allow_historical", goal_ref: wireGoalRef(sessionGoalRef) }; +} + +export function decideChatSessionLifecycle( + value: unknown, +): ChatSessionLifecycleDecision { + const facts = decodeFacts(value); + if (facts.profileId !== SOURCE_SESSION_PROFILE_ID) { + return { kind: "legacy" }; + } + switch (facts.operation) { + case "select": + return decideSelection(facts); + case "admit": + case "claim": + case "replay_claim": + case "complete": + return decideAdmission(facts); + default: + return assertNever(facts, "unsupported Chat session lifecycle facts"); + } +} diff --git a/loopx/semantics/goal_instance_binding_inventory_v1.json b/loopx/semantics/goal_instance_binding_inventory_v1.json index feb4c455bf..d829e0915f 100644 --- a/loopx/semantics/goal_instance_binding_inventory_v1.json +++ b/loopx/semantics/goal_instance_binding_inventory_v1.json @@ -16,19 +16,23 @@ "locator": "/chat/sessions//session.json", "revision_signal": "last_activity_at", "content_digest_signal": "none", - "observed_reference": "goal_id", + "observed_reference": "goal_id + goal_instance_id", "producer_sites": [ "loopx/attached_session.py::bind_attached_agent_session", - "loopx/chat_store.py::ChatSessionStore.create_session" + "loopx/attached_session.py::enqueue_attached_agent_turn", + "loopx/chat_store.py::ChatSessionStore.create_session", + "loopx/control_plane/goals/chat_session_lifecycle.ts::decideChatSessionLifecycle" ], "consumer_sites": [ + "loopx/attached_session.py::claim_attached_agent_turn", + "loopx/attached_session.py::complete_attached_agent_turn", "loopx/chat_runtime.py::ChatRuntimeController.resume_session" ], - "effect_boundary": "loopx/attached_session.py::bind_attached_agent_session", + "effect_boundary": "loopx/control_plane/goals/chat_session_lifecycle.ts::decideChatSessionLifecycle", "authority_role": "runtime_session_binding", - "current_identity_strength": "goal_alias_only", - "cleanup_support": "session_close", - "m1_disposition": "alias_only_inventory", + "current_identity_strength": "exact_goal_ref_enforced", + "cleanup_support": "historical_completion_and_session_close", + "m1_disposition": "m3_qualified", "target_milestone": "M3" }, { diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 09c271ee70..fefadf4945 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -37,6 +37,46 @@ "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/attached_session.py::._attached_session_lifetime::codec_read:load_project_registry#1", + "line": 495, + "column": 20, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/attached_session.py::._bind_source_session::codec_read:load_project_registry#1", + "line": 208, + "column": 20, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/attached_session.py::._strict_registry::codec_read:load_project_registry#1", + "line": 143, + "column": 16, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/attached_session.py::.load_attached_session_registry::codec_read:load_project_registry#1", + "line": 150, + "column": 12, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, + { + "site": "loopx/attached_session.py::.select_current_attached_session::codec_read:load_project_registry#1", + "line": 544, + "column": 20, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/authority.py::.import_doc_registry_authority::codec_read:load_project_registry#1", "line": 575, @@ -853,14 +893,6 @@ "api": "load_registry", "classification": "codec_api" }, - { - "site": "loopx/cli_commands/worker_bridge.py::.handle_worker_bridge_command::codec_read:load_registry#1", - "line": 407, - "column": 24, - "kind": "codec_read", - "api": "load_registry", - "classification": "codec_api" - }, { "site": "loopx/cli_rollout.py::.append_cli_rollout_event::codec_read:load_registry#1", "line": 50, diff --git a/tests/architecture/test_goal_instance_binding_inventory.py b/tests/architecture/test_goal_instance_binding_inventory.py index 497afb2203..4aab393574 100644 --- a/tests/architecture/test_goal_instance_binding_inventory.py +++ b/tests/architecture/test_goal_instance_binding_inventory.py @@ -53,6 +53,7 @@ PARTIALLY_ENFORCED_OWNER_IDS = { "first_party_host_runtime", } +QUALIFIED_OWNER_IDS = {"attached_host_chat_session"} TYPESCRIPT_DECLARATION = re.compile( r"^(?:export\s+)?(?:async\s+)?(?:function|class)\s+([A-Za-z_$][A-Za-z0-9_$]*)", re.MULTILINE, @@ -162,6 +163,9 @@ def test_m1_observation_claims_are_bounded_to_the_selected_lifecycle() -> None: elif owner["owner_id"] in PARTIALLY_ENFORCED_OWNER_IDS: assert owner["m1_disposition"] == "source_exact_partial_enforcement" assert owner["current_identity_strength"] == "source_exact_partial" + elif owner["owner_id"] in QUALIFIED_OWNER_IDS: + assert owner["m1_disposition"] == "m3_qualified" + assert owner["current_identity_strength"] == "exact_goal_ref_enforced" assert owner["target_milestone"] == "M3" else: assert owner["m1_disposition"] == "alias_only_inventory" diff --git a/tests/architecture/test_source_session_registry_denial.py b/tests/architecture/test_source_session_registry_denial.py index d4887729f4..717b41255d 100644 --- a/tests/architecture/test_source_session_registry_denial.py +++ b/tests/architecture/test_source_session_registry_denial.py @@ -7,6 +7,7 @@ REPO_ROOT = Path(__file__).resolve().parents[2] DIRECT_LOADER_ALLOWLIST = { "loopx/authority.py", + "loopx/attached_session.py", "loopx/bootstrap.py", "loopx/claude_goal_mode/scripts/connect.py", "loopx/configure_goal.py", @@ -33,12 +34,18 @@ def test_direct_project_registry_loaders_have_source_session_denial() -> None: assert callers == DIRECT_LOADER_ALLOWLIST source_session_owners = { + "loopx/attached_session.py", "loopx/control_plane/goals/first_party_host_admission.py", "loopx/control_plane/projects/registry.py", } for relative in callers - source_session_owners: source = (REPO_ROOT / relative).read_text(encoding="utf-8") assert "require_runtime_compatible_project_registry(" in source, relative + attached_owner = (REPO_ROOT / "loopx/attached_session.py").read_text( + encoding="utf-8" + ) + assert "SOURCE_SESSION_PROFILE_ID" in attached_owner + assert "source_session_goal_lifetime" in attached_owner def test_generic_registry_decoder_enforces_source_session_denial() -> None: diff --git a/tests/control_plane_ts/chat_session_lifecycle.test.ts b/tests/control_plane_ts/chat_session_lifecycle.test.ts new file mode 100644 index 0000000000..b4635aab75 --- /dev/null +++ b/tests/control_plane_ts/chat_session_lifecycle.test.ts @@ -0,0 +1,167 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { + decideChatSessionLifecycle, +} from "../../loopx/control_plane/goals/chat_session_lifecycle.ts"; + +const goalA = { + goal_id: "release", + goal_instance_id: "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", +}; +const goalB = { + goal_id: "release", + goal_instance_id: "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb", +}; + +function session( + sessionId: string, + goalRef: typeof goalA | null, + updatedAt = "2026-09-26T00:00:00Z", +) { + return { + session_id: sessionId, + goal_id: "release", + goal_instance_id: goalRef?.goal_instance_id ?? null, + updated_at: updatedAt, + }; +} + +test("legacy profiles preserve the existing path", () => { + assert.deepEqual( + decideChatSessionLifecycle({ + operation: "select", + profile_id: null, + current_goal_ref: null, + candidates: [], + }), + { kind: "legacy" }, + ); +}); + +test("selection reuses only the newest exact current session", () => { + assert.deepEqual( + decideChatSessionLifecycle({ + operation: "select", + profile_id: "source_session_v1", + current_goal_ref: goalB, + candidates: [ + session("legacy", null, "2026-09-26T03:00:00Z"), + session("stale-a", goalA, "2026-09-26T04:00:00Z"), + session("current-old", goalB, "2026-09-26T01:00:00Z"), + session("current-new", goalB, "2026-09-26T02:00:00Z"), + ], + }), + { kind: "reuse", session_id: "current-new", goal_ref: goalB }, + ); +}); + +test("selection creates a new exact session when only stale state exists", () => { + assert.deepEqual( + decideChatSessionLifecycle({ + operation: "select", + profile_id: "source_session_v1", + current_goal_ref: goalB, + candidates: [session("stale-a", goalA), session("legacy", null)], + }), + { kind: "create", goal_ref: goalB }, + ); +}); + +test("new work requires the session GoalRef to be current", () => { + const base = { + profile_id: "source_session_v1", + current_goal_ref: goalB, + session: session("session-a", goalA), + turn: null, + }; + for (const operation of ["admit", "claim"] as const) { + assert.deepEqual( + decideChatSessionLifecycle({ ...base, operation }), + { kind: "reject", code: "stale_goal_instance" }, + ); + } + assert.deepEqual( + decideChatSessionLifecycle({ + ...base, + operation: "admit", + session: session("session-b", goalB), + }), + { kind: "allow_current", goal_ref: goalB }, + ); +}); + +test("claim replay and completion allow exact historical admitted work", () => { + const historical = { + profile_id: "source_session_v1", + current_goal_ref: goalB, + session: session("session-a", goalA), + turn: { + goal_id: goalA.goal_id, + goal_instance_id: goalA.goal_instance_id, + admitted_goal_instance_id: goalA.goal_instance_id, + }, + }; + for (const operation of ["replay_claim", "complete"] as const) { + assert.deepEqual( + decideChatSessionLifecycle({ ...historical, operation }), + { kind: "allow_historical", goal_ref: goalA }, + ); + } +}); + +test("historical completion rejects unstamped and mismatched turns", () => { + const base = { + operation: "complete", + profile_id: "source_session_v1", + current_goal_ref: goalB, + session: session("session-a", goalA), + }; + assert.deepEqual( + decideChatSessionLifecycle({ + ...base, + turn: { + goal_id: goalA.goal_id, + goal_instance_id: goalA.goal_instance_id, + admitted_goal_instance_id: null, + }, + }), + { kind: "reject", code: "turn_not_admitted" }, + ); + assert.deepEqual( + decideChatSessionLifecycle({ + ...base, + turn: { + goal_id: goalA.goal_id, + goal_instance_id: goalB.goal_instance_id, + admitted_goal_instance_id: goalB.goal_instance_id, + }, + }), + { kind: "reject", code: "turn_goal_instance_mismatch" }, + ); +}); + +test("strict operations reject unstamped sessions and malformed GoalRefs", () => { + assert.deepEqual( + decideChatSessionLifecycle({ + operation: "admit", + profile_id: "source_session_v1", + current_goal_ref: goalA, + session: session("legacy", null), + turn: null, + }), + { kind: "reject", code: "goal_instance_id_missing" }, + ); + assert.throws( + () => decideChatSessionLifecycle({ + operation: "select", + profile_id: "source_session_v1", + current_goal_ref: { + goal_id: "release", + goal_instance_id: "bad", + }, + candidates: [], + }), + /invalid_goal_instance_id/, + ); +}); diff --git a/tests/control_plane_ts/effect_runtime_handlers.test.ts b/tests/control_plane_ts/effect_runtime_handlers.test.ts index 03b9bf22bd..22b6b52ddf 100644 --- a/tests/control_plane_ts/effect_runtime_handlers.test.ts +++ b/tests/control_plane_ts/effect_runtime_handlers.test.ts @@ -201,6 +201,19 @@ test("runtime exposes the source-session lifetime decisions", async () => { ).kind, "commit", ); + assert.deepEqual( + await dispatchEffectRuntimeMethod( + handlers, + "goal.chat_session.lifecycle.decide", + { + operation: "select", + profile_id: "source_session_v1", + current_goal_ref: goalRef, + candidates: [], + }, + ), + { kind: "create", goal_ref: goalRef }, + ); }); test("runtime boundary registers the quota monitor-poll transaction", async () => { diff --git a/tests/test_attached_session_goal_instance.py b/tests/test_attached_session_goal_instance.py new file mode 100644 index 0000000000..b13ab987d2 --- /dev/null +++ b/tests/test_attached_session_goal_instance.py @@ -0,0 +1,570 @@ +from __future__ import annotations + +import json +import threading +import time +from pathlib import Path + +import pytest + +from loopx.attached_session import ( + bind_attached_agent_session, + claim_attached_agent_turn, + complete_attached_agent_turn, + select_current_attached_session, +) +from loopx.chat_agent import CodexChatAgentError +from loopx.cli import main +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_store import ChatSessionStore +from loopx.control_plane.goals.source_session_recreation import ( + RecreateGoalRequest, + recreate_goal_instance, +) +from loopx.control_plane.goals.source_session_registration import ( + FreshSourceSessionRegistration, + register_fresh_source_session_project, +) +from loopx.control_plane.projects.registry_codec import load_project_registry + + +GOAL_ID = "release" +AGENT_ID = "release-worker" +HOST_SURFACE = "codex-app-ssh" +HOST_SESSION_ID = "host-thread" + + +def _register(tmp_path: Path) -> tuple[Path, Path, ChatSessionStore, str]: + project_root = tmp_path / "project" + project_root.mkdir() + registry_path = project_root / ".loopx" / "registry.json" + runtime_root = tmp_path / "runtime" + result = register_fresh_source_session_project( + FreshSourceSessionRegistration( + registry_path=registry_path, + runtime_root=runtime_root, + operation_id="register-release", + project_id="project", + goal_id=GOAL_ID, + objective="Ship the release.", + non_goals=[], + acceptance=["The release is verified."], + unknowns=[], + next_effect="Inspect the release.", + stop_condition="Stop when authority is absent.", + project_record={ + "id": "project", + "kind": "work", + "path": str(project_root), + }, + goal_record={ + "id": GOAL_ID, + "project_id": "project", + "title": "Release", + "status": "active", + "repo": str(project_root), + "coordination": { + "registered_agents": [AGENT_ID], + "thread_agent_bindings": [ + { + "host_surface": HOST_SURFACE, + "thread_id": HOST_SESSION_ID, + "agent_id": AGENT_ID, + } + ], + }, + }, + state_file=project_root / "GOAL.md", + ) + ) + return ( + registry_path, + runtime_root, + ChatSessionStore(runtime_root), + str(result["goal_ref"]["goal_instance_id"]), + ) + + +def _bind( + store: ChatSessionStore, + registry_path: Path, +) -> dict[str, object]: + return bind_attached_agent_session( + store=store, + registry=load_project_registry(registry_path), + registry_path=registry_path, + goal_id=GOAL_ID, + agent_id=AGENT_ID, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + executor_endpoint_id="codex", + execute=True, + ) + + +def _recreate(registry_path: Path, instance_id: str) -> str: + result = recreate_goal_instance( + RecreateGoalRequest( + registry_path=registry_path, + goal_id=GOAL_ID, + goal_instance_id=instance_id, + operation_id=f"recreate-{instance_id}", + ) + ) + return str(result["goal_ref"]["goal_instance_id"]) + + +def test_recreated_goal_cannot_reuse_or_claim_from_attached_session( + tmp_path: Path, +) -> None: + registry_path, _runtime_root, store, instance_a = _register(tmp_path) + bound_a = _bind(store, registry_path) + session_a = str(bound_a["session"]["session_id"]) # type: ignore[index] + persisted_a = store.load_session(session_a) + assert persisted_a is not None + assert persisted_a["goal_instance_id"] == instance_a + assert "goal_instance_id" not in bound_a["session"] # type: ignore[operator] + + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + first, created = runtime.submit_turn( + session_id=session_a, + client_turn_id="first", + message="first message", + work_dir=tmp_path, + objective="Ship the release.", + ) + second, queued = runtime.enqueue_turn( + session_id=session_a, + client_turn_id="second", + message="second message", + work_dir=tmp_path, + objective="Ship the release.", + ) + assert created and queued + assert first["goal_instance_id"] == instance_a + assert second["goal_instance_id"] == instance_a + + claimed = claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_a, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="claim-a", + ) + assert claimed["turn"]["turn_id"] == first["turn_id"] + admitted = store.load_turn(session_a, str(first["turn_id"])) + assert admitted is not None + assert admitted["admitted_goal_instance_id"] == instance_a + + instance_b = _recreate(registry_path, instance_a) + assert instance_b != instance_a + bound_b = _bind(store, registry_path) + session_b = str(bound_b["session"]["session_id"]) # type: ignore[index] + assert session_b != session_a + persisted_b = store.load_session(session_b) + assert persisted_b is not None + assert persisted_b["goal_instance_id"] == instance_b + before_b = persisted_b.copy() + + replay = claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_a, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="claim-a", + ) + assert replay["turn"]["turn_id"] == first["turn_id"] + + completed = complete_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_a, + turn_id=str(first["turn_id"]), + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="claim-a", + completion_id="completion-a", + response={"message": "historical A result"}, + ) + assert completed["created"] is True + assert store.load_session(session_b) == before_b + assert store.load_turn(session_a, str(first["turn_id"]))["status"] == "completed" # type: ignore[index] + + with pytest.raises(ValueError, match="stale_goal_instance"): + claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_a, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="claim-second", + ) + with pytest.raises(ValueError, match="stale_goal_instance"): + runtime.resume_session( + session_id=session_a, + work_dir=tmp_path, + objective="Ship the release.", + ) + assert store.load_turn(session_a, str(second["turn_id"]))["status"] == "queued" # type: ignore[index] + + +def test_stale_attached_enqueue_is_rejected_without_mutation( + tmp_path: Path, +) -> None: + registry_path, _runtime_root, store, instance_a = _register(tmp_path) + session_a = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] + _recreate(registry_path, instance_a) + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + before = store.load_session(session_a) + + with pytest.raises(ValueError, match="stale_goal_instance"): + runtime.submit_turn( + session_id=session_a, + client_turn_id="late", + message="must not enter A", + work_dir=tmp_path, + objective="Ship the release.", + ) + + assert store.load_session(session_a) == before + assert store.turn_for_client(session_a, "late") is None + + +@pytest.mark.parametrize( + "operation", + ["submit", "enqueue", "resume", "claim", "complete"], +) +def test_strict_profile_rejects_unstamped_attached_session( + tmp_path: Path, + operation: str, +) -> None: + registry_path, _runtime_root, store, _instance_a = _register(tmp_path) + session = store.create_session( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + adapter_kind="attached_host_session", + upstream_thread_id=HOST_SESSION_ID, + session_mode="attached_host", + host_surface=HOST_SURFACE, + ) + session_id = str(session["session_id"]) + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + turn: dict[str, object] | None = None + if operation in {"claim", "complete"}: + turn, _created = store.create_queued_turn( + session_id, + client_turn_id=f"{operation}-turn", + message="must reject", + ) + if operation == "complete": + claimed = store.claim_next_queued_turn( + session_id, + host_claim_id="legacy-claim", + ) + assert claimed is not None + turn = claimed + + before_session = store.load_session(session_id) + before_turn = ( + store.load_turn(session_id, str(turn["turn_id"])) if turn is not None else None + ) + with pytest.raises(ValueError, match="goal_instance_id_missing"): + if operation == "submit": + runtime.submit_turn( + session_id=session_id, + client_turn_id="strict-submit", + message="must reject", + work_dir=tmp_path, + objective="Ship the release.", + ) + elif operation == "enqueue": + runtime.enqueue_turn( + session_id=session_id, + client_turn_id="strict-enqueue", + message="must reject", + work_dir=tmp_path, + objective="Ship the release.", + ) + elif operation == "resume": + runtime.resume_session( + session_id=session_id, + work_dir=tmp_path, + objective="Ship the release.", + ) + elif operation == "claim": + claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_id, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="strict-claim", + ) + else: + assert turn is not None + complete_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_id, + turn_id=str(turn["turn_id"]), + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="legacy-claim", + completion_id="strict-completion", + response={"message": "must reject"}, + ) + + assert store.load_session(session_id) == before_session + if turn is not None: + assert store.load_turn(session_id, str(turn["turn_id"])) == before_turn + assert store.turn_for_client(session_id, "strict-submit") is None + assert store.turn_for_client(session_id, "strict-enqueue") is None + + +def test_claim_wait_does_not_hold_goal_lifetime_guard(tmp_path: Path) -> None: + registry_path, _runtime_root, store, instance_a = _register(tmp_path) + session_a = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] + finished = threading.Event() + errors: list[BaseException] = [] + + def wait_for_claim() -> None: + try: + claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_a, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="waiting-claim", + wait_seconds=0.5, + ) + except ValueError as exc: + if "stale_goal_instance" not in str(exc): + errors.append(exc) + finally: + finished.set() + + thread = threading.Thread(target=wait_for_claim) + thread.start() + time.sleep(0.1) + started = time.monotonic() + _recreate(registry_path, instance_a) + elapsed = time.monotonic() - started + thread.join(timeout=2) + + assert elapsed < 0.4 + assert finished.is_set() + assert errors == [] + + +def test_resume_writeback_holds_goal_lifetime_guard( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry_path, _runtime_root, store, instance_a = _register(tmp_path) + session_a = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + update_entered = threading.Event() + release_update = threading.Event() + recreation_finished = threading.Event() + resumed: list[dict[str, object]] = [] + recreated: list[str] = [] + errors: list[BaseException] = [] + original_update = store.update_session + + def blocked_update(session_id: str, **changes: object) -> dict[str, object]: + update_entered.set() + if not release_update.wait(timeout=2): + raise TimeoutError("resume writeback was not released") + return original_update(session_id, **changes) + + monkeypatch.setattr(store, "update_session", blocked_update) + + def resume() -> None: + try: + resumed.append( + runtime.resume_session( + session_id=session_a, + work_dir=tmp_path, + objective="Ship the release.", + ) + ) + except BaseException as exc: + errors.append(exc) + + def recreate() -> None: + try: + recreated.append(_recreate(registry_path, instance_a)) + except BaseException as exc: + errors.append(exc) + finally: + recreation_finished.set() + + resume_thread = threading.Thread(target=resume) + recreate_thread = threading.Thread(target=recreate) + resume_thread.start() + assert update_entered.wait(timeout=2) + recreate_thread.start() + try: + assert not recreation_finished.wait(timeout=0.1) + finally: + release_update.set() + resume_thread.join(timeout=2) + recreate_thread.join(timeout=2) + + assert not resume_thread.is_alive() + assert not recreate_thread.is_alive() + assert errors == [] + assert resumed[0]["session_id"] == session_a + assert recreated[0] != instance_a + + +def test_strict_managed_open_fails_before_provider_start( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry_path, _runtime_root, store, _instance_a = _register(tmp_path) + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + monkeypatch.setattr( + runtime, + "_start_adapter", + lambda **_kwargs: pytest.fail("provider start must remain fenced"), + ) + + with pytest.raises(CodexChatAgentError) as raised: + runtime.open_session( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + work_dir=tmp_path, + objective="Ship the release.", + mode="new", + ) + + assert raised.value.error_code == "source_session_managed_chat_unsupported" + + +def test_strict_session_does_not_fall_back_when_registry_disappears( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + registry_path, _runtime_root, store, _instance_a = _register(tmp_path) + _bind(store, registry_path) + registry_path.unlink() + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + monkeypatch.setattr( + runtime, + "_start_adapter", + lambda **_kwargs: pytest.fail("provider start must remain fenced"), + ) + + with pytest.raises(FileNotFoundError): + runtime.open_session( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + work_dir=tmp_path, + objective="Ship the release.", + mode="resume_latest", + ) + + +def test_missing_registry_without_exact_history_preserves_legacy_selection( + tmp_path: Path, +) -> None: + selected, strict = select_current_attached_session( + store=ChatSessionStore(tmp_path / "runtime"), + registry_path=tmp_path / "missing-registry.json", + goal_id=GOAL_ID, + agent_id=AGENT_ID, + channel_id=f"goal.{GOAL_ID}", + ) + + assert selected is None + assert strict is False + + +def test_worker_bridge_cli_binds_source_session_without_identity_input( + tmp_path: Path, + capsys: pytest.CaptureFixture[str], +) -> None: + registry_path, runtime_root, _store, instance_a = _register(tmp_path) + + assert ( + main( + [ + "--registry", + str(registry_path), + "--runtime-root", + str(runtime_root), + "--format", + "json", + "worker-bridge", + "attached-session-bind", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--host-surface", + HOST_SURFACE, + "--host-session-id", + HOST_SESSION_ID, + "--executor-endpoint-id", + "codex", + "--execute", + ] + ) + == 0 + ) + + payload = json.loads(capsys.readouterr().out) + session_id = str(payload["session"]["session_id"]) + persisted = ChatSessionStore(runtime_root).load_session(session_id) + assert persisted is not None + assert persisted["goal_instance_id"] == instance_a + assert "goal_instance_id" not in payload["session"] + + +def test_legacy_writers_omit_goal_instance_fields(tmp_path: Path) -> None: + store = ChatSessionStore(tmp_path) + session = store.create_session( + goal_id="legacy", + agent_id="worker", + adapter_kind="attached_host_session", + upstream_thread_id="host", + session_mode="attached_host", + host_surface="host", + ) + turn, _created = store.create_queued_turn( + str(session["session_id"]), + client_turn_id="legacy-turn", + message="legacy", + ) + + assert "goal_instance_id" not in session + assert "goal_instance_id" not in turn + assert "admitted_goal_instance_id" not in turn diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 1a9a48fc8e..cece05dccc 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -51,6 +51,7 @@ "loopx/control_plane/goals/shared_goal_alignment.ts", "loopx/control_plane/goals/goal_amendment_proposal.ts", "loopx/control_plane/goals/goal_instance_identity.ts", + "loopx/control_plane/goals/chat_session_lifecycle.ts", "loopx/control_plane/quota/settlement_workspace_causality.ts", "loopx/control_plane/quota/settlement_readback.ts", "loopx/control_plane/quota/heartbeat_receipt_identity.ts", @@ -160,6 +161,7 @@ "tests/control_plane_ts/replan_history.test.ts", "tests/control_plane_ts/goal_amendment_proposal.test.ts", "tests/control_plane_ts/goal_instance_identity.test.ts", + "tests/control_plane_ts/chat_session_lifecycle.test.ts", "tests/control_plane_ts/test_python_runtime.test.ts" ] } From 5ed5743a14f0b105dbf950d49bc7560423789a2c Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 28 Sep 2026 14:36:36 +0800 Subject: [PATCH 2/4] fix(ci): refresh registry I/O manifest after main sync Signed-off-by: duanjialing.777 --- loopx/semantics/project_registry_io_manifest_v1.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index fefadf4945..ecd675df71 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -919,7 +919,7 @@ }, { "site": "loopx/contract.py::.check_contract::codec_read:load_registry#1", - "line": 1014, + "line": 1027, "column": 20, "kind": "codec_read", "api": "load_registry", From 57699598dea45ad53e7c1262776a409b5d02616a Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 28 Sep 2026 14:46:42 +0800 Subject: [PATCH 3/4] fix(chat): fail closed when source context disappears Signed-off-by: duanjialing.777 --- loopx/attached_session.py | 84 +++++--- .../project_registry_io_manifest_v1.json | 8 +- .../references/repair-patterns.md | 1 + tests/test_attached_session_goal_instance.py | 198 +++++++++++++++++- 4 files changed, 255 insertions(+), 36 deletions(-) diff --git a/loopx/attached_session.py b/loopx/attached_session.py index 6c1835776b..2b537ebc28 100644 --- a/loopx/attached_session.py +++ b/loopx/attached_session.py @@ -144,6 +144,32 @@ def _strict_registry( return registry if registry.get("profile_id") == SOURCE_SESSION_PROFILE_ID else None +def _source_registry_path_or_legacy( + *, + store: ChatSessionStore, + registry_path: Path | None, + goal_id: str, + agent_id: str, + channel_id: str, +) -> Path | None: + """Resolve source context without downgrading known exact history.""" + + if registry_path is not None and registry_path.exists(): + return registry_path + candidates = store.session_candidates( + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + if not any( + candidate.get("goal_instance_id") is not None for candidate in candidates + ): + return None + if registry_path is None: + raise ValueError("source-session attached operation requires registry_path") + raise FileNotFoundError(registry_path) + + def load_attached_session_registry(registry_path: Path) -> dict[str, Any]: """Load the source-aware registry for this qualified owner only.""" @@ -465,9 +491,7 @@ def _attached_session_lifetime( registry_path: Path | None, session_id: str, allow_closed: bool = False, -) -> Iterator[ - tuple[dict[str, Any], dict[str, Any] | None, dict[str, str] | None] -]: +) -> Iterator[tuple[dict[str, Any], dict[str, Any] | None, dict[str, str] | None]]: """Load authoritative Goal context before choosing strict or legacy behavior.""" session = store.load_session(session_id) @@ -475,28 +499,25 @@ def _attached_session_lifetime( raise KeyError("attached Agent session was not found") if session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED: raise ValueError("the selected Session is not an attached host session") - exact_session = session.get("goal_instance_id") is not None - if registry_path is None: - if exact_session: - raise ValueError("source-session attached operation requires registry_path") - yield session, None, None - return - if not registry_path.exists(): - if exact_session: - raise FileNotFoundError(registry_path) + goal_id = str(session.get("goal_id") or "") + source_registry_path = _source_registry_path_or_legacy( + store=store, + registry_path=registry_path, + goal_id=goal_id, + agent_id=str(session.get("agent_id") or ""), + channel_id=str(session.get("channel_id") or f"goal.{goal_id}"), + ) + if source_registry_path is None: yield session, None, None return - goal_id = str(session.get("goal_id") or "") with exclusive_cross_runtime_file_lock( - guard_path(registry_path, goal_id), + guard_path(source_registry_path, goal_id), operation="source_session_goal_lifetime", ): - registry = load_project_registry(registry_path) + registry = load_project_registry(source_registry_path) session = store.load_session(session_id) - if session is None or ( - session.get("status") == "closed" and not allow_closed - ): + if session is None or (session.get("status") == "closed" and not allow_closed): raise KeyError("attached Agent session was not found") if session.get("session_mode") != CHAT_SESSION_MODE_ATTACHED: raise ValueError("the selected Session is not an attached host session") @@ -517,15 +538,14 @@ def select_current_attached_session( ) -> tuple[dict[str, Any] | None, bool]: """Select an attached Session and report whether strict identity is active.""" - if registry_path is not None and not registry_path.exists(): - candidates = store.session_candidates( - goal_id=goal_id, - agent_id=agent_id, - channel_id=channel_id, - ) - if any(candidate.get("goal_instance_id") is not None for candidate in candidates): - raise FileNotFoundError(registry_path) - registry = _strict_registry(registry_path) + source_registry_path = _source_registry_path_or_legacy( + store=store, + registry_path=registry_path, + goal_id=goal_id, + agent_id=agent_id, + channel_id=channel_id, + ) + registry = _strict_registry(source_registry_path) if registry is None: return ( store.latest_session( @@ -535,8 +555,8 @@ def select_current_attached_session( ), False, ) - assert registry_path is not None - guard = guard_path(registry_path, goal_id) + assert source_registry_path is not None + guard = guard_path(source_registry_path, goal_id) with exclusive_cross_runtime_file_lock( guard, operation="source_session_goal_lifetime", @@ -873,7 +893,11 @@ def render_attached_session_broker_markdown(payload: dict[str, Any]) -> str: f"{turn.get('message')}" ) session = payload.get("session") - session_id = session.get("session_id") if isinstance(session, dict) else payload.get("session_id") + session_id = ( + session.get("session_id") + if isinstance(session, dict) + else payload.get("session_id") + ) return ( "# Attached Agent Session\n\n" f"- Action: `{action}`\n" diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index ecd675df71..f07dc4516c 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -39,7 +39,7 @@ }, { "site": "loopx/attached_session.py::._attached_session_lifetime::codec_read:load_project_registry#1", - "line": 495, + "line": 518, "column": 20, "kind": "codec_read", "api": "load_project_registry", @@ -47,7 +47,7 @@ }, { "site": "loopx/attached_session.py::._bind_source_session::codec_read:load_project_registry#1", - "line": 208, + "line": 234, "column": 20, "kind": "codec_read", "api": "load_project_registry", @@ -63,7 +63,7 @@ }, { "site": "loopx/attached_session.py::.load_attached_session_registry::codec_read:load_project_registry#1", - "line": 150, + "line": 176, "column": 12, "kind": "codec_read", "api": "load_project_registry", @@ -71,7 +71,7 @@ }, { "site": "loopx/attached_session.py::.select_current_attached_session::codec_read:load_project_registry#1", - "line": 544, + "line": 564, "column": 20, "kind": "codec_read", "api": "load_project_registry", diff --git a/skills/loopx-self-repair/references/repair-patterns.md b/skills/loopx-self-repair/references/repair-patterns.md index 7a7f5a91ee..e00ca3309b 100644 --- a/skills/loopx-self-repair/references/repair-patterns.md +++ b/skills/loopx-self-repair/references/repair-patterns.md @@ -13,6 +13,7 @@ teaches a reusable control-plane lesson. | `acceptance_validation_failure_flattened` | A bound Todo is `ready` and its ordinary validator passes, yet completion reports only `goal_acceptance_validation_rejected`; the agent searches contract bindings before discovering a workspace or runner failure. | Exact Todo acceptance state, completion's typed criterion failure, recorded delivery workspace, current worktree cleanliness, and the privacy-safe runner receipt. | The completion boundary collapsed a failed criterion receipt into a generic contract error, hiding the workspace or command status. | Project the failed criterion ID, allowlisted validation status, safe exit code and bounded next action without command output or local paths. Repair the execution context and retry under the same Turn/lease; do not rebind owner criteria or infer a Goal-wide hold. | | `review_compatibility_assumption_gap` | A correct fix retains parallel protocol paths, and the review treats historical receipt recovery as proof that every old request decoder is needed. | Published review, structured compatibility rationale, real caller/deployment inventory, stored request versus receipt shape, and smaller-design readback. | Re-review verifies the last bug but does not separate compatibility obligations or test consolidation; free-text claims pass as evidence. | Replace the capability's existing compatibility rationale with a bounded structured assessment. Distinguish transient requests, independent client rollout and persisted replay formats; compare one typed current contract; preserve real legacy consumers. Cover needless retention and unsafe removal with positive twins. Keep optional simplifications advisory and do not make field completeness certify evidence truth. | | `review_outcome_continuity_gap` | Local feature checks pass while later work starves, recovery only records a blocker, or ordinary users face new repeated intervention. | Accepted product outcome, published review, real later invocations, affected user surfaces, durable progress and authorized recovery. | Review proved one operation or prevention of bypass without checking sustained progress and the complete user journey. | Make long-horizon continuity and user experience explicit judgments in the existing problem context. Reuse bounded real walkthroughs; apply scope counterexamples when gates are involved. Accept deliberate waits only with an independent basis and recovery or safe terminal route. Qualify with exact-source historical counterexamples and repaired positive controls; keep old verdicts out of model inputs and distinguish reported historical checks from newly executed tests. Section length and declared evidence cannot certify truth. | +| `source_context_legacy_downgrade` | An unstamped record mutates state through a legacy path after the authoritative source disappears, although another record in the same owner context preserves an exact source identity. | Selected record identity, all same-context durable candidates, source path availability, operation entrypoint, and business-file bytes before and after refusal. | The adapter treated one selected record's missing identity as proof that the whole context was legacy and ignored exact historical evidence when the current authority could not be read. | Resolve source-context availability once for selection and mutation. Exact history may prove that legacy fallback is unsafe, but it must never supply or guess the current identity. Fail closed before writes when authority is unavailable; preserve contexts with no exact history through opposite legacy controls across every mutation entrypoint. | | `archive_capture_classification_gap` | Whole-Goal capture rejects a reachable archived Agent record although its active read was valid. | Recorded role/class, existing legacy read classification, transitive dependency closure, bootstrap and writer-outbox readback. | Archive storage preserved the role but omitted the resolved class; capture treated missing class as missing authority. | Keep recorded identity separate from compatibility classification. Only a recorded Agent role can adopt the existing read class; use it consistently for closure and materialization. Preserve the class on new archive moves, keep user authority fail-closed, and validate full source capture in disposable real providers without changing the active Goal. | | `shadow_proof_transport_amplification` | A large Goal cannot capture its first mutation or finish drain although the same tiny Goal succeeds; RPC rejects an oversized response. | Same source population, base/head serialized response sizes, actual sequence/drain callers, qualified lineage and cursor readback. | A consumer needing progress or partition markers received the full head and every historical projection across the language boundary. | Keep full history verification in the typed owner and return a purpose-specific compact proof. Preserve receipts, sequence and lineage checks; do not raise transport limits, truncate source records or weaken qualification to make the test pass. Cover the old oversized response and real CLI capture, drain and reviewed cutover on a disposable snapshot. | | `qualification_host_contract_mismatch` | A model stops on ordinary inspection, shell composition or draft correction and the result is reported as a semantic-control failure. | Actual synthetic operation, advertised tool contract, OS isolation, subprocess exit status, returned diagnostics and durable writeback attempts. | A shell-labelled host imposed a separate command language or ended execution without returning normal tool errors. | Use a normal shell inside an isolated execution environment, with real CLI effects supervised at their existing authority boundary. Observe source evidence and durable outcomes instead of requiring a command spelling or read ritual. Return errors within a disclosed scenario budget; keep original inputs and authority stores protected. Separate host rejection, budget exhaustion and core semantic admission, retaining earlier failures. Do not insert model answers, waive evidence or grow a command whitelist one failed trajectory at a time. | diff --git a/tests/test_attached_session_goal_instance.py b/tests/test_attached_session_goal_instance.py index b13ab987d2..4a71910d64 100644 --- a/tests/test_attached_session_goal_instance.py +++ b/tests/test_attached_session_goal_instance.py @@ -4,6 +4,7 @@ import threading import time from pathlib import Path +from typing import Any import pytest @@ -32,6 +33,7 @@ AGENT_ID = "release-worker" HOST_SURFACE = "codex-app-ssh" HOST_SESSION_ID = "host-thread" +ATTACHED_MUTATION_OPERATIONS = ("submit", "enqueue", "resume", "claim", "complete") def _register(tmp_path: Path) -> tuple[Path, Path, ChatSessionStore, str]: @@ -114,6 +116,117 @@ def _recreate(registry_path: Path, instance_id: str) -> str: return str(result["goal_ref"]["goal_instance_id"]) +def _create_unstamped_attached_session(store: ChatSessionStore) -> str: + session = store.create_session( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + adapter_kind="attached_host_session", + upstream_thread_id=HOST_SESSION_ID, + session_mode="attached_host", + host_surface=HOST_SURFACE, + ) + return str(session["session_id"]) + + +def _prepare_attached_operation( + store: ChatSessionStore, + *, + session_id: str, + operation: str, +) -> dict[str, Any] | None: + if operation == "resume": + store.update_session( + session_id, + status="resume_failed", + last_error_code="fixture", + ) + return None + if operation not in {"claim", "complete"}: + return None + turn, _created = store.create_queued_turn( + session_id, + client_turn_id=f"{operation}-turn", + message="exercise attached operation", + ) + if operation == "complete": + claimed = store.claim_next_queued_turn( + session_id, + host_claim_id="fixture-claim", + ) + assert claimed is not None + return claimed + return turn + + +def _invoke_attached_operation( + *, + store: ChatSessionStore, + registry_path: Path, + session_id: str, + operation: str, + turn: dict[str, Any] | None, + work_dir: Path, +) -> Any: + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + if operation == "submit": + return runtime.submit_turn( + session_id=session_id, + client_turn_id="missing-registry-submit", + message="exercise attached operation", + work_dir=work_dir, + objective="Ship the release.", + ) + if operation == "enqueue": + return runtime.enqueue_turn( + session_id=session_id, + client_turn_id="missing-registry-enqueue", + message="exercise attached operation", + work_dir=work_dir, + objective="Ship the release.", + ) + if operation == "resume": + return runtime.resume_session( + session_id=session_id, + work_dir=work_dir, + objective="Ship the release.", + ) + if operation == "claim": + return claim_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_id, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="fixture-claim", + ) + if operation == "complete": + assert turn is not None + return complete_attached_agent_turn( + store=store, + registry_path=registry_path, + session_id=session_id, + turn_id=str(turn["turn_id"]), + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="fixture-claim", + completion_id="fixture-completion", + response={"message": "exercise attached operation"}, + ) + raise AssertionError(f"unsupported attached operation fixture: {operation}") + + +def _chat_business_state(store: ChatSessionStore) -> dict[str, bytes]: + return { + path.relative_to(store.root).as_posix(): path.read_bytes() + for path in sorted(store.root.rglob("*")) + if path.is_file() and path.suffix in {".json", ".jsonl"} + } + + def test_recreated_goal_cannot_reuse_or_claim_from_attached_session( tmp_path: Path, ) -> None: @@ -242,7 +355,7 @@ def test_stale_attached_enqueue_is_rejected_without_mutation( @pytest.mark.parametrize( "operation", - ["submit", "enqueue", "resume", "claim", "complete"], + ATTACHED_MUTATION_OPERATIONS, ) def test_strict_profile_rejects_unstamped_attached_session( tmp_path: Path, @@ -335,6 +448,83 @@ def test_strict_profile_rejects_unstamped_attached_session( assert store.turn_for_client(session_id, "strict-enqueue") is None +@pytest.mark.parametrize("operation", ATTACHED_MUTATION_OPERATIONS) +@pytest.mark.parametrize("target", ["exact", "unstamped"]) +def test_missing_registry_rejects_mutation_after_exact_session_history( + tmp_path: Path, + operation: str, + target: str, +) -> None: + registry_path, _runtime_root, store, _instance_a = _register(tmp_path) + exact_session_id = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] + session_id = ( + exact_session_id + if target == "exact" + else _create_unstamped_attached_session(store) + ) + turn = _prepare_attached_operation( + store, + session_id=session_id, + operation=operation, + ) + registry_path.unlink() + before = _chat_business_state(store) + + with pytest.raises(FileNotFoundError): + _invoke_attached_operation( + store=store, + registry_path=registry_path, + session_id=session_id, + operation=operation, + turn=turn, + work_dir=tmp_path, + ) + + assert _chat_business_state(store) == before + + +@pytest.mark.parametrize("operation", ATTACHED_MUTATION_OPERATIONS) +def test_missing_registry_without_exact_history_preserves_legacy_mutation( + tmp_path: Path, + operation: str, +) -> None: + store = ChatSessionStore(tmp_path / "runtime") + session_id = _create_unstamped_attached_session(store) + turn = _prepare_attached_operation( + store, + session_id=session_id, + operation=operation, + ) + result = _invoke_attached_operation( + store=store, + registry_path=tmp_path / "missing-registry.json", + session_id=session_id, + operation=operation, + turn=turn, + work_dir=tmp_path, + ) + + if operation in {"submit", "enqueue"}: + queued_turn, created = result + assert created is True + assert "goal_instance_id" not in queued_turn + elif operation == "resume": + assert result["status"] == "ready" + assert result["last_error_code"] is None + elif operation == "claim": + assert result["claimed"] is True + claimed_turn = store.load_turn(session_id, str(turn["turn_id"])) # type: ignore[index] + assert claimed_turn is not None + assert claimed_turn["status"] == "running" + assert "admitted_goal_instance_id" not in claimed_turn + else: + assert result["created"] is True + completed_turn = store.load_turn(session_id, str(turn["turn_id"])) # type: ignore[index] + assert completed_turn is not None + assert completed_turn["status"] == "completed" + assert "admitted_goal_instance_id" not in completed_turn + + def test_claim_wait_does_not_hold_goal_lifetime_guard(tmp_path: Path) -> None: registry_path, _runtime_root, store, instance_a = _register(tmp_path) session_a = str(_bind(store, registry_path)["session"]["session_id"]) # type: ignore[index] @@ -465,12 +655,16 @@ def test_strict_managed_open_fails_before_provider_start( assert raised.value.error_code == "source_session_managed_chat_unsupported" -def test_strict_session_does_not_fall_back_when_registry_disappears( +@pytest.mark.parametrize("latest_session", ["exact", "unstamped"]) +def test_strict_context_does_not_fall_back_when_registry_disappears( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + latest_session: str, ) -> None: registry_path, _runtime_root, store, _instance_a = _register(tmp_path) _bind(store, registry_path) + if latest_session == "unstamped": + _create_unstamped_attached_session(store) registry_path.unlink() runtime = ChatRuntimeController( store=store, From efba85051557f36844d20b87a60abaf32ddb56bf Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:45:57 +0800 Subject: [PATCH 4/4] chore(ci): refresh the registry I/O census after the main sync Five main commits landed registry-reading sites before this branch could merge, so the checked-in census named stale metadata for loopx/cli.py. Regenerated with scripts/generate_project_registry_io_manifest.py; the census and Goal instance inventory ratchets pass again. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 036ec69882..b23c94c153 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -487,7 +487,7 @@ }, { "site": "loopx/cli.py::.main::codec_read:load_project_registry#1", - "line": 837, + "line": 846, "column": 17, "kind": "codec_read", "api": "load_project_registry",