From dca405b69a93cb829849c2377848ef123bc36ada Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 01:57:05 +0800 Subject: [PATCH 1/6] fix(chat): recover interrupted turn acceptance Signed-off-by: duanjialing.777 --- apps/presentation/dashboard/package.json | 1 + .../smoke/chat-turn-acceptance-retry-smoke.ts | 155 ++++ apps/presentation/dashboard/src/data/chat.ts | 90 ++- loopx/chat_runtime.py | 176 ++++- loopx/chat_server.py | 1 + loopx/chat_store.py | 416 ++++++++-- loopx/chat_turn_acceptance.py | 307 ++++++++ .../control_plane/effect_runtime_handlers.ts | 2 + .../turn_driver/chat_turn_acceptance.ts | 731 ++++++++++++++++++ .../chat_turn_acceptance.test.ts | 367 +++++++++ tests/test_chat_session_active_turn.py | 584 +++++++++++++- tsconfig.control-plane.json | 2 + 12 files changed, 2696 insertions(+), 136 deletions(-) create mode 100644 apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts create mode 100644 loopx/chat_turn_acceptance.py create mode 100644 loopx/control_plane/turn_driver/chat_turn_acceptance.ts create mode 100644 tests/control_plane_ts/chat_turn_acceptance.test.ts diff --git a/apps/presentation/dashboard/package.json b/apps/presentation/dashboard/package.json index 4fabe78542..92a4b8d1be 100644 --- a/apps/presentation/dashboard/package.json +++ b/apps/presentation/dashboard/package.json @@ -20,6 +20,7 @@ "smoke:benchmark-study-browser": "node ../../../examples/dashboard-benchmark-study-browser-smoke.mjs", "smoke:capability-configuration": "rm -rf node_modules/.cache/loopx-capability-configuration-smoke && tsc --ignoreConfig --target ES2022 --module NodeNext --moduleResolution NodeNext --types node --skipLibCheck --strict --outDir node_modules/.cache/loopx-capability-configuration-smoke smoke/capability-configuration-smoke.ts src/data/capability-configuration.ts && node node_modules/.cache/loopx-capability-configuration-smoke/smoke/capability-configuration-smoke.js", "smoke:chat-route": "tsc --ignoreConfig --target ES2022 --module ES2022 --moduleResolution Bundler --ignoreDeprecations 6.0 --skipLibCheck --strict --outDir node_modules/.cache/loopx-chat-route-smoke smoke/chat-route-smoke.ts src/data/chat-model.ts src/vite-env.d.ts && node node_modules/.cache/loopx-chat-route-smoke/smoke/chat-route-smoke.js", + "smoke:chat-turn-acceptance-retry": "vite build --ssr smoke/chat-turn-acceptance-retry-smoke.ts --outDir node_modules/.cache/loopx-chat-turn-acceptance-retry --emptyOutDir && node node_modules/.cache/loopx-chat-turn-acceptance-retry/chat-turn-acceptance-retry-smoke.js", "smoke:demo-readiness": "bash ../../../scripts/loopx-python.sh --exec ../../../examples/dashboard-demo-readiness-smoke.py", "smoke:frontstage-browser": "node ../../../examples/dashboard-frontstage-browser-smoke.mjs", "smoke:frontstage-design-baseline": "node ../../../examples/dashboard-frontstage-design-baseline-smoke.mjs", diff --git a/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts new file mode 100644 index 0000000000..c4262ae8d3 --- /dev/null +++ b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts @@ -0,0 +1,155 @@ +import assert from "node:assert/strict"; + +import { + ChatApiError, + acceptChatTurn, +} from "../src/data/chat.ts"; + +const acceptedResponse = { + ok: true, + session_id: "session-one", + turn_id: "turn-one", + created: false, + status: "queued", + events_url: "/api/chat/sessions/session-one/turns/turn-one/events", +}; +const originalFetch = globalThis.fetch; + +try { + const bodies: string[] = []; + globalThis.fetch = async (_input, init) => { + bodies.push(String(init?.body)); + if (bodies.length === 1) { + throw new TypeError("response connection closed"); + } + return new Response(JSON.stringify(acceptedResponse), { + status: 202, + headers: {"Content-Type": "application/json"}, + }); + }; + + const accepted = await acceptChatTurn( + "session-one", + "continue", + "stable-client-turn", + ); + assert.deepEqual(accepted, acceptedResponse); + assert.equal(bodies.length, 2); + assert.deepEqual( + bodies.map((body) => JSON.parse(body).client_turn_id), + ["stable-client-turn", "stable-client-turn"], + ); + + const interruptedBodies: string[] = []; + globalThis.fetch = async (_input, init) => { + interruptedBodies.push(String(init?.body)); + if (interruptedBodies.length === 1) { + const body = new ReadableStream({ + start(controller) { + controller.error(new TypeError("response body connection closed")); + }, + }); + return new Response(body, { + status: 202, + headers: {"Content-Type": "application/json"}, + }); + } + return new Response(JSON.stringify(acceptedResponse), { + status: 202, + headers: {"Content-Type": "application/json"}, + }); + }; + + await acceptChatTurn( + "session-one", + "continue after body loss", + "stable-body-loss-turn", + ); + assert.deepEqual( + interruptedBodies.map((body) => JSON.parse(body).client_turn_id), + ["stable-body-loss-turn", "stable-body-loss-turn"], + ); + + const unavailableBodies: string[] = []; + globalThis.fetch = async (_input, init) => { + unavailableBodies.push(String(init?.body)); + if (unavailableBodies.length === 1) { + return new Response(JSON.stringify({error: "temporarily unavailable"}), { + status: 503, + headers: {"Content-Type": "application/json"}, + }); + } + return new Response(JSON.stringify(acceptedResponse), { + status: 202, + headers: {"Content-Type": "application/json"}, + }); + }; + + await acceptChatTurn( + "session-one", + "continue after unavailable", + "stable-unavailable-turn", + ); + assert.deepEqual( + unavailableBodies.map((body) => JSON.parse(body).client_turn_id), + ["stable-unavailable-turn", "stable-unavailable-turn"], + ); + + const resumeFailureBodies: string[] = []; + globalThis.fetch = async (_input, init) => { + resumeFailureBodies.push(String(init?.body)); + if (resumeFailureBodies.length === 1) { + return new Response(JSON.stringify({ + error: "adapter startup interrupted", + error_code: "resume_failed", + }), { + status: 424, + headers: {"Content-Type": "application/json"}, + }); + } + return new Response(JSON.stringify(acceptedResponse), { + status: 202, + headers: {"Content-Type": "application/json"}, + }); + }; + + await acceptChatTurn( + "session-one", + "continue after adapter recovery", + "stable-resume-turn", + ); + assert.deepEqual( + resumeFailureBodies.map((body) => JSON.parse(body).client_turn_id), + ["stable-resume-turn", "stable-resume-turn"], + ); + + let conflictCalls = 0; + globalThis.fetch = async () => { + conflictCalls += 1; + return new Response( + JSON.stringify({ + ok: false, + error: "another turn is active", + active_turn_id: "turn-two", + }), + { + status: 409, + headers: {"Content-Type": "application/json"}, + }, + ); + }; + await assert.rejects( + acceptChatTurn( + "session-one", + "different request", + "different-client-turn", + ), + (error: unknown) => ( + error instanceof ChatApiError + && error.payload.http_status === 409 + ), + ); + assert.equal(conflictCalls, 1); +} finally { + globalThis.fetch = originalFetch; +} diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 26bf36e492..76fb529c64 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -555,7 +555,18 @@ async function requestJson(url: string, init?: RequestInit): Promise { { error_code: "chat_api_unavailable" }, ); } - const responseText = await response.text(); + let responseText: string; + try { + responseText = await response.text(); + } catch { + throw new ChatApiError( + "LoopX Chat 服务响应中断。请重试当前操作。", + { + error_code: "chat_api_unavailable", + http_status: response.status, + }, + ); + } let parsedPayload: unknown = null; if (responseText.trim()) { try { @@ -577,9 +588,15 @@ async function requestJson(url: string, init?: RequestInit): Promise { const serviceMessage = response.status >= 500 ? `LoopX Chat 服务暂时不可用(HTTP ${response.status})。请确认 Dashboard 与 Chat 服务已启动且来自同一版本。` : `LoopX Chat 请求失败(HTTP ${response.status})。`; - throw new ChatApiError(staleMessage ?? String(payload.error || serviceMessage), Object.keys(payload).length - ? payload - : { error_code: "chat_api_unavailable", http_status: response.status }); + throw new ChatApiError( + staleMessage ?? String(payload.error || serviceMessage), + Object.keys(payload).length + ? { ...payload, http_status: response.status } + : { + error_code: "chat_api_unavailable", + http_status: response.status, + }, + ); } if (parsedPayload === null) { throw new ChatApiError( @@ -779,28 +796,49 @@ export async function acceptChatTurn( message: string, clientTurnId: string, attachments: ChatImageAttachmentInput[] = [], + signal?: AbortSignal, ) { - return requestJson<{ - ok: true; - session_id: string; - turn_id: string; - created: boolean; - status: string; - events_url: string; - }>(`/api/chat/sessions/${sessionId}/turns`, { - method: "POST", - body: JSON.stringify({ - message, - client_turn_id: clientTurnId, - ...(attachments.length ? { attachments: attachments.map((attachment) => ({ - data_url: attachment.dataUrl, - id: attachment.id, - mime_type: attachment.mimeType, - name: attachment.name, - size: attachment.size, - })) } : {}), - }), + const body = JSON.stringify({ + message, + client_turn_id: clientTurnId, + ...(attachments.length ? { attachments: attachments.map((attachment) => ({ + data_url: attachment.dataUrl, + id: attachment.id, + mime_type: attachment.mimeType, + name: attachment.name, + size: attachment.size, + })) } : {}), }); + for (let attempt = 0; attempt < 2; attempt += 1) { + try { + return await requestJson<{ + ok: true; + session_id: string; + turn_id: string; + created: boolean; + status: string; + events_url: string; + }>(`/api/chat/sessions/${sessionId}/turns`, { + method: "POST", + body, + signal, + }); + } catch (error) { + const status = error instanceof ChatApiError + ? Number(error.payload.http_status ?? 0) + : 0; + const retryable = error instanceof ChatApiError && ( + error.payload.error_code === "chat_api_unavailable" + || status >= 500 + || ( + status === 424 + && error.payload.error_code === "resume_failed" + ) + ); + if (attempt > 0 || signal?.aborted || !retryable) throw error; + } + } + throw new Error("unreachable Chat turn acceptance retry state"); } function parseSseBlock(block: string): ChatStreamEvent | null { @@ -1009,11 +1047,13 @@ export async function sendChatTurnStreaming( signal?: AbortSignal; } = {}, ) { + const clientTurnId = options.clientTurnId ?? crypto.randomUUID(); const accepted = await acceptChatTurn( sessionId, message, - options.clientTurnId ?? crypto.randomUUID(), + clientTurnId, options.attachments, + options.signal, ); options.onPhase?.("turn.accepted", accepted.turn_id); return receiveChatTurnStreaming( diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 60997503d1..b95b21b2a6 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -651,6 +651,7 @@ def _ensure_adapter_locked( work_dir: Path, objective: str, interrupted_turn_id: str | None = None, + accepted_turn_id: str | None = None, ) -> ChatRuntimeAdapter: session_id = str(session["session_id"]) current_session = self.store.load_session(session_id) @@ -718,9 +719,28 @@ def _ensure_adapter_locked( **manager_runtime_session_fields(manager_runtime), ) return reusable - self.store.update_session(session_id, status="resuming", last_error_code=None) - active_turn_id = interrupted_turn_id or session.get("active_turn_id") - if active_turn_id: + if accepted_turn_id is not None: + accepted_turn = self.store.load_turn(session_id, accepted_turn_id) + if ( + session.get("active_turn_id") != accepted_turn_id + or accepted_turn is None + or accepted_turn.get("status") != "queued" + ): + raise ValueError( + "accepted chat turn is no longer queued and active" + ) + else: + self.store.update_session( + session_id, + status="resuming", + last_error_code=None, + ) + active_turn_id = ( + None + if accepted_turn_id is not None + else interrupted_turn_id or session.get("active_turn_id") + ) + if active_turn_id is not None: active = self.store.load_turn(session_id, str(active_turn_id)) if active and active.get("status") not in TERMINAL_TURN_STATES: failed = self.store.update_turn( @@ -748,6 +768,10 @@ def _ensure_adapter_locked( } for item in stored_messages if item.get("role") in {"user", "agent"} + and ( + accepted_turn_id is None + or item.get("turn_id") != accepted_turn_id + ) ] legacy_manager_context = ( is_manager_channel(session.get("channel_id")) @@ -803,12 +827,20 @@ def _ensure_adapter_locked( adapter.close_session() raise except Exception as exc: - self.store.update_session( - session_id, - status="resume_failed", - active_turn_id=None, - last_error_code="resume_failed", - ) + if accepted_turn_id is not None: + self.store.update_session( + session_id, + status="busy", + active_turn_id=accepted_turn_id, + last_error_code="resume_failed", + ) + else: + self.store.update_session( + session_id, + status="resume_failed", + active_turn_id=None, + last_error_code="resume_failed", + ) gate = exc.gate if isinstance(exc, CodexChatAgentError) else None raise CodexChatAgentError( "The previous Agent conversation could not be restored.", @@ -878,38 +910,90 @@ def submit_turn( with self._session_adapter_lock(session_id): if loopx_execution and session.get("loopx_tools") is not True: session = self.loopx_mode.activate_tools(session, work_dir=work_dir, objective=objective) - adapter = self._ensure_adapter_locked( - session, + accepted = self.store.accept_managed_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + attachments=attachments, + display_message=( + ( + "开启 LoopX 模式,持续推进当前 Goal。" + if (loopx_request or {}).get("operation") == "start" + else "恢复 LoopX 模式。" + ) + if loopx_execution + else None + ), + loopx_execution=loopx_execution, + loopx_request=loopx_request, + ) + if accepted.dispatch_required: + adapter = self._ensure_adapter_locked( + session, + work_dir=work_dir, + objective=objective, + accepted_turn_id=str(accepted.turn["turn_id"]), + ) + if accepted.dispatch_reason == "already_started": + self.resume_session( + session_id=session_id, work_dir=work_dir, objective=objective, ) - turn, created = self.store.create_turn( + current = self.store.load_turn( session_id, - client_turn_id=client_turn_id, + str(accepted.turn["turn_id"]), + ) + return current or accepted.turn, accepted.created + if accepted.dispatch_required: + self._start_accepted_turn_worker( + session_id=session_id, + turn_id=str(accepted.turn["turn_id"]), message=message, - attachments=attachments, - **({"display_message": "开启 LoopX 模式,持续推进当前 Goal。" if (loopx_request or {}).get("operation") == "start" else "恢复 LoopX 模式。"} if loopx_execution else {}), + attachments=attachments or [], + adapter=adapter, + loopx_execution=loopx_execution, ) - if created and loopx_execution: - self.store.update_turn(session_id, turn["turn_id"], loopx_execution=True, loopx_request=loopx_request) - if not created: - return turn, False + return accepted.turn, accepted.created + + def _start_accepted_turn_worker( + self, + *, + session_id: str, + turn_id: str, + message: str, + attachments: list[dict[str, Any]], + adapter: ChatRuntimeAdapter, + loopx_execution: bool, + ) -> bool: + key = (session_id, turn_id) + done_event = threading.Event() + with self.lock: + if key in self.turn_done_events: + return False + self.turn_done_events[key] = done_event worker = threading.Thread( target=self._run_turn, kwargs={ "session_id": session_id, - "turn_id": str(turn["turn_id"]), + "turn_id": turn_id, "message": message, - "attachments": attachments or [], + "attachments": attachments, "adapter": adapter, "loopx_execution": loopx_execution, + "done_event": done_event, }, daemon=True, ) - with self.lock: - self.turn_done_events[(session_id, str(turn["turn_id"]))] = threading.Event() - worker.start() - return turn, True + try: + worker.start() + except Exception: + with self.lock: + if self.turn_done_events.get(key) is done_event: + self.turn_done_events.pop(key, None) + done_event.set() + raise + return True def steer_active_turn( self, @@ -1164,14 +1248,16 @@ def _drain_session_queue( if preparation_error is not None: self._fail_queue_preparation(session_id, turn_id, preparation_error) continue + done_event = threading.Event() with self.lock: - self.turn_done_events[(session_id, turn_id)] = threading.Event() + self.turn_done_events[(session_id, turn_id)] = done_event self._run_turn( session_id=session_id, turn_id=turn_id, message=str(turn.get("message") or ""), attachments=[], adapter=adapter, + done_event=done_event, ) finally: with self.lock: @@ -1208,8 +1294,16 @@ def _run_turn( message: str, attachments: list[dict[str, Any]], adapter: ChatRuntimeAdapter, + done_event: threading.Event | None = None, loopx_execution: bool = False, ) -> None: + key = (session_id, turn_id) + if done_event is None: + with self.lock: + done_event = self.turn_done_events.get(key) + if done_event is None: + done_event = threading.Event() + self.turn_done_events[key] = done_event started = utc_now() started_turn = self.store.update_turn( session_id, @@ -1221,9 +1315,9 @@ def _run_turn( if started_turn is None: with self.lock: self.cancelled_turns.discard((session_id, turn_id)) - done_event = self.turn_done_events.pop((session_id, turn_id), None) - if done_event is not None: - done_event.set() + if self.turn_done_events.get(key) is done_event: + self.turn_done_events.pop(key, None) + done_event.set() return event_buffer = _TurnEventBuffer( store=self.store, @@ -1409,10 +1503,10 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None: finally: event_buffer.close() with self.lock: - self.turn_event_buffers.pop((session_id, turn_id), None) - done_event = self.turn_done_events.pop((session_id, turn_id), None) - if done_event is not None: - done_event.set() + self.turn_event_buffers.pop(key, None) + if self.turn_done_events.get(key) is done_event: + self.turn_done_events.pop(key, None) + done_event.set() def _fail_turn( self, @@ -1584,6 +1678,22 @@ def close_session(self, session_id: str) -> bool: return True def resume_session(self, *, session_id: str, work_dir: Path, objective: str) -> dict[str, Any]: + prepared = self.store.prepared_managed_turn_request(session_id) + if prepared is not None: + self.submit_turn( + session_id=session_id, + client_turn_id=str(prepared["client_turn_id"]), + message=str(prepared["message"]), + attachments=prepared.get("attachments"), + work_dir=work_dir, + objective=objective, + loopx_execution=bool(prepared["loopx_execution"]), + loopx_request=prepared.get("loopx_request"), + ) + restored = self.store.load_session(session_id) + if restored is None: + raise KeyError("chat session was not found") + return restored with self._session_adapter_lock(session_id): session = self.store.load_session(session_id) if session is None or session.get("status") == "closed": diff --git a/loopx/chat_server.py b/loopx/chat_server.py index fab3fc546d..9e86efa7c0 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -721,6 +721,7 @@ def _session_turn(self, session_id: str) -> None: str(exc), status=424, gate=exc.gate, + error_code=exc.error_code, ) return except RuntimeError as exc: diff --git a/loopx/chat_store.py b/loopx/chat_store.py index a182ed84c5..7d44233e10 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -19,6 +19,11 @@ ) from .chat_event_cache import ChatEventCache from .chat_ingress import ChatIngressStore +from .chat_turn_acceptance import ( + CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA, + AcceptedManagedTurn, + plan_managed_turn_acceptance, +) from .file_lock import exclusive_file_lock @@ -82,7 +87,7 @@ def _read_jsonl(path: Path) -> list[dict[str, Any]]: for line in lines: try: item = json.loads(line) - except json.JSONDecodeError: + except (UnicodeDecodeError, json.JSONDecodeError): continue if isinstance(item, dict): rows.append(item) @@ -113,10 +118,49 @@ def _append_jsonl(path: Path, payload: dict[str, Any]) -> None: _append_jsonl_rows(path, [payload]) +def _repair_incomplete_jsonl_tail(path: Path) -> None: + try: + handle = path.open("r+b") + except FileNotFoundError: + return + with handle: + end = handle.seek(0, os.SEEK_END) + if end == 0: + return + handle.seek(end - 1) + if handle.read(1) == b"\n": + return + cursor = end + truncate_at = 0 + while cursor > 0: + start = max(0, cursor - 8192) + handle.seek(start) + chunk = handle.read(cursor - start) + if (newline := chunk.rfind(b"\n")) >= 0: + truncate_at = start + newline + 1 + break + cursor = start + handle.seek(truncate_at) + tail = handle.read(end - truncate_at) + try: + item = json.loads(tail) + except (UnicodeDecodeError, json.JSONDecodeError): + handle.truncate(truncate_at) + else: + if isinstance(item, dict): + handle.seek(end) + handle.write(b"\n") + else: + handle.truncate(truncate_at) + handle.flush() + os.fsync(handle.fileno()) + + def _append_jsonl_rows(path: Path, rows: list[dict[str, Any]]) -> None: if not rows: return path.parent.mkdir(parents=True, exist_ok=True, mode=0o700) + _repair_incomplete_jsonl_tail(path) with path.open("a", encoding="utf-8") as handle: for payload in rows: handle.write(json.dumps(payload, ensure_ascii=False, separators=(",", ":"))) @@ -579,6 +623,89 @@ def append_message( def messages(self, session_id: str) -> list[dict[str, Any]]: return _read_jsonl(self._session_dir(session_id) / "messages.jsonl") + def prepared_managed_turn_request( + self, + session_id: str, + ) -> dict[str, Any] | None: + session_token = _opaque_id(session_id, field="session_id") + session_path = self._session_path(session_token) + with self._session_lock(session_token): + with exclusive_file_lock( + session_path, + agent_id="loopx-chat", + operation="read_prepared_managed_chat_turn", + ): + prepared = [ + payload + for path in ( + self._session_dir(session_token) / "turns" + ).glob("*.json") + if ( + payload := _read_json(path) + ).get("schema_version") == CHAT_TURN_SCHEMA_VERSION + and "_acceptance" in payload + ] + if not prepared: + return None + if len(prepared) != 1: + raise ValueError( + "chat turn acceptance state is inconsistent" + ) + turn = prepared[0] + acceptance = turn.get("_acceptance") + if ( + not isinstance(acceptance, dict) + or acceptance.get("schema_version") + != CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA + or acceptance.get("phase") != "prepared" + or not isinstance(turn.get("message"), str) + or not isinstance( + acceptance.get("display_message"), + str, + ) + ): + raise ValueError( + "chat turn acceptance state is inconsistent" + ) + attachments = acceptance.get("attachments") + if attachments is not None and ( + not isinstance(attachments, list) + or any( + not isinstance(attachment, dict) + for attachment in attachments + ) + ): + raise ValueError( + "chat turn acceptance state is inconsistent" + ) + loopx_execution = turn.get("loopx_execution", False) + loopx_request = turn.get("loopx_request") + if ( + not isinstance(loopx_execution, bool) + or ( + loopx_request is not None + and not isinstance(loopx_request, dict) + ) + ): + raise ValueError( + "chat turn acceptance state is inconsistent" + ) + return { + "client_turn_id": _opaque_id( + turn.get("client_turn_id"), + field="client_turn_id", + ), + "message": turn["message"], + "attachments": attachments, + "origin": _opaque_id( + turn.get("origin"), + field="origin", + ), + "display_message": acceptance["display_message"], + "loopx_execution": loopx_execution, + "loopx_request": loopx_request, + } + def create_turn( self, session_id: str, @@ -589,90 +716,221 @@ def create_turn( origin: str = "web", display_message: str | None = None, ) -> tuple[dict[str, Any], bool]: + accepted = self.accept_managed_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + attachments=attachments, + origin=origin, + display_message=display_message, + ) + return accepted.turn, accepted.created + + def accept_managed_turn( + self, + session_id: str, + *, + client_turn_id: str, + message: str, + attachments: list[dict[str, Any]] | None = None, + origin: str = "web", + display_message: str | None = None, + loopx_execution: bool = False, + loopx_request: dict[str, object] | None = None, + ) -> AcceptedManagedTurn: + session_token = _opaque_id(session_id, field="session_id") client_id = _opaque_id(client_turn_id, field="client_turn_id") - session_path = self._session_path(session_id) - with self._session_lock(session_id): + normalized_origin = _opaque_id(origin, field="origin") + execution_message = str(message) + visible_message = ( + str(display_message) + if display_message is not None + else execution_message + ) + normalized_attachments = attachments or None + candidate_turn_id = uuid.uuid4().hex + candidate_message_id = uuid.uuid4().hex + accepted_at = utc_now() + session_path = self._session_path(session_token) + with self._session_lock(session_token): with exclusive_file_lock( session_path, agent_id="loopx-chat", - operation="create_chat_turn", + operation="accept_managed_chat_turn", ): - existing = self.turn_for_client(session_id, client_id) - if existing is not None: - # Attachments live in the transcript, not the Turn record. - # This also supports pre-existing Turns without a migration. - original = next( - (row for row in self.messages(session_id) - if row.get("role") == "user" - and row.get("turn_id") == existing["turn_id"]), - None, + turns = [ + payload + for path in ( + self._session_dir(session_token) / "turns" + ).glob("*.json") + if ( + payload := _read_json(path) + ).get("schema_version") == CHAT_TURN_SCHEMA_VERSION + ] + existing = next( + ( + turn + for turn in turns + if turn.get("client_turn_id") == client_id + ), + None, + ) + session = self.load_session(session_token) + active_turn = None + if session is not None and session.get("active_turn_id"): + active_turn = self.load_turn( + session_token, + str(session["active_turn_id"]), ) - if original is None: - raise ValueError("client_turn_id original request is unavailable") - require_matching_replay( - {**existing, "attachments": original.get("attachments") or None}, - identity="client_turn_id", - request={ - "message": str(message), - "origin": _opaque_id(origin, field="origin"), - "attachments": attachments or None, + existing_turn_id = ( + str(existing.get("turn_id")) + if existing is not None + else None + ) + if existing_turn_id is not None: + self.flush_events(session_token, existing_turn_id) + plan = plan_managed_turn_acceptance( + session_id=session_token, + client_turn_id=client_id, + message=execution_message, + display_message=visible_message, + attachments=normalized_attachments, + origin=normalized_origin, + loopx_execution=loopx_execution, + loopx_request=loopx_request, + candidate_turn_id=candidate_turn_id, + candidate_message_id=candidate_message_id, + accepted_at=accepted_at, + session=session, + active_turn=active_turn, + matching_turn=existing, + turns=turns, + messages=self.messages(session_token), + queued_events=( + _read_jsonl( + self._event_path( + session_token, + existing_turn_id, + ) + ) + if existing_turn_id is not None + else [] + ), + ) + if plan.writes.prepare_turn: + if plan.turn_id != candidate_turn_id: + raise ValueError( + "chat turn acceptance selected an invalid candidate" + ) + turn_id = candidate_turn_id + payload = { + "schema_version": CHAT_TURN_SCHEMA_VERSION, + "turn_id": turn_id, + "session_id": session_token, + "client_turn_id": client_id, + "status": "queued", + "message": execution_message, + "origin": normalized_origin, + "upstream_turn_id": None, + "response": None, + "error_code": None, + "error": None, + "created_at": accepted_at, + "started_at": None, + "first_event_at": None, + "completed_at": None, + "last_activity_at": accepted_at, + "delta_count": 0, + "sse_reconnect_count": 0, + **( + { + "loopx_execution": True, + "loopx_request": loopx_request, + } + if loopx_execution + else {} + ), + "_acceptance": { + "schema_version": ( + CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA + ), + "phase": "prepared", + "request_sha256": plan.request_sha256, + "message_id": plan.message_id, + "display_message": visible_message, + "attachments": normalized_attachments, }, - ) - return existing, False - session = self.load_session(session_id) - if session is None or session.get("status") == "closed": - raise KeyError("chat session was not found") - active_turn_id = session.get("active_turn_id") - if active_turn_id: - active = self.load_turn(session_id, str(active_turn_id)) - if active and active.get("status") in { - "queued", "starting", "running", "completing", "interrupting" - }: - raise RuntimeError(str(active_turn_id)) - now = utc_now() - turn_id = uuid.uuid4().hex - payload = { - "schema_version": CHAT_TURN_SCHEMA_VERSION, - "turn_id": turn_id, - "session_id": session_id, - "client_turn_id": client_id, - "status": "queued", - "message": str(message), - "origin": _opaque_id(origin, field="origin"), - "upstream_turn_id": None, - "response": None, - "error_code": None, - "error": None, - "created_at": now, - "started_at": None, - "first_event_at": None, - "completed_at": None, - "last_activity_at": now, - "delta_count": 0, - "sse_reconnect_count": 0, - } - path = self._turn_path(session_id, turn_id) - _atomic_write_json(path, payload) - os.chmod(path, 0o600) - session.update( - { - "status": "busy", - "active_turn_id": turn_id, - "last_activity_at": now, - "updated_at": utc_now(), } + path = self._turn_path(session_token, turn_id) + _atomic_write_json(path, payload) + os.chmod(path, 0o600) + else: + if existing is None or plan.turn_id != existing_turn_id: + raise ValueError( + "chat turn acceptance selected an invalid replay" + ) + payload = existing + + if plan.writes.activate_session: + if session is None: + raise KeyError("chat session was not found") + session.update( + { + "status": "busy", + "active_turn_id": plan.turn_id, + "last_activity_at": accepted_at, + "updated_at": utc_now(), + } + ) + _atomic_write_json( + session_path, + session, + preserve_mode=True, + ) + + if plan.writes.append_message: + self.append_message( + session_token, + role="user", + text=visible_message, + turn_id=plan.turn_id, + attachments=normalized_attachments, + origin=normalized_origin, + message_id=plan.message_id, + ) + if plan.writes.append_queued_event: + self.append_event( + session_token, + plan.turn_id, + kind="turn.queued", + payload={}, + ) + if plan.writes.settle_turn: + settled = self.load_turn(session_token, plan.turn_id) + if settled is None: + raise ValueError( + "chat turn acceptance lost its prepared turn" + ) + settled.pop("_acceptance", None) + _atomic_write_json( + self._turn_path(session_token, plan.turn_id), + settled, + preserve_mode=True, + ) + payload = settled + else: + current = self.load_turn(session_token, plan.turn_id) + if current is None: + raise ValueError( + "chat turn acceptance lost its replayed turn" + ) + payload = current + return AcceptedManagedTurn( + turn=payload, + created=plan.created, + dispatch_required=plan.dispatch_required, + dispatch_reason=plan.dispatch_reason, ) - _atomic_write_json(session_path, session, preserve_mode=True) - self.append_message( - session_id, - role="user", - text=display_message if display_message is not None else message, - turn_id=turn_id, - attachments=attachments, - origin=origin, - ) - self.append_event(session_id, turn_id, kind="turn.queued", payload={}) - return payload, True def create_queued_turn( self, @@ -1436,6 +1694,12 @@ def session_snapshot(self, session_id: str) -> dict[str, Any]: active_turn = None if payload.get("active_turn_id"): active_turn = self.load_turn(session_id, str(payload["active_turn_id"])) + if active_turn is not None: + active_turn = { + key: value + for key, value in active_turn.items() + if key != "_acceptance" + } return { "ok": True, "schema_version": CHAT_STORE_SCHEMA_VERSION, diff --git a/loopx/chat_turn_acceptance.py b/loopx/chat_turn_acceptance.py new file mode 100644 index 0000000000..7f4187a1d2 --- /dev/null +++ b/loopx/chat_turn_acceptance.py @@ -0,0 +1,307 @@ +"""TypeScript-owned admission planning for managed Chat turns.""" + +from __future__ import annotations + +from dataclasses import dataclass +import hashlib +import json +from typing import Any + +from .control_plane.effect_runtime import effect_runtime_result + + +CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA = "loopx_chat_turn_acceptance_request_v0" +CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA = "loopx_chat_turn_acceptance_result_v0" +CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA = "loopx_chat_turn_acceptance_v0" + + +def _sha256(value: str) -> str: + return f"sha256:{hashlib.sha256(value.encode('utf-8')).hexdigest()}" + + +def _text_digest(value: Any) -> str: + if not isinstance(value, str): + return "invalid" + return _sha256(value) + + +def _json_digest(value: Any) -> str: + try: + encoded = json.dumps( + value, + ensure_ascii=False, + allow_nan=False, + separators=(",", ":"), + sort_keys=True, + ) + except (TypeError, ValueError): + return "invalid" + return _sha256(encoded) + + +def _acceptance_facts(turn: dict[str, Any]) -> dict[str, Any] | None: + acceptance = turn.get("_acceptance") + if acceptance is None: + return None + if not isinstance(acceptance, dict): + return {"schema_version": None} + return { + "schema_version": acceptance.get("schema_version"), + "phase": acceptance.get("phase"), + "request_sha256": acceptance.get("request_sha256"), + "message_id": acceptance.get("message_id"), + "display_message_sha256": _text_digest( + acceptance.get("display_message") + ), + "attachments_sha256": _json_digest( + acceptance.get("attachments") or None + ), + } + + +def _turn_facts(turn: dict[str, Any] | None) -> dict[str, Any] | None: + if turn is None: + return None + return { + "turn_id": turn.get("turn_id"), + "client_turn_id": turn.get("client_turn_id"), + "status": turn.get("status"), + "execution_message_sha256": _text_digest(turn.get("message")), + "origin": turn.get("origin"), + "loopx_execution": turn.get("loopx_execution", False), + "loopx_request_sha256": _json_digest(turn.get("loopx_request")), + "acceptance": _acceptance_facts(turn), + } + + +def _transcript_facts( + messages: list[dict[str, Any]], + turn_id: str | None, +) -> dict[str, Any]: + if turn_id is None: + return {"count": 0} + matching = [ + row + for row in messages + if row.get("role") == "user" and row.get("turn_id") == turn_id + ] + if len(matching) != 1: + return {"count": len(matching)} + message = matching[0] + return { + "count": 1, + "message_id": message.get("message_id"), + "display_message_sha256": _text_digest(message.get("text")), + "attachments_sha256": _json_digest( + message.get("attachments") or None + ), + "origin": message.get("origin"), + } + + +def _queued_event_facts(events: list[dict[str, Any]]) -> dict[str, Any]: + queued = [event for event in events if event.get("kind") == "turn.queued"] + if len(queued) != 1: + return {"count": len(queued)} + return { + "count": 1, + "payload_sha256": _json_digest(queued[0].get("payload")), + } + + +def _prepared_turn_facts( + turns: list[dict[str, Any]], +) -> dict[str, Any]: + prepared = [turn for turn in turns if "_acceptance" in turn] + if len(prepared) != 1: + return {"count": len(prepared)} + return { + "count": 1, + "turn_id": prepared[0].get("turn_id"), + "client_turn_id": prepared[0].get("client_turn_id"), + } + + +@dataclass(frozen=True) +class ManagedTurnAcceptanceWrites: + prepare_turn: bool + activate_session: bool + append_message: bool + append_queued_event: bool + settle_turn: bool + + +@dataclass(frozen=True) +class ManagedTurnAcceptancePlan: + disposition: str + turn_id: str + message_id: str + request_sha256: str + created: bool + writes: ManagedTurnAcceptanceWrites + dispatch_required: bool + dispatch_reason: str | None + + +@dataclass(frozen=True) +class AcceptedManagedTurn: + turn: dict[str, Any] + created: bool + dispatch_required: bool + dispatch_reason: str | None + + +def _required_string(payload: dict[str, Any], field: str) -> str: + value = payload.get(field) + if not isinstance(value, str) or not value: + raise ValueError(f"chat turn acceptance result has invalid {field}") + return value + + +def _required_bool(payload: dict[str, Any], field: str) -> bool: + value = payload.get(field) + if not isinstance(value, bool): + raise ValueError(f"chat turn acceptance result has invalid {field}") + return value + + +def _raise_rejection(result: dict[str, Any]) -> None: + code = result.get("code") + if code in {"session_not_found", "session_closed"}: + raise KeyError("chat session was not found") + if code == "request_conflict": + raise ValueError( + "client_turn_id already belongs to a different request" + ) + if code == "active_turn_conflict": + active_turn_id = result.get("active_turn_id") + raise RuntimeError( + str(active_turn_id or "chat session already has an active turn") + ) + if code == "original_request_unavailable": + raise ValueError("client_turn_id original request is unavailable") + if code == "durable_state_conflict": + raise ValueError("chat turn acceptance state is inconsistent") + raise ValueError("chat turn acceptance returned an unsupported rejection") + + +def _decode_plan(value: Any) -> ManagedTurnAcceptancePlan: + if not isinstance(value, dict): + raise ValueError("chat turn acceptance result must be an object") + if value.get("schema_version") != CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA: + raise ValueError("chat turn acceptance result schema is unsupported") + if value.get("kind") == "rejected": + _raise_rejection(value) + if value.get("kind") != "accepted": + raise ValueError("chat turn acceptance result kind is unsupported") + disposition = _required_string(value, "disposition") + if disposition not in {"created", "repaired", "replayed"}: + raise ValueError("chat turn acceptance disposition is unsupported") + writes = value.get("writes") + dispatch = value.get("dispatch") + if not isinstance(writes, dict) or not isinstance(dispatch, dict): + raise ValueError("chat turn acceptance result is incomplete") + dispatch_kind = dispatch.get("kind") + if dispatch_kind not in {"required", "not_required"}: + raise ValueError("chat turn acceptance dispatch kind is unsupported") + dispatch_reason = dispatch.get("reason") + if dispatch_kind == "not_required": + if dispatch_reason not in { + "already_started", + "completion_in_progress", + "terminal", + }: + raise ValueError("chat turn acceptance dispatch reason is unsupported") + elif dispatch_reason is not None: + raise ValueError("chat turn acceptance dispatch reason is unsupported") + return ManagedTurnAcceptancePlan( + disposition=disposition, + turn_id=_required_string(value, "turn_id"), + message_id=_required_string(value, "message_id"), + request_sha256=_required_string(value, "request_sha256"), + created=_required_bool(value, "created"), + writes=ManagedTurnAcceptanceWrites( + prepare_turn=_required_bool(writes, "prepare_turn"), + activate_session=_required_bool(writes, "activate_session"), + append_message=_required_bool(writes, "append_message"), + append_queued_event=_required_bool( + writes, + "append_queued_event", + ), + settle_turn=_required_bool(writes, "settle_turn"), + ), + dispatch_required=dispatch_kind == "required", + dispatch_reason=( + str(dispatch_reason) + if dispatch_kind == "not_required" + else None + ), + ) + + +def plan_managed_turn_acceptance( + *, + session_id: str, + client_turn_id: str, + message: str, + display_message: str, + attachments: list[dict[str, Any]] | None, + origin: str, + loopx_execution: bool, + loopx_request: dict[str, object] | None, + candidate_turn_id: str, + candidate_message_id: str, + accepted_at: str, + session: dict[str, Any] | None, + active_turn: dict[str, Any] | None, + matching_turn: dict[str, Any] | None, + turns: list[dict[str, Any]], + messages: list[dict[str, Any]], + queued_events: list[dict[str, Any]], +) -> ManagedTurnAcceptancePlan: + request = { + "schema_version": CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA, + "request": { + "session_id": session_id, + "client_turn_id": client_turn_id, + "execution_message_sha256": _text_digest(message), + "display_message_sha256": _text_digest(display_message), + "attachments_sha256": _json_digest(attachments or None), + "origin": origin, + "loopx_execution": loopx_execution, + "loopx_request_sha256": _json_digest(loopx_request), + }, + "candidate": { + "turn_id": candidate_turn_id, + "message_id": candidate_message_id, + "accepted_at": accepted_at, + }, + "session": ( + { + "status": session.get("status"), + "active_turn_id": session.get("active_turn_id"), + } + if session is not None + else None + ), + "active_turn": ( + { + "turn_id": active_turn.get("turn_id"), + "status": active_turn.get("status"), + } + if active_turn is not None + else None + ), + "matching_turn": _turn_facts(matching_turn), + "transcript": _transcript_facts( + messages, + ( + str(matching_turn.get("turn_id")) + if matching_turn is not None + else None + ), + ), + "queued_event": _queued_event_facts(queued_events), + "prepared_turn": _prepared_turn_facts(turns), + } + return _decode_plan(effect_runtime_result("chat.turn.accept", request)) diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 4ecc815623..5002d0c379 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -17,6 +17,7 @@ import {recordDelegationAdoption, delegationInventoryItem, delegationInventoryQu import {resolveConversationTrigger} from "./collaboration/conversation_trigger.ts"; import {planChatMode} from "./collaboration/chat_mode.ts"; import {resolveConversationScope} from "./collaboration/conversation_scope.ts"; +import {planChatTurnAcceptance} from "./turn_driver/chat_turn_acceptance.ts"; import {previewTeamPlan, planTeamTransaction, teamTransactionIdentity} from "./work_items/team_plan.ts"; import {commitLocalTeamPlan} from "./work_items/team_plan_authority.ts"; import {inspectLocalGoalAcceptance, commitLocalGoalAcceptance, @@ -724,6 +725,7 @@ export function createEffectRuntimeHandlers( ["collaboration.chat_mode", planChatMode], ["collaboration.conversation.trigger", resolveConversationTrigger], ["collaboration.conversation.scope", resolveConversationScope], + ["chat.turn.accept", planChatTurnAcceptance], ["collaboration.delegation.observe", transitionDelegationObservation], ["collaboration.delegation.recover_validated_settlement", recoverValidatedDelegationSettlement], ["collaboration.delegation.adoption", recordDelegationAdoption], diff --git a/loopx/control_plane/turn_driver/chat_turn_acceptance.ts b/loopx/control_plane/turn_driver/chat_turn_acceptance.ts new file mode 100644 index 0000000000..879a105baa --- /dev/null +++ b/loopx/control_plane/turn_driver/chat_turn_acceptance.ts @@ -0,0 +1,731 @@ +import {createHash} from "node:crypto"; + +import type {JsonObject} from "../effect_program.ts"; +import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; +import { + requireBoolean, + requireInteger, + requireJsonObject, + requireNonEmptyString, + requireStringLiteral, +} from "../runtime_decode.ts"; + +export const CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA = + "loopx_chat_turn_acceptance_request_v0"; +export const CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA = + "loopx_chat_turn_acceptance_result_v0"; +export const CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA = + "loopx_chat_turn_acceptance_v0"; + +const TURN_STATUSES = [ + "queued", + "starting", + "running", + "completing", + "interrupting", + "completed", + "interrupted", + "timed_out", + "failed", +] as const; +const TERMINAL_TURN_STATUSES = new Set([ + "completed", + "interrupted", + "timed_out", + "failed", +]); +const OPAQUE_ID = /^[A-Za-z0-9._-]{1,160}$/; +const SHA256 = /^sha256:[0-9a-f]{64}$/; +const EMPTY_OBJECT_SHA256 = sha256("{}"); + +type OpaqueId = string & {readonly __brand: "OpaqueId"}; +type Sha256 = string & {readonly __brand: "Sha256"}; +type TurnStatus = (typeof TURN_STATUSES)[number]; + +interface RequestFacts { + readonly sessionId: OpaqueId; + readonly clientTurnId: OpaqueId; + readonly executionMessageSha256: Sha256; + readonly displayMessageSha256: Sha256; + readonly attachmentsSha256: Sha256; + readonly origin: OpaqueId; + readonly loopxExecution: boolean; + readonly loopxRequestSha256: Sha256; +} + +interface CandidateFacts { + readonly turnId: OpaqueId; + readonly messageId: OpaqueId; + readonly acceptedAt: string; +} + +interface SessionFacts { + readonly status: string; + readonly activeTurnId: OpaqueId | null; +} + +interface ActiveTurnFacts { + readonly turnId: OpaqueId; + readonly status: TurnStatus; +} + +interface AcceptanceCapsuleFacts { + readonly requestSha256: Sha256; + readonly messageId: OpaqueId; + readonly displayMessageSha256: Sha256; + readonly attachmentsSha256: Sha256; +} + +interface MatchingTurnFacts { + readonly turnId: OpaqueId; + readonly clientTurnId: OpaqueId; + readonly status: TurnStatus; + readonly executionMessageSha256: Sha256; + readonly origin: OpaqueId; + readonly loopxExecution: boolean; + readonly loopxRequestSha256: Sha256; + readonly acceptance: AcceptanceCapsuleFacts | null; +} + +type TranscriptFacts = + | {readonly kind: "absent"} + | { + readonly kind: "single"; + readonly messageId: OpaqueId; + readonly displayMessageSha256: Sha256; + readonly attachmentsSha256: Sha256; + readonly origin: OpaqueId; + } + | {readonly kind: "ambiguous"; readonly count: number}; + +type QueuedEventFacts = + | {readonly kind: "absent"} + | {readonly kind: "single"; readonly payloadSha256: Sha256} + | {readonly kind: "ambiguous"; readonly count: number}; + +type PreparedTurnFacts = + | {readonly kind: "absent"} + | { + readonly kind: "single"; + readonly turnId: OpaqueId; + readonly clientTurnId: OpaqueId; + } + | {readonly kind: "ambiguous"; readonly count: number}; + +interface AcceptanceFacts { + readonly request: RequestFacts; + readonly candidate: CandidateFacts; + readonly session: SessionFacts | null; + readonly activeTurn: ActiveTurnFacts | null; + readonly matchingTurn: MatchingTurnFacts | null; + readonly transcript: TranscriptFacts; + readonly queuedEvent: QueuedEventFacts; + readonly preparedTurn: PreparedTurnFacts; +} + +export type ChatTurnAcceptanceRejectionCode = + | "session_not_found" + | "session_closed" + | "request_conflict" + | "active_turn_conflict" + | "original_request_unavailable" + | "durable_state_conflict"; + +type ChatTurnDispatch = + | {readonly kind: "required"} + | { + readonly kind: "not_required"; + readonly reason: + | "already_started" + | "completion_in_progress" + | "terminal"; + }; + +interface AcceptanceWrites { + readonly prepare_turn: boolean; + readonly activate_session: boolean; + readonly append_message: boolean; + readonly append_queued_event: boolean; + readonly settle_turn: boolean; +} + +export type ChatTurnAcceptancePlan = + | { + readonly schema_version: typeof CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA; + readonly kind: "accepted"; + readonly disposition: "created" | "repaired" | "replayed"; + readonly turn_id: string; + readonly message_id: string; + readonly request_sha256: string; + readonly created: boolean; + readonly writes: AcceptanceWrites; + readonly dispatch: ChatTurnDispatch; + } + | { + readonly schema_version: typeof CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA; + readonly kind: "rejected"; + readonly code: ChatTurnAcceptanceRejectionCode; + readonly active_turn_id?: string; + }; + +function sha256(value: string): Sha256 { + return `sha256:${createHash("sha256").update(value).digest("hex")}` as Sha256; +} + +function opaqueId(value: unknown, label: string): OpaqueId { + const candidate = requireNonEmptyString(value, label); + if (!OPAQUE_ID.test(candidate)) { + throw new EffectRuntimeRequestError(`${label} must be a compact opaque id`); + } + return candidate as OpaqueId; +} + +function digest(value: unknown, label: string): Sha256 { + const candidate = requireNonEmptyString(value, label); + if (!SHA256.test(candidate)) { + throw new EffectRuntimeRequestError(`${label} must be a SHA-256 digest`); + } + return candidate as Sha256; +} + +function optionalOpaqueId(value: unknown, label: string): OpaqueId | null { + if (value === null || value === undefined) return null; + return opaqueId(value, label); +} + +function turnStatus(value: unknown, label: string): TurnStatus { + return requireStringLiteral(value, TURN_STATUSES, label); +} + +function decodeRequest(value: unknown): RequestFacts { + const request = requireJsonObject(value, "request"); + return { + sessionId: opaqueId(request.session_id, "request.session_id"), + clientTurnId: opaqueId( + request.client_turn_id, + "request.client_turn_id", + ), + executionMessageSha256: digest( + request.execution_message_sha256, + "request.execution_message_sha256", + ), + displayMessageSha256: digest( + request.display_message_sha256, + "request.display_message_sha256", + ), + attachmentsSha256: digest( + request.attachments_sha256, + "request.attachments_sha256", + ), + origin: opaqueId(request.origin, "request.origin"), + loopxExecution: requireBoolean( + request.loopx_execution, + "request.loopx_execution", + ), + loopxRequestSha256: digest( + request.loopx_request_sha256, + "request.loopx_request_sha256", + ), + }; +} + +function decodeCandidate(value: unknown): CandidateFacts { + const candidate = requireJsonObject(value, "candidate"); + return { + turnId: opaqueId(candidate.turn_id, "candidate.turn_id"), + messageId: opaqueId(candidate.message_id, "candidate.message_id"), + acceptedAt: requireNonEmptyString( + candidate.accepted_at, + "candidate.accepted_at", + ), + }; +} + +function decodeSession(value: unknown): SessionFacts | null { + if (value === null || value === undefined) return null; + const session = requireJsonObject(value, "session"); + return { + status: requireNonEmptyString(session.status, "session.status"), + activeTurnId: optionalOpaqueId( + session.active_turn_id, + "session.active_turn_id", + ), + }; +} + +function decodeActiveTurn(value: unknown): ActiveTurnFacts | null { + if (value === null || value === undefined) return null; + const turn = requireJsonObject(value, "active_turn"); + return { + turnId: opaqueId(turn.turn_id, "active_turn.turn_id"), + status: turnStatus(turn.status, "active_turn.status"), + }; +} + +function decodeAcceptance(value: unknown): AcceptanceCapsuleFacts | null { + if (value === null || value === undefined) return null; + const acceptance = requireJsonObject(value, "matching_turn.acceptance"); + requireStringLiteral( + acceptance.schema_version, + [CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA] as const, + "matching_turn.acceptance.schema_version", + ); + requireStringLiteral( + acceptance.phase, + ["prepared"] as const, + "matching_turn.acceptance.phase", + ); + return { + requestSha256: digest( + acceptance.request_sha256, + "matching_turn.acceptance.request_sha256", + ), + messageId: opaqueId( + acceptance.message_id, + "matching_turn.acceptance.message_id", + ), + displayMessageSha256: digest( + acceptance.display_message_sha256, + "matching_turn.acceptance.display_message_sha256", + ), + attachmentsSha256: digest( + acceptance.attachments_sha256, + "matching_turn.acceptance.attachments_sha256", + ), + }; +} + +function decodeMatchingTurn(value: unknown): MatchingTurnFacts | null { + if (value === null || value === undefined) return null; + const turn = requireJsonObject(value, "matching_turn"); + return { + turnId: opaqueId(turn.turn_id, "matching_turn.turn_id"), + clientTurnId: opaqueId( + turn.client_turn_id, + "matching_turn.client_turn_id", + ), + status: turnStatus(turn.status, "matching_turn.status"), + executionMessageSha256: digest( + turn.execution_message_sha256, + "matching_turn.execution_message_sha256", + ), + origin: opaqueId(turn.origin, "matching_turn.origin"), + loopxExecution: requireBoolean( + turn.loopx_execution, + "matching_turn.loopx_execution", + ), + loopxRequestSha256: digest( + turn.loopx_request_sha256, + "matching_turn.loopx_request_sha256", + ), + acceptance: decodeAcceptance(turn.acceptance), + }; +} + +function observedCount(value: JsonObject, label: string): number { + const count = requireInteger(value.count, `${label}.count`); + if (count < 0) { + throw new EffectRuntimeRequestError(`${label}.count must not be negative`); + } + return count; +} + +function decodeTranscript(value: unknown): TranscriptFacts { + const transcript = requireJsonObject(value, "transcript"); + const count = observedCount(transcript, "transcript"); + if (count === 0) return {kind: "absent"}; + if (count > 1) return {kind: "ambiguous", count}; + return { + kind: "single", + messageId: opaqueId(transcript.message_id, "transcript.message_id"), + displayMessageSha256: digest( + transcript.display_message_sha256, + "transcript.display_message_sha256", + ), + attachmentsSha256: digest( + transcript.attachments_sha256, + "transcript.attachments_sha256", + ), + origin: opaqueId(transcript.origin, "transcript.origin"), + }; +} + +function decodeQueuedEvent(value: unknown): QueuedEventFacts { + const event = requireJsonObject(value, "queued_event"); + const count = observedCount(event, "queued_event"); + if (count === 0) return {kind: "absent"}; + if (count > 1) return {kind: "ambiguous", count}; + return { + kind: "single", + payloadSha256: digest( + event.payload_sha256, + "queued_event.payload_sha256", + ), + }; +} + +function decodePreparedTurn(value: unknown): PreparedTurnFacts { + const prepared = requireJsonObject(value, "prepared_turn"); + const count = observedCount(prepared, "prepared_turn"); + if (count === 0) return {kind: "absent"}; + if (count > 1) return {kind: "ambiguous", count}; + return { + kind: "single", + turnId: opaqueId(prepared.turn_id, "prepared_turn.turn_id"), + clientTurnId: opaqueId( + prepared.client_turn_id, + "prepared_turn.client_turn_id", + ), + }; +} + +function decodeFacts(input: JsonObject): AcceptanceFacts { + requireStringLiteral( + input.schema_version, + [CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA] as const, + "schema_version", + ); + return { + request: decodeRequest(input.request), + candidate: decodeCandidate(input.candidate), + session: decodeSession(input.session), + activeTurn: decodeActiveTurn(input.active_turn), + matchingTurn: decodeMatchingTurn(input.matching_turn), + transcript: decodeTranscript(input.transcript), + queuedEvent: decodeQueuedEvent(input.queued_event), + preparedTurn: decodePreparedTurn(input.prepared_turn), + }; +} + +function requestIdentity( + request: RequestFacts, + displayMessageSha256 = request.displayMessageSha256, + attachmentsSha256 = request.attachmentsSha256, +): Sha256 { + return sha256(JSON.stringify({ + schema_version: "loopx_chat_turn_request_identity_v0", + session_id: request.sessionId, + client_turn_id: request.clientTurnId, + execution_message_sha256: request.executionMessageSha256, + display_message_sha256: displayMessageSha256, + attachments_sha256: attachmentsSha256, + origin: request.origin, + loopx_execution: request.loopxExecution, + loopx_request_sha256: request.loopxRequestSha256, + })); +} + +function storedRequestIdentity( + request: RequestFacts, + turn: MatchingTurnFacts, + displayMessageSha256: Sha256, + attachmentsSha256: Sha256, +): Sha256 { + return requestIdentity( + { + sessionId: request.sessionId, + clientTurnId: turn.clientTurnId, + executionMessageSha256: turn.executionMessageSha256, + displayMessageSha256, + attachmentsSha256, + origin: turn.origin, + loopxExecution: turn.loopxExecution, + loopxRequestSha256: turn.loopxRequestSha256, + }, + displayMessageSha256, + attachmentsSha256, + ); +} + +function rejected( + code: ChatTurnAcceptanceRejectionCode, + activeTurnId?: OpaqueId, +): ChatTurnAcceptancePlan { + return { + schema_version: CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA, + kind: "rejected", + code, + ...(activeTurnId ? {active_turn_id: activeTurnId} : {}), + }; +} + +function dispatchFor(status: TurnStatus): ChatTurnDispatch { + if (status === "queued") return {kind: "required"}; + if (status === "starting" || status === "running") { + return {kind: "not_required", reason: "already_started"}; + } + if (status === "completing" || status === "interrupting") { + return {kind: "not_required", reason: "completion_in_progress"}; + } + return {kind: "not_required", reason: "terminal"}; +} + +function activeObservationIsConsistent( + session: SessionFacts, + activeTurn: ActiveTurnFacts | null, +): boolean { + return session.activeTurnId === null + ? activeTurn === null + : activeTurn !== null && activeTurn.turnId === session.activeTurnId; +} + +function activeTurnBlocks( + activeTurn: ActiveTurnFacts | null, + acceptedTurnId: OpaqueId, +): OpaqueId | null { + if ( + activeTurn === null || + activeTurn.turnId === acceptedTurnId || + TERMINAL_TURN_STATUSES.has(activeTurn.status) + ) { + return null; + } + return activeTurn.turnId; +} + +function acceptedPlan(input: { + disposition: "created" | "repaired" | "replayed"; + turnId: OpaqueId; + messageId: OpaqueId; + requestSha256: Sha256; + writes: AcceptanceWrites; + dispatch: ChatTurnDispatch; +}): ChatTurnAcceptancePlan { + return { + schema_version: CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA, + kind: "accepted", + disposition: input.disposition, + turn_id: input.turnId, + message_id: input.messageId, + request_sha256: input.requestSha256, + created: input.disposition === "created", + writes: input.writes, + dispatch: input.dispatch, + }; +} + +function planNewAcceptance( + facts: AcceptanceFacts, + session: SessionFacts, + requestSha256: Sha256, +): ChatTurnAcceptancePlan { + if ( + facts.transcript.kind !== "absent" || + facts.queuedEvent.kind !== "absent" + ) { + return rejected("durable_state_conflict"); + } + if (facts.preparedTurn.kind === "ambiguous") { + return rejected("durable_state_conflict"); + } + if (facts.preparedTurn.kind === "single") { + return rejected("active_turn_conflict", facts.preparedTurn.turnId); + } + const blockingTurnId = activeTurnBlocks( + facts.activeTurn, + facts.candidate.turnId, + ); + if (blockingTurnId !== null) { + return rejected("active_turn_conflict", blockingTurnId); + } + return acceptedPlan({ + disposition: "created", + turnId: facts.candidate.turnId, + messageId: facts.candidate.messageId, + requestSha256, + writes: { + prepare_turn: true, + activate_session: + session.status !== "busy" || + session.activeTurnId !== facts.candidate.turnId, + append_message: true, + append_queued_event: true, + settle_turn: true, + }, + dispatch: {kind: "required"}, + }); +} + +function planPreparedAcceptance( + facts: AcceptanceFacts, + session: SessionFacts, + turn: MatchingTurnFacts, + requestSha256: Sha256, +): ChatTurnAcceptancePlan { + const capsule = turn.acceptance; + if ( + capsule === null || + turn.status !== "queued" || + facts.preparedTurn.kind !== "single" || + facts.preparedTurn.turnId !== turn.turnId || + facts.preparedTurn.clientTurnId !== turn.clientTurnId + ) { + return rejected("durable_state_conflict"); + } + const storedIdentity = storedRequestIdentity( + facts.request, + turn, + capsule.displayMessageSha256, + capsule.attachmentsSha256, + ); + if (storedIdentity !== capsule.requestSha256) { + return rejected("durable_state_conflict"); + } + if (requestSha256 !== capsule.requestSha256) { + return rejected("request_conflict"); + } + const blockingTurnId = activeTurnBlocks(facts.activeTurn, turn.turnId); + if (blockingTurnId !== null) { + return rejected("active_turn_conflict", blockingTurnId); + } + if (facts.transcript.kind === "ambiguous") { + return rejected("durable_state_conflict"); + } + if ( + facts.transcript.kind === "single" && + ( + facts.transcript.messageId !== capsule.messageId || + facts.transcript.displayMessageSha256 !== + capsule.displayMessageSha256 || + facts.transcript.attachmentsSha256 !== capsule.attachmentsSha256 || + facts.transcript.origin !== turn.origin + ) + ) { + return rejected("durable_state_conflict"); + } + if ( + facts.queuedEvent.kind === "ambiguous" || + ( + facts.queuedEvent.kind === "single" && + facts.queuedEvent.payloadSha256 !== EMPTY_OBJECT_SHA256 + ) + ) { + return rejected("durable_state_conflict"); + } + const sessionOwnsPreparedTurn = + session.status === "busy" && session.activeTurnId === turn.turnId; + if ( + ( + facts.transcript.kind === "single" || + facts.queuedEvent.kind === "single" + ) && + !sessionOwnsPreparedTurn + ) { + return rejected("durable_state_conflict"); + } + if ( + facts.queuedEvent.kind === "single" && + facts.transcript.kind !== "single" + ) { + return rejected("durable_state_conflict"); + } + return acceptedPlan({ + disposition: "repaired", + turnId: turn.turnId, + messageId: capsule.messageId, + requestSha256, + writes: { + prepare_turn: false, + activate_session: !sessionOwnsPreparedTurn, + append_message: facts.transcript.kind === "absent", + append_queued_event: facts.queuedEvent.kind === "absent", + settle_turn: true, + }, + dispatch: {kind: "required"}, + }); +} + +function planSettledReplay( + facts: AcceptanceFacts, + session: SessionFacts, + turn: MatchingTurnFacts, + requestSha256: Sha256, +): ChatTurnAcceptancePlan { + if (facts.transcript.kind === "absent") { + return rejected("original_request_unavailable"); + } + if (facts.transcript.kind === "ambiguous") { + return rejected("durable_state_conflict"); + } + const storedIdentity = storedRequestIdentity( + facts.request, + turn, + facts.transcript.displayMessageSha256, + facts.transcript.attachmentsSha256, + ); + if ( + requestSha256 !== storedIdentity || + facts.transcript.origin !== turn.origin + ) { + return rejected("request_conflict"); + } + if ( + facts.queuedEvent.kind !== "single" || + facts.queuedEvent.payloadSha256 !== EMPTY_OBJECT_SHA256 + ) { + return rejected("durable_state_conflict"); + } + const dispatch = dispatchFor(turn.status); + if (dispatch.kind === "required") { + const blockingTurnId = activeTurnBlocks(facts.activeTurn, turn.turnId); + if (blockingTurnId !== null) { + return rejected("active_turn_conflict", blockingTurnId); + } + if (session.activeTurnId !== turn.turnId || session.status !== "busy") { + return rejected("durable_state_conflict"); + } + } else if ( + dispatch.reason === "already_started" && + session.activeTurnId !== turn.turnId + ) { + return rejected("durable_state_conflict"); + } + return acceptedPlan({ + disposition: "replayed", + turnId: turn.turnId, + messageId: facts.transcript.messageId, + requestSha256, + writes: { + prepare_turn: false, + activate_session: false, + append_message: false, + append_queued_event: false, + settle_turn: false, + }, + dispatch, + }); +} + +export function planChatTurnAcceptance( + input: JsonObject, +): ChatTurnAcceptancePlan { + const facts = decodeFacts(input); + if (facts.session === null) return rejected("session_not_found"); + if (facts.session.status === "closed") return rejected("session_closed"); + if (!activeObservationIsConsistent(facts.session, facts.activeTurn)) { + return rejected("durable_state_conflict"); + } + + const requestSha256 = requestIdentity(facts.request); + const matchingTurn = facts.matchingTurn; + if (matchingTurn === null) { + return planNewAcceptance(facts, facts.session, requestSha256); + } + if (matchingTurn.clientTurnId !== facts.request.clientTurnId) { + return rejected("durable_state_conflict"); + } + if (matchingTurn.acceptance !== null) { + return planPreparedAcceptance( + facts, + facts.session, + matchingTurn, + requestSha256, + ); + } + return planSettledReplay( + facts, + facts.session, + matchingTurn, + requestSha256, + ); +} diff --git a/tests/control_plane_ts/chat_turn_acceptance.test.ts b/tests/control_plane_ts/chat_turn_acceptance.test.ts new file mode 100644 index 0000000000..4af1d57a42 --- /dev/null +++ b/tests/control_plane_ts/chat_turn_acceptance.test.ts @@ -0,0 +1,367 @@ +import assert from "node:assert/strict"; +import {createHash} from "node:crypto"; +import test from "node:test"; + +import { + CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA, + CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA, + planChatTurnAcceptance, +} from "../../loopx/control_plane/turn_driver/chat_turn_acceptance.ts"; + +function sha256(value: string): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +function textDigest(value: string): string { + return sha256(value); +} + +function jsonDigest(value: unknown): string { + return sha256(JSON.stringify(value)); +} + +function packet() { + return { + schema_version: CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA, + request: { + session_id: "session-one", + client_turn_id: "client-one", + execution_message_sha256: textDigest("inspect"), + display_message_sha256: textDigest("inspect"), + attachments_sha256: jsonDigest(null), + origin: "web", + loopx_execution: false, + loopx_request_sha256: jsonDigest(null), + }, + candidate: { + turn_id: "candidate-turn", + message_id: "candidate-message", + accepted_at: "2026-09-27T00:00:00Z", + }, + session: { + status: "ready", + active_turn_id: null, + }, + active_turn: null, + matching_turn: null, + transcript: {count: 0}, + queued_event: {count: 0}, + prepared_turn: {count: 0}, + }; +} + +function preparedPacket(options: { + transcript?: "absent" | "single"; + queuedEvent?: "absent" | "single"; + active?: boolean; +} = {}) { + const input = packet(); + const created = planChatTurnAcceptance(input); + assert.equal(created.kind, "accepted"); + if (created.kind !== "accepted") throw new Error("expected accepted plan"); + const active = options.active ?? true; + return { + ...input, + session: { + status: active ? "busy" : "ready", + active_turn_id: active ? "accepted-turn" : null, + }, + active_turn: active + ? {turn_id: "accepted-turn", status: "queued"} + : null, + matching_turn: { + turn_id: "accepted-turn", + client_turn_id: "client-one", + status: "queued", + execution_message_sha256: input.request.execution_message_sha256, + origin: "web", + loopx_execution: false, + loopx_request_sha256: input.request.loopx_request_sha256, + acceptance: { + schema_version: CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA, + phase: "prepared", + request_sha256: created.request_sha256, + message_id: "accepted-message", + display_message_sha256: input.request.display_message_sha256, + attachments_sha256: input.request.attachments_sha256, + }, + }, + transcript: options.transcript === "single" + ? { + count: 1, + message_id: "accepted-message", + display_message_sha256: input.request.display_message_sha256, + attachments_sha256: input.request.attachments_sha256, + origin: "web", + } + : {count: 0}, + queued_event: options.queuedEvent === "single" + ? {count: 1, payload_sha256: jsonDigest({})} + : {count: 0}, + prepared_turn: { + count: 1, + turn_id: "accepted-turn", + client_turn_id: "client-one", + }, + }; +} + +function settledPacket(status: string) { + const input = preparedPacket({ + transcript: "single", + queuedEvent: "single", + }); + return { + ...input, + session: { + status: status === "completed" ? "ready" : "busy", + active_turn_id: status === "completed" ? null : "accepted-turn", + }, + active_turn: status === "completed" + ? null + : {turn_id: "accepted-turn", status}, + matching_turn: { + ...input.matching_turn, + status, + acceptance: null, + }, + prepared_turn: {count: 0}, + }; +} + +test("a new request receives one complete acceptance plan", () => { + const result = planChatTurnAcceptance(packet()); + + assert.equal(result.kind, "accepted"); + if (result.kind !== "accepted") return; + assert.equal(result.disposition, "created"); + assert.equal(result.created, true); + assert.equal(result.turn_id, "candidate-turn"); + assert.equal(result.message_id, "candidate-message"); + assert.match(result.request_sha256, /^sha256:[0-9a-f]{64}$/); + assert.deepEqual(result.writes, { + prepare_turn: true, + activate_session: true, + append_message: true, + append_queued_event: true, + settle_turn: true, + }); + assert.deepEqual(result.dispatch, {kind: "required"}); +}); + +for (const testCase of [ + { + name: "prepared Turn only", + input: preparedPacket({active: false}), + expected: { + prepare_turn: false, + activate_session: true, + append_message: true, + append_queued_event: true, + settle_turn: true, + }, + }, + { + name: "active Turn and transcript", + input: preparedPacket({transcript: "single"}), + expected: { + prepare_turn: false, + activate_session: false, + append_message: false, + append_queued_event: true, + settle_turn: true, + }, + }, + { + name: "complete durable prefix", + input: preparedPacket({ + transcript: "single", + queuedEvent: "single", + }), + expected: { + prepare_turn: false, + activate_session: false, + append_message: false, + append_queued_event: false, + settle_turn: true, + }, + }, +]) { + test(`an exact retry repairs ${testCase.name}`, () => { + const result = planChatTurnAcceptance(testCase.input); + + assert.equal(result.kind, "accepted"); + if (result.kind !== "accepted") return; + assert.equal(result.disposition, "repaired"); + assert.equal(result.created, false); + assert.deepEqual(result.writes, testCase.expected); + assert.deepEqual(result.dispatch, {kind: "required"}); + }); +} + +test("the durable request digest rejects every changed identity field", () => { + const original = preparedPacket(); + const changes = [ + {execution_message_sha256: textDigest("changed")}, + {display_message_sha256: textDigest("changed")}, + {attachments_sha256: jsonDigest([{id: "image-two"}])}, + {origin: "external"}, + {loopx_execution: true}, + {loopx_request_sha256: jsonDigest({operation: "start"})}, + ]; + + for (const change of changes) { + const result = planChatTurnAcceptance({ + ...original, + request: {...original.request, ...change}, + }); + assert.deepEqual(result, { + schema_version: "loopx_chat_turn_acceptance_result_v0", + kind: "rejected", + code: "request_conflict", + }); + } +}); + +test("capsule corruption fails closed before retry identity is considered", () => { + const input = preparedPacket(); + const result = planChatTurnAcceptance({ + ...input, + matching_turn: { + ...input.matching_turn, + execution_message_sha256: textDigest("corrupted"), + }, + request: { + ...input.request, + execution_message_sha256: textDigest("other retry"), + }, + }); + + assert.deepEqual(result, { + schema_version: "loopx_chat_turn_acceptance_result_v0", + kind: "rejected", + code: "durable_state_conflict", + }); +}); + +test("another request cannot displace a prepared Turn", () => { + const input = packet(); + const result = planChatTurnAcceptance({ + ...input, + session: {status: "busy", active_turn_id: "other-turn"}, + active_turn: {turn_id: "other-turn", status: "queued"}, + prepared_turn: { + count: 1, + turn_id: "other-turn", + client_turn_id: "other-client", + }, + }); + + assert.deepEqual(result, { + schema_version: "loopx_chat_turn_acceptance_result_v0", + kind: "rejected", + code: "active_turn_conflict", + active_turn_id: "other-turn", + }); +}); + +for (const [status, dispatch] of [ + ["queued", {kind: "required"}], + ["starting", {kind: "not_required", reason: "already_started"}], + ["running", {kind: "not_required", reason: "already_started"}], + ["completing", {kind: "not_required", reason: "completion_in_progress"}], + ["interrupting", {kind: "not_required", reason: "completion_in_progress"}], + ["completed", {kind: "not_required", reason: "terminal"}], +] as const) { + test(`settled ${status} replay has an independent dispatch decision`, () => { + const result = planChatTurnAcceptance(settledPacket(status)); + + assert.equal(result.kind, "accepted"); + if (result.kind !== "accepted") return; + assert.equal(result.disposition, "replayed"); + assert.equal(result.created, false); + assert.deepEqual(result.dispatch, dispatch); + assert.equal(Object.values(result.writes).some(Boolean), false); + }); +} + +test("a legacy Turn without its original transcript remains unavailable", () => { + const input = settledPacket("queued"); + const result = planChatTurnAcceptance({ + ...input, + transcript: {count: 0}, + }); + + assert.deepEqual(result, { + schema_version: "loopx_chat_turn_acceptance_result_v0", + kind: "rejected", + code: "original_request_unavailable", + }); +}); + +for (const status of ["completing", "interrupting"] as const) { + test(`${status} replay remains valid after the Session releases the Turn`, () => { + const input = settledPacket(status); + const result = planChatTurnAcceptance({ + ...input, + session: {status: "ready", active_turn_id: null}, + active_turn: null, + }); + + assert.equal(result.kind, "accepted"); + if (result.kind !== "accepted") return; + assert.equal(result.disposition, "replayed"); + assert.deepEqual(result.dispatch, { + kind: "not_required", + reason: "completion_in_progress", + }); + }); +} + +test("duplicate transcript rows or queued events fail closed", () => { + const input = preparedPacket(); + for (const change of [ + {transcript: {count: 2}}, + {queued_event: {count: 2}}, + {queued_event: {count: 1, payload_sha256: jsonDigest({changed: true})}}, + ]) { + const result = planChatTurnAcceptance({...input, ...change}); + assert.equal(result.kind, "rejected"); + if (result.kind === "rejected") { + assert.equal(result.code, "durable_state_conflict"); + } + } +}); + +test("repair accepts only a durable prefix of the declared write order", () => { + const messageWithoutSession = preparedPacket({ + active: false, + transcript: "single", + }); + const eventWithoutMessage = preparedPacket({queuedEvent: "single"}); + + for (const input of [messageWithoutSession, eventWithoutMessage]) { + const result = planChatTurnAcceptance(input); + assert.equal(result.kind, "rejected"); + if (result.kind === "rejected") { + assert.equal(result.code, "durable_state_conflict"); + } + } +}); + +test("the Effect boundary rejects malformed digests and observations", () => { + assert.throws( + () => planChatTurnAcceptance({ + ...packet(), + request: {...packet().request, attachments_sha256: "not-a-digest"}, + }), + /attachments_sha256 must be a SHA-256 digest/, + ); + assert.throws( + () => planChatTurnAcceptance({ + ...packet(), + transcript: {count: -1}, + }), + /count must not be negative/, + ); +}); diff --git a/tests/test_chat_session_active_turn.py b/tests/test_chat_session_active_turn.py index cc08f2bef0..ff3b2fc1f0 100644 --- a/tests/test_chat_session_active_turn.py +++ b/tests/test_chat_session_active_turn.py @@ -286,7 +286,7 @@ def test_managed_replay_compares_durable_attachments( assert len(store.messages(session_id)) == 1 -def test_managed_replay_rejects_missing_original_message(tmp_path: Path, monkeypatch) -> None: +def test_managed_replay_repairs_missing_original_message(tmp_path: Path, monkeypatch) -> None: store = ChatSessionStore(tmp_path) session_id = store.create_session( goal_id="goal-one", agent_id="codex", executor_endpoint_id="codex", @@ -300,11 +300,591 @@ def interrupted_append(*args, **kwargs): monkeypatch.setattr(store, "append_message", interrupted_append) with pytest.raises(OSError, match="interrupted transcript"): store.create_turn(session_id, client_turn_id="request", message="inspect") + replay, created = ChatSessionStore(tmp_path).create_turn( + session_id, + client_turn_id="request", + message="inspect", + ) + + assert created is False + assert replay["status"] == "queued" + assert "_acceptance" not in replay + + +def test_legacy_turn_without_original_message_still_fails_closed( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + + def interrupted_append(*args, **kwargs): + raise OSError("interrupted transcript write") + + monkeypatch.setattr(store, "append_message", interrupted_append) + with pytest.raises(OSError, match="interrupted transcript"): + store.create_turn( + session_id, + client_turn_id="legacy-request", + message="inspect", + ) + interrupted = store.turn_for_client(session_id, "legacy-request") + assert interrupted is not None + interrupted.pop("_acceptance") + chat_store._atomic_write_json( + store._turn_path(session_id, str(interrupted["turn_id"])), + interrupted, + preserve_mode=True, + ) + with pytest.raises(ValueError, match="original request is unavailable"): ChatSessionStore(tmp_path).create_turn( - session_id, client_turn_id="request", message="inspect", + session_id, + client_turn_id="legacy-request", + message="inspect", + ) + + +@pytest.mark.parametrize( + "boundary", + [ + "prepare_turn", + "activate_session", + "append_message", + "append_queued_event", + "settle_turn", + ], +) +def test_managed_acceptance_repairs_every_durable_prefix( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + boundary: str, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + failed = False + + if boundary in {"prepare_turn", "activate_session", "settle_turn"}: + original_atomic_write = chat_store._atomic_write_json + + def fail_after_atomic_write( + path: Path, + payload: dict[str, object], + *, + preserve_mode: bool = False, + ) -> None: + nonlocal failed + original_atomic_write( + path, + payload, + preserve_mode=preserve_mode, + ) + is_turn = path.parent.name == "turns" + matches = { + "prepare_turn": is_turn and "_acceptance" in payload, + "activate_session": ( + path.name == "session.json" + and payload.get("active_turn_id") is not None + ), + "settle_turn": ( + is_turn + and preserve_mode + and payload.get("client_turn_id") == "fault-request" + and "_acceptance" not in payload + ), + }[boundary] + if matches and not failed: + failed = True + raise OSError(f"interrupted after {boundary}") + + monkeypatch.setattr( + chat_store, + "_atomic_write_json", + fail_after_atomic_write, + ) + elif boundary == "append_message": + original_append_message = store.append_message + + def fail_after_message(*args, **kwargs): + nonlocal failed + result = original_append_message(*args, **kwargs) + if not failed: + failed = True + raise OSError("interrupted after append_message") + return result + + monkeypatch.setattr(store, "append_message", fail_after_message) + else: + original_append_event = store.append_event + + def fail_after_event(*args, **kwargs): + nonlocal failed + result = original_append_event(*args, **kwargs) + if kwargs.get("kind") == "turn.queued" and not failed: + failed = True + raise OSError("interrupted after append_queued_event") + return result + + monkeypatch.setattr(store, "append_event", fail_after_event) + + with pytest.raises(OSError, match=f"interrupted after {boundary}"): + store.create_turn( + session_id, + client_turn_id="fault-request", + message="recover this request", + attachments=[{"id": "image-one", "mime_type": "image/png"}], + ) + assert failed + + restarted = ChatSessionStore(tmp_path) + partial = restarted.session_snapshot(session_id) + if partial["active_turn"] is not None: + assert "_acceptance" not in partial["active_turn"] + replay, created = restarted.create_turn( + session_id, + client_turn_id="fault-request", + message="recover this request", + attachments=[{"id": "image-one", "mime_type": "image/png"}], + ) + + assert created is False + assert replay["status"] == "queued" + assert "_acceptance" not in replay + assert restarted.load_session(session_id)["active_turn_id"] == replay["turn_id"] # type: ignore[index] + assert [ + message["text"] + for message in restarted.messages(session_id) + if message["role"] == "user" + ] == ["recover this request"] + assert [ + event["kind"] + for event in restarted.events_after( + session_id, + str(replay["turn_id"]), + None, + ) + ] == ["turn.queued"] + + +@pytest.mark.parametrize("interrupted_log", ["message", "event"]) +def test_managed_acceptance_repairs_an_incomplete_jsonl_tail( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + interrupted_log: str, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + original_append = chat_store._append_jsonl_rows + failed = False + + def append_with_interruption( + path: Path, + rows: list[dict[str, object]], + ) -> None: + nonlocal failed + selected = ( + path.name == "messages.jsonl" + if interrupted_log == "message" + else path.name.endswith(".events.jsonl") + ) + if selected and not failed: + failed = True + path.parent.mkdir(parents=True, exist_ok=True) + with path.open("ab") as handle: + handle.write(b'{"text":"\xe4') + handle.flush() + raise OSError(f"interrupted {interrupted_log} append") + original_append(path, rows) + + monkeypatch.setattr( + chat_store, + "_append_jsonl_rows", + append_with_interruption, + ) + with pytest.raises( + OSError, + match=f"interrupted {interrupted_log} append", + ): + store.create_turn( + session_id, + client_turn_id="partial-jsonl", + message="repair the incomplete tail", ) + restarted = ChatSessionStore(tmp_path) + replay, created = restarted.create_turn( + session_id, + client_turn_id="partial-jsonl", + message="repair the incomplete tail", + ) + + assert created is False + assert "_acceptance" not in replay + assert [ + message["text"] + for message in restarted.messages(session_id) + if message["role"] == "user" + ] == ["repair the incomplete tail"] + assert [ + event["kind"] + for event in restarted.events_after( + session_id, + str(replay["turn_id"]), + None, + ) + ] == ["turn.queued"] + + +def test_jsonl_append_preserves_a_valid_final_record_without_newline( + tmp_path: Path, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + store.append_message(session_id, role="user", text="first") + path = store._session_dir(session_id) / "messages.jsonl" + path.write_bytes(path.read_bytes().removesuffix(b"\n")) + + store.append_message(session_id, role="agent", text="second") + + assert [ + message["text"] + for message in store.messages(session_id) + ] == ["first", "second"] + + +def test_same_process_queued_replays_register_one_worker( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + runtime.adapters[session_id] = _HealthyChatAdapter() # type: ignore[assignment] + started = threading.Event() + release = threading.Event() + starts = 0 + + def blocked_run_turn(**kwargs: object) -> None: + nonlocal starts + starts += 1 + started.set() + release.wait(timeout=2) + done_event = kwargs["done_event"] + assert isinstance(done_event, threading.Event) + with runtime.lock: + runtime.turn_done_events.pop( + (session_id, str(kwargs["turn_id"])), + None, + ) + done_event.set() + + monkeypatch.setattr(runtime, "_run_turn", blocked_run_turn) + first, first_created = runtime.submit_turn( + session_id=session_id, + client_turn_id="one-local-worker", + message="run once", + work_dir=tmp_path, + objective="sample objective", + ) + assert first_created is True + assert started.wait(timeout=2) + + replay, replay_created = runtime.submit_turn( + session_id=session_id, + client_turn_id="one-local-worker", + message="run once", + work_dir=tmp_path, + objective="sample objective", + ) + + assert replay_created is False + assert replay["turn_id"] == first["turn_id"] + assert starts == 1 + with runtime.lock: + assert (session_id, str(first["turn_id"])) in runtime.turn_done_events + release.set() + + +def test_concurrent_managed_retries_dispatch_once( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + accepted, created = store.create_turn( + session_id, + client_turn_id="response-lost", + message="continue exactly once", + ) + assert created is True + + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + adapter = _BlockingChatAdapter() + start_calls = 0 + start_lock = threading.Lock() + original_start_turn = adapter.start_turn + + def counted_start_turn(message: str, event_sink): + nonlocal start_calls + with start_lock: + start_calls += 1 + return original_start_turn(message, event_sink) + + monkeypatch.setattr(adapter, "start_turn", counted_start_turn) + runtime.adapters[session_id] = adapter # type: ignore[assignment] + results: list[tuple[dict[str, object], bool]] = [] + + def retry() -> None: + results.append( + runtime.submit_turn( + session_id=session_id, + client_turn_id="response-lost", + message="continue exactly once", + work_dir=tmp_path, + objective="sample objective", + ) + ) + + retries = [threading.Thread(target=retry) for _ in range(2)] + for thread in retries: + thread.start() + for thread in retries: + thread.join(timeout=2) + + assert not any(thread.is_alive() for thread in retries) + assert len(results) == 2 + assert all(created is False for _turn, created in results) + assert {turn["turn_id"] for turn, _created in results} == { + accepted["turn_id"] + } + assert adapter.started.wait(timeout=2) + assert start_calls == 1 + + adapter.release.set() + assert runtime.wait_for_turn( + session_id=session_id, + turn_id=str(accepted["turn_id"]), + timeout_sec=2, + )["status"] == "completed" + + +def test_adapter_start_failure_keeps_accepted_turn_retryable( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + adapter = _BlockingChatAdapter() + attempts = 0 + + def start_adapter(**_kwargs): + nonlocal attempts + attempts += 1 + if attempts == 1: + raise OSError("adapter startup interrupted") + return adapter + + monkeypatch.setattr(runtime, "_start_adapter", start_adapter) + with pytest.raises(CodexChatAgentError, match="could not be restored"): + runtime.submit_turn( + session_id=session_id, + client_turn_id="adapter-retry", + message="survive adapter startup", + work_dir=tmp_path, + objective="sample objective", + ) + + persisted = store.turn_for_client(session_id, "adapter-retry") + assert persisted is not None + assert persisted["status"] == "queued" + assert "_acceptance" not in persisted + assert store.load_session(session_id)["active_turn_id"] == persisted["turn_id"] # type: ignore[index] + + replay, created = runtime.submit_turn( + session_id=session_id, + client_turn_id="adapter-retry", + message="survive adapter startup", + work_dir=tmp_path, + objective="sample objective", + ) + + assert created is False + assert replay["turn_id"] == persisted["turn_id"] + assert adapter.started.wait(timeout=2) + adapter.release.set() + assert runtime.wait_for_turn( + session_id=session_id, + turn_id=str(replay["turn_id"]), + timeout_sec=2, + )["status"] == "completed" + + +def test_resume_repairs_and_dispatches_a_prepared_turn( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + original_append_message = store.append_message + + def interrupted_append(*args, **kwargs): + raise OSError("interrupted before transcript") + + monkeypatch.setattr(store, "append_message", interrupted_append) + with pytest.raises(OSError, match="interrupted before transcript"): + store.create_turn( + session_id, + client_turn_id="resume-prepared", + message="resume the accepted request", + ) + monkeypatch.setattr(store, "append_message", original_append_message) + + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + adapter = _BlockingChatAdapter() + monkeypatch.setattr(runtime, "_start_adapter", lambda **_kwargs: adapter) + + restored = runtime.resume_session( + session_id=session_id, + work_dir=tmp_path, + objective="sample objective", + ) + + prepared = store.turn_for_client(session_id, "resume-prepared") + assert prepared is not None + assert "_acceptance" not in prepared + assert restored["status"] == "busy" + assert restored["active_turn_id"] == prepared["turn_id"] + assert adapter.started.wait(timeout=2) + adapter.release.set() + assert runtime.wait_for_turn( + session_id=session_id, + turn_id=str(prepared["turn_id"]), + timeout_sec=2, + )["status"] == "completed" + + +def test_replay_after_process_loss_fails_an_unowned_started_turn( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + )["session_id"] + ) + turn, _created = store.create_turn( + session_id, + client_turn_id="lost-running-worker", + message="do not leave this running", + ) + store.update_turn( + session_id, + str(turn["turn_id"]), + expected_statuses={"queued"}, + status="starting", + ) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + monkeypatch.setattr( + runtime, + "_start_adapter", + lambda **_kwargs: _HealthyChatAdapter(), + ) + + replay, created = runtime.submit_turn( + session_id=session_id, + client_turn_id="lost-running-worker", + message="do not leave this running", + work_dir=tmp_path, + objective="sample objective", + ) + + assert created is False + assert replay["status"] == "failed" + assert replay["error_code"] == "server_restarted" + restored = store.load_session(session_id) + assert restored is not None + assert restored["status"] == "ready" + assert restored["active_turn_id"] is None + def test_completed_turn_cannot_release_a_newer_active_turn(tmp_path: Path) -> None: store = ChatSessionStore(tmp_path) diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 37cd8ab98d..1a9a48fc8e 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -73,6 +73,7 @@ "loopx/control_plane/todos/next_action.ts", "loopx/control_plane/turn_driver/turn_journal.ts", "loopx/control_plane/turn_driver/turn_journal_effects.ts", + "loopx/control_plane/turn_driver/chat_turn_acceptance.ts", "loopx/control_plane/turn_driver/delivery_continuity.ts", "loopx/control_plane/work_items/delivery_outcome.ts", "loopx/control_plane/work_items/interaction_contract.ts", @@ -148,6 +149,7 @@ "tests/control_plane_ts/legacy_writer_fence_read_error.test.ts", "tests/control_plane_ts/turn_journal.test.ts", "tests/control_plane_ts/turn_journal_effects.test.ts", + "tests/control_plane_ts/chat_turn_acceptance.test.ts", "tests/control_plane_ts/vision_checkpoint.test.ts", "tests/control_plane_ts/checkpoint_read_context.test.ts", "tests/control_plane_ts/checkpoint_provider_head.test.ts", From 063de90e5cacffa78cdd4a93b8fd0bee4ef06c0b Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 12:22:48 +0800 Subject: [PATCH 2/6] fix(chat): retire interrupted acceptance capsules Signed-off-by: duanjialing.777 --- loopx/chat_runtime.py | 123 ++++++++++++------ loopx/chat_store.py | 20 ++- loopx/chat_turn_acceptance.py | 10 +- .../turn_driver/chat_turn_acceptance.ts | 26 +++- .../chat_turn_acceptance.test.ts | 72 ++++++++++ tests/test_chat_queue_preparation.py | 54 ++++++++ tests/test_chat_session_active_turn.py | 50 +++++++ 7 files changed, 306 insertions(+), 49 deletions(-) diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index b95b21b2a6..f927fdbe8e 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -908,6 +908,10 @@ def submit_turn( origin="web", ) with self._session_adapter_lock(session_id): + session = self.store.load_session(session_id) + if session is None or session.get("status") == "closed": + raise KeyError("chat session was not found") + self._check_codex_home(session) if loopx_execution and session.get("loopx_tools") is not True: session = self.loopx_mode.activate_tools(session, work_dir=work_dir, objective=objective) accepted = self.store.accept_managed_turn( @@ -1270,21 +1274,37 @@ def _drain_session_queue( def _fail_queue_preparation( self, session_id: str, turn_id: str, error: Exception, ) -> None: + key = (session_id, turn_id) + done_event = threading.Event() + with self.lock: + current = self.turn_done_events.get(key) + owns_done_event = current is None + if owns_done_event: + self.turn_done_events[key] = done_event + else: + done_event = current logging.getLogger(__name__).error( "Chat queue runtime preparation failed", exc_info=(type(error), error, error.__traceback__), ) - self._fail_turn( - session_id, turn_id, - error.error_code if isinstance(error, CodexChatAgentError) else "runtime_unavailable", - str(error) if isinstance(error, CodexChatAgentError) else ( - "The Agent runtime could not be prepared. Check the runtime installation " - "and configuration, then retry." - ), - status="failed", - gate=error.gate if isinstance(error, CodexChatAgentError) else None, - expected_statuses={"queued", "starting", "running"}, - ) + try: + self._fail_turn( + session_id, turn_id, + error.error_code if isinstance(error, CodexChatAgentError) else "runtime_unavailable", + str(error) if isinstance(error, CodexChatAgentError) else ( + "The Agent runtime could not be prepared. Check the runtime installation " + "and configuration, then retry." + ), + status="failed", + gate=error.gate if isinstance(error, CodexChatAgentError) else None, + expected_statuses={"queued", "starting", "running"}, + ) + finally: + done_event.set() + if owns_done_event: + with self.lock: + if self.turn_done_events.get(key) is done_event: + self.turn_done_events.pop(key, None) def _run_turn( self, @@ -1545,32 +1565,49 @@ def _fail_turn( ) def interrupt_turn(self, *, session_id: str, turn_id: str) -> dict[str, Any]: - turn = self.store.load_turn(session_id, turn_id) - if turn is None: - raise KeyError("chat turn was not found") - if turn.get("status") in TERMINAL_TURN_STATES: - return turn - session = self.store.load_session(session_id) - if ( - session - and session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED - and session.get("active_turn_id") == turn_id - ): - raise CodexChatAgentError( - "The attached host does not expose interrupt control to LoopX Chat.", - error_code="attached_session_interrupt_unavailable", - gate={ - "kind": "host_tool_gate", - "summary": "The active Turn is owned by the attached host.", - "next_action": "Stop the Turn in the attached host, then retry.", - }, + with self._session_adapter_lock(session_id): + turn = self.store.load_turn(session_id, turn_id) + if turn is None: + raise KeyError("chat turn was not found") + session = self.store.load_session(session_id) + if ( + session + and session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED + and session.get("active_turn_id") == turn_id + ): + raise CodexChatAgentError( + "The attached host does not expose interrupt control to LoopX Chat.", + error_code="attached_session_interrupt_unavailable", + gate={ + "kind": "host_tool_gate", + "summary": "The active Turn is owned by the attached host.", + "next_action": "Stop the Turn in the attached host, then retry.", + }, + ) + prepared = self.store.prepared_managed_turn_request( + session_id, + turn_id=turn_id, + ) + if prepared is not None: + accepted = self.store.accept_managed_turn( + session_id, + client_turn_id=str(prepared["client_turn_id"]), + message=str(prepared["message"]), + attachments=prepared.get("attachments"), + origin=str(prepared["origin"]), + display_message=str(prepared["display_message"]), + loopx_execution=bool(prepared["loopx_execution"]), + loopx_request=prepared.get("loopx_request"), + ) + turn = accepted.turn + if turn.get("status") in TERMINAL_TURN_STATES: + return turn + interrupting = self.store.update_turn( + session_id, + turn_id, + expected_statuses={"queued", "starting", "running"}, + status="interrupting", ) - interrupting = self.store.update_turn( - session_id, - turn_id, - expected_statuses={"queued", "starting", "running"}, - status="interrupting", - ) if interrupting is None: current = self.store.load_turn(session_id, turn_id) if current is None: @@ -1646,12 +1683,18 @@ def wait_for_turn(self, *, session_id: str, turn_id: str, timeout_sec: float = 9 while True: if (turn := self.store.load_turn(session_id, turn_id)) is None: raise KeyError("chat turn was not found") - if turn.get("status") in TERMINAL_TURN_STATES: - return turn - if (remaining := deadline - time.monotonic()) <= 0: - raise TimeoutError("chat turn wait timed out") with self.lock: done_event = self.turn_done_events.get((session_id, turn_id)) + remaining = deadline - time.monotonic() + if turn.get("status") in TERMINAL_TURN_STATES: + if done_event is None or done_event.is_set(): + return turn + if remaining <= 0: + raise TimeoutError("chat turn wait timed out") + done_event.wait(remaining) + continue + if remaining <= 0: + raise TimeoutError("chat turn wait timed out") (done_event.wait if done_event else time.sleep)(remaining if done_event else min(0.02, remaining)) def close_session(self, session_id: str) -> bool: diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 7d44233e10..3ad939aa75 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -80,7 +80,7 @@ def _read_json(path: Path) -> dict[str, Any]: def _read_jsonl(path: Path) -> list[dict[str, Any]]: try: - lines = path.read_text(encoding="utf-8").split("\n") + lines = path.read_bytes().split(b"\n") except OSError: return [] rows: list[dict[str, Any]] = [] @@ -626,8 +626,15 @@ def messages(self, session_id: str) -> list[dict[str, Any]]: def prepared_managed_turn_request( self, session_id: str, + *, + turn_id: str | None = None, ) -> dict[str, Any] | None: session_token = _opaque_id(session_id, field="session_id") + turn_token = ( + _opaque_id(turn_id, field="turn_id") + if turn_id is not None + else None + ) session_path = self._session_path(session_token) with self._session_lock(session_token): with exclusive_file_lock( @@ -644,6 +651,17 @@ def prepared_managed_turn_request( payload := _read_json(path) ).get("schema_version") == CHAT_TURN_SCHEMA_VERSION and "_acceptance" in payload + and ( + ( + turn_token is not None + and payload.get("turn_id") == turn_token + ) + or ( + turn_token is None + and payload.get("status") + not in TERMINAL_TURN_STATES + ) + ) ] if not prepared: return None diff --git a/loopx/chat_turn_acceptance.py b/loopx/chat_turn_acceptance.py index 7f4187a1d2..e15df8b7cc 100644 --- a/loopx/chat_turn_acceptance.py +++ b/loopx/chat_turn_acceptance.py @@ -13,6 +13,7 @@ CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA = "loopx_chat_turn_acceptance_request_v0" CHAT_TURN_ACCEPTANCE_RESULT_SCHEMA = "loopx_chat_turn_acceptance_result_v0" CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA = "loopx_chat_turn_acceptance_v0" +TERMINAL_TURN_STATUSES = {"completed", "interrupted", "timed_out", "failed"} def _sha256(value: str) -> str: @@ -112,7 +113,14 @@ def _queued_event_facts(events: list[dict[str, Any]]) -> dict[str, Any]: def _prepared_turn_facts( turns: list[dict[str, Any]], ) -> dict[str, Any]: - prepared = [turn for turn in turns if "_acceptance" in turn] + prepared = [ + turn + for turn in turns + if ( + "_acceptance" in turn + and turn.get("status") not in TERMINAL_TURN_STATUSES + ) + ] if len(prepared) != 1: return {"count": len(prepared)} return { diff --git a/loopx/control_plane/turn_driver/chat_turn_acceptance.ts b/loopx/control_plane/turn_driver/chat_turn_acceptance.ts index 879a105baa..c30959b250 100644 --- a/loopx/control_plane/turn_driver/chat_turn_acceptance.ts +++ b/loopx/control_plane/turn_driver/chat_turn_acceptance.ts @@ -553,8 +553,15 @@ function planPreparedAcceptance( requestSha256: Sha256, ): ChatTurnAcceptancePlan { const capsule = turn.acceptance; - if ( - capsule === null || + if (capsule === null) { + return rejected("durable_state_conflict"); + } + const terminal = TERMINAL_TURN_STATUSES.has(turn.status); + if (terminal) { + if (facts.preparedTurn.kind !== "absent") { + return rejected("durable_state_conflict"); + } + } else if ( turn.status !== "queued" || facts.preparedTurn.kind !== "single" || facts.preparedTurn.turnId !== turn.turnId || @@ -574,9 +581,11 @@ function planPreparedAcceptance( if (requestSha256 !== capsule.requestSha256) { return rejected("request_conflict"); } - const blockingTurnId = activeTurnBlocks(facts.activeTurn, turn.turnId); - if (blockingTurnId !== null) { - return rejected("active_turn_conflict", blockingTurnId); + if (!terminal) { + const blockingTurnId = activeTurnBlocks(facts.activeTurn, turn.turnId); + if (blockingTurnId !== null) { + return rejected("active_turn_conflict", blockingTurnId); + } } if (facts.transcript.kind === "ambiguous") { return rejected("durable_state_conflict"); @@ -605,6 +614,7 @@ function planPreparedAcceptance( const sessionOwnsPreparedTurn = session.status === "busy" && session.activeTurnId === turn.turnId; if ( + !terminal && ( facts.transcript.kind === "single" || facts.queuedEvent.kind === "single" @@ -626,12 +636,14 @@ function planPreparedAcceptance( requestSha256, writes: { prepare_turn: false, - activate_session: !sessionOwnsPreparedTurn, + activate_session: !terminal && !sessionOwnsPreparedTurn, append_message: facts.transcript.kind === "absent", append_queued_event: facts.queuedEvent.kind === "absent", settle_turn: true, }, - dispatch: {kind: "required"}, + dispatch: terminal + ? {kind: "not_required", reason: "terminal"} + : {kind: "required"}, }); } diff --git a/tests/control_plane_ts/chat_turn_acceptance.test.ts b/tests/control_plane_ts/chat_turn_acceptance.test.ts index 4af1d57a42..348dfcad09 100644 --- a/tests/control_plane_ts/chat_turn_acceptance.test.ts +++ b/tests/control_plane_ts/chat_turn_acceptance.test.ts @@ -129,6 +129,25 @@ function settledPacket(status: string) { }; } +function terminalPreparedPacket(options: { + transcript?: "absent" | "single"; + queuedEvent?: "absent" | "single"; +} = {}) { + const input = preparedPacket({ + transcript: options.transcript, + queuedEvent: options.queuedEvent, + active: false, + }); + return { + ...input, + matching_turn: { + ...input.matching_turn, + status: "interrupted", + }, + prepared_turn: {count: 0}, + }; +} + test("a new request receives one complete acceptance plan", () => { const result = planChatTurnAcceptance(packet()); @@ -199,6 +218,59 @@ for (const testCase of [ }); } +for (const testCase of [ + { + name: "prepared Turn only", + input: terminalPreparedPacket(), + expected: { + prepare_turn: false, + activate_session: false, + append_message: true, + append_queued_event: true, + settle_turn: true, + }, + }, + { + name: "terminal Turn and transcript", + input: terminalPreparedPacket({transcript: "single"}), + expected: { + prepare_turn: false, + activate_session: false, + append_message: false, + append_queued_event: true, + settle_turn: true, + }, + }, + { + name: "complete terminal durable prefix", + input: terminalPreparedPacket({ + transcript: "single", + queuedEvent: "single", + }), + expected: { + prepare_turn: false, + activate_session: false, + append_message: false, + append_queued_event: false, + settle_turn: true, + }, + }, +]) { + test(`an exact retry retires ${testCase.name} without dispatch`, () => { + const result = planChatTurnAcceptance(testCase.input); + + assert.equal(result.kind, "accepted"); + if (result.kind !== "accepted") return; + assert.equal(result.disposition, "repaired"); + assert.equal(result.created, false); + assert.deepEqual(result.writes, testCase.expected); + assert.deepEqual(result.dispatch, { + kind: "not_required", + reason: "terminal", + }); + }); +} + test("the durable request digest rejects every changed identity field", () => { const original = preparedPacket(); const changes = [ diff --git a/tests/test_chat_queue_preparation.py b/tests/test_chat_queue_preparation.py index a970d6ae29..d8ecfd561c 100644 --- a/tests/test_chat_queue_preparation.py +++ b/tests/test_chat_queue_preparation.py @@ -130,6 +130,60 @@ def test_removed_manager_release_fails_before_starting_provider(tmp_path, monkey assert store.load_session(sid)["active_turn_id"] is None +def test_wait_for_preparation_failure_observes_released_session( + tmp_path, + monkeypatch, +): + store, sid, runtime = session_runtime(tmp_path) + queued, _ = store.create_queued_turn( + sid, + client_turn_id="settlement-publication", + message="Research cash flow", + ) + claimed = store.claim_next_queued_turn(sid) + assert claimed is not None + assert claimed["turn_id"] == queued["turn_id"] + + release_entered = threading.Event() + allow_release = threading.Event() + original_release = store.release_active_turn + + def blocked_release(*args, **kwargs): + release_entered.set() + assert allow_release.wait(3) + return original_release(*args, **kwargs) + + monkeypatch.setattr(store, "release_active_turn", blocked_release) + failure = threading.Thread( + target=runtime._fail_queue_preparation, + args=(sid, str(queued["turn_id"]), OSError("runtime removed")), + ) + failure.start() + assert release_entered.wait(3) + + observed = [] + waiter = threading.Thread( + target=lambda: observed.append( + runtime.wait_for_turn( + session_id=sid, + turn_id=str(queued["turn_id"]), + timeout_sec=3, + ) + ) + ) + waiter.start() + waiter.join(0.05) + assert waiter.is_alive() + + allow_release.set() + failure.join(3) + waiter.join(3) + assert not failure.is_alive() + assert not waiter.is_alive() + assert observed[0]["status"] == "failed" + assert store.load_session(sid)["active_turn_id"] is None + + @pytest.mark.parametrize("error,code", [ (OSError("private installation path"), "runtime_unavailable"), (CodexChatAgentError("Cannot restore", error_code="resume_failed", gate=None), "resume_failed"), diff --git a/tests/test_chat_session_active_turn.py b/tests/test_chat_session_active_turn.py index ff3b2fc1f0..82de98c2e5 100644 --- a/tests/test_chat_session_active_turn.py +++ b/tests/test_chat_session_active_turn.py @@ -364,10 +364,12 @@ def interrupted_append(*args, **kwargs): "settle_turn", ], ) +@pytest.mark.parametrize("interrupt_before_retry", [False, True]) def test_managed_acceptance_repairs_every_durable_prefix( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, boundary: str, + interrupt_before_retry: bool, ) -> None: store = ChatSessionStore(tmp_path) session_id = str( @@ -455,6 +457,54 @@ def fail_after_event(*args, **kwargs): assert failed restarted = ChatSessionStore(tmp_path) + if interrupt_before_retry: + partial_turn = restarted.turn_for_client( + session_id, + "fault-request", + ) + assert partial_turn is not None + runtime = ChatRuntimeController( + store=restarted, + codex_bin="missing-codex", + ) + interrupted = runtime.interrupt_turn( + session_id=session_id, + turn_id=str(partial_turn["turn_id"]), + ) + assert interrupted["status"] == "interrupted" + assert "_acceptance" not in interrupted + + replay = restarted.accept_managed_turn( + session_id, + client_turn_id="fault-request", + message="recover this request", + attachments=[ + {"id": "image-one", "mime_type": "image/png"}, + ], + ) + assert replay.created is False + assert replay.turn["status"] == "interrupted" + assert replay.dispatch_required is False + assert replay.dispatch_reason == "terminal" + + next_turn, created = restarted.create_turn( + session_id, + client_turn_id="next-request", + message="continue after interruption", + ) + assert created is True + assert next_turn["status"] == "queued" + assert restarted.load_session(session_id)["active_turn_id"] == next_turn["turn_id"] # type: ignore[index] + assert [ + message["text"] + for message in restarted.messages(session_id) + if ( + message["role"] == "user" + and message["turn_id"] == interrupted["turn_id"] + ) + ] == ["recover this request"] + return + partial = restarted.session_snapshot(session_id) if partial["active_turn"] is not None: assert "_acceptance" not in partial["active_turn"] From 5a64fdeaa12d24653c8301a4a399735f0be6722b Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 19:55:05 +0800 Subject: [PATCH 3/6] fix(chat): classify acceptance write failures Signed-off-by: duanjialing.777 --- loopx/chat_runtime.py | 44 ++++++++++++------- loopx/chat_server.py | 15 ++++++- .../project_registry_io_manifest_v1.json | 6 +-- 3 files changed, 44 insertions(+), 21 deletions(-) diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index f927fdbe8e..fde39ee609 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -68,6 +68,11 @@ "live_steering_turn_not_started", }) + +class ChatTurnAcceptanceUnavailableError(Exception): + """The durable acceptance attempt can be retried with the same request.""" + + class ChatRuntimeAdapter(Protocol): @property def upstream_thread_id(self) -> str: ... @@ -914,23 +919,28 @@ def submit_turn( self._check_codex_home(session) if loopx_execution and session.get("loopx_tools") is not True: session = self.loopx_mode.activate_tools(session, work_dir=work_dir, objective=objective) - accepted = self.store.accept_managed_turn( - session_id, - client_turn_id=client_turn_id, - message=message, - attachments=attachments, - display_message=( - ( - "开启 LoopX 模式,持续推进当前 Goal。" - if (loopx_request or {}).get("operation") == "start" - else "恢复 LoopX 模式。" - ) - if loopx_execution - else None - ), - loopx_execution=loopx_execution, - loopx_request=loopx_request, - ) + try: + accepted = self.store.accept_managed_turn( + session_id, + client_turn_id=client_turn_id, + message=message, + attachments=attachments, + display_message=( + ( + "开启 LoopX 模式,持续推进当前 Goal。" + if (loopx_request or {}).get("operation") == "start" + else "恢复 LoopX 模式。" + ) + if loopx_execution + else None + ), + loopx_execution=loopx_execution, + loopx_request=loopx_request, + ) + except OSError as exc: + raise ChatTurnAcceptanceUnavailableError( + "Chat turn acceptance is temporarily unavailable." + ) from exc if accepted.dispatch_required: adapter = self._ensure_adapter_locked( session, diff --git a/loopx/chat_server.py b/loopx/chat_server.py index 9e86efa7c0..6f1feab2e7 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -29,7 +29,12 @@ add_goal_subagent_routes, ) from .chat_status_api import ChatStatusRequestMixin -from .chat_runtime import ChatRuntimeController, STEERING_NOT_DELIVERED_CODES, TERMINAL_TURN_STATES +from .chat_runtime import ( + STEERING_NOT_DELIVERED_CODES, + TERMINAL_TURN_STATES, + ChatRuntimeController, + ChatTurnAcceptanceUnavailableError, +) from .chat_manager import ( MANAGER_AGENT_GOAL_ID, MANAGER_AGENT_OBJECTIVE, is_manager_channel, manager_capabilities_projection, manager_workspace, @@ -716,6 +721,14 @@ def _session_turn(self, session_id: str) -> None: ) return response = completed.get("response") + except ChatTurnAcceptanceUnavailableError: + self._send_error( + "Chat turn acceptance is temporarily unavailable.", + status=503, + error_code="chat_turn_acceptance_unavailable", + turn_replay_safe=True, + ) + return except CodexChatAgentError as exc: self._send_error( str(exc), diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 34539c92a7..c881453cca 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -383,7 +383,7 @@ }, { "site": "loopx/chat_server.py::.ChatRequestHandler._goal_channel_extension_ready::codec_read:load_registry#1", - "line": 949, + "line": 963, "column": 24, "kind": "codec_read", "api": "load_registry", @@ -391,7 +391,7 @@ }, { "site": "loopx/chat_server.py::.ChatRequestHandler._registry_and_goal::codec_read:load_registry#1", - "line": 506, + "line": 511, "column": 20, "kind": "codec_read", "api": "load_registry", @@ -399,7 +399,7 @@ }, { "site": "loopx/chat_server.py::.serve_chat::codec_read:load_registry#1", - "line": 1471, + "line": 1485, "column": 16, "kind": "codec_read", "api": "load_registry", From 62f375af463de7a6cfb32a83794e8d6a38f6b22a Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 19:55:24 +0800 Subject: [PATCH 4/6] test(chat): cover HTTP acceptance recovery Signed-off-by: duanjialing.777 --- .github/workflows/python-tests.yml | 1 + .../chat-turn-acceptance-http-fixture.py | 224 ++++++++++++++++++ .../smoke/chat-turn-acceptance-retry-smoke.ts | 194 +++++++++++++++ tests/test_python_ci_workflow.py | 7 +- 4 files changed, 425 insertions(+), 1 deletion(-) create mode 100644 apps/presentation/dashboard/smoke/chat-turn-acceptance-http-fixture.py diff --git a/.github/workflows/python-tests.yml b/.github/workflows/python-tests.yml index 5d685b0997..08d3a1db8d 100644 --- a/.github/workflows/python-tests.yml +++ b/.github/workflows/python-tests.yml @@ -107,6 +107,7 @@ jobs: run: | ./node_modules/.bin/playwright install --with-deps chromium npm run smoke:personal-workspace-packaged + npm run smoke:chat-turn-acceptance-retry npm run smoke:chat-upgrade - uses: actions/upload-artifact@v7 with: diff --git a/apps/presentation/dashboard/smoke/chat-turn-acceptance-http-fixture.py b/apps/presentation/dashboard/smoke/chat-turn-acceptance-http-fixture.py new file mode 100644 index 0000000000..015d4c04c1 --- /dev/null +++ b/apps/presentation/dashboard/smoke/chat-turn-acceptance-http-fixture.py @@ -0,0 +1,224 @@ +from __future__ import annotations + +import json +from pathlib import Path +import sys +import tempfile +import threading +from typing import Any + +sys.path.insert(0, str(Path(__file__).resolve().parents[4])) + +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_server import ChatHTTPServer, ChatRequestHandler +from loopx.chat_store import ChatSessionStore + + +class _HealthyAdapter: + upstream_thread_id = "fixture-upstream" + + def capabilities(self) -> dict[str, Any]: + return {} + + def start_turn( + self, + message: str, + event_sink: Any, + ) -> dict[str, Any]: + del message, event_sink + raise AssertionError("the fixture replaces runtime dispatch") + + def interrupt_turn(self, turn_id: str | None = None) -> None: + del turn_id + + def healthcheck(self) -> bool: + return True + + def close_session(self) -> None: + return None + + +def _turn_count(store: ChatSessionStore, session_id: str) -> int: + return sum( + 1 + for path in (store.sessions_root / session_id / "turns").glob("*.json") + if not path.name.endswith(".events.json") + ) + + +def main() -> int: + scenario = sys.argv[1] if len(sys.argv) > 1 else "" + if scenario not in {"before_transcript", "after_queued"}: + raise ValueError("unknown acceptance fault scenario") + + root = Path(tempfile.mkdtemp(prefix="loopx-chat-acceptance-http-")) + project = root / "project" + project.mkdir() + registry_path = root / "registry.json" + registry_path.write_text( + json.dumps( + { + "schema_version": "0.1", + "goals": [ + { + "id": "goal-one", + "repo": str(project), + "status": "active", + } + ], + } + ), + encoding="utf-8", + ) + store = ChatSessionStore(root / "runtime") + runtime = ChatRuntimeController( + store=store, + codex_bin="missing-codex", + registry_path=registry_path, + ) + session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="fixture-upstream", + upstream_mode="chat", + codex_home=str(runtime.codex_home), + )["session_id"] + ) + mismatched_session_id = str( + store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="other-upstream", + upstream_mode="chat", + codex_home=str(root / "other-codex-home"), + )["session_id"] + ) + runtime.adapters[session_id] = _HealthyAdapter() + + dispatch_started = threading.Event() + dispatch_release = threading.Event() + dispatch_count = 0 + + def blocked_run_turn(**kwargs: object) -> None: + nonlocal dispatch_count + dispatch_count += 1 + dispatch_started.set() + dispatch_release.wait(timeout=10) + done_event = kwargs["done_event"] + if not isinstance(done_event, threading.Event): + raise TypeError("turn dispatch must carry a completion event") + with runtime.lock: + runtime.turn_done_events.pop( + (session_id, str(kwargs["turn_id"])), + None, + ) + done_event.set() + + runtime._run_turn = blocked_run_turn + + failed = False + if scenario == "before_transcript": + original_append_message = store.append_message + + def fail_before_transcript(*args: Any, **kwargs: Any) -> dict[str, Any]: + nonlocal failed + if not failed and kwargs.get("role") == "user": + failed = True + raise OSError("private transcript fault detail") + result = original_append_message(*args, **kwargs) + if not isinstance(result, dict): + raise TypeError("append_message returned an invalid result") + return result + + store.append_message = fail_before_transcript + else: + original_append_event = store.append_event + + def fail_after_queued(*args: Any, **kwargs: Any) -> dict[str, Any]: + nonlocal failed + result = original_append_event(*args, **kwargs) + if not isinstance(result, dict): + raise TypeError("append_event returned an invalid result") + if not failed and kwargs.get("kind") == "turn.queued": + failed = True + raise OSError("private queued event fault detail") + return result + + store.append_event = fail_after_queued + + server = ChatHTTPServer(("127.0.0.1", 0), ChatRequestHandler) + server.verbose = False + server.registry_path = registry_path + server.chat_store = store + server.runtime_controller = runtime + server_thread = threading.Thread(target=server.serve_forever, daemon=True) + server_thread.start() + print( + json.dumps( + { + "origin": f"http://127.0.0.1:{server.server_port}", + "session_id": session_id, + "mismatched_session_id": mismatched_session_id, + } + ), + flush=True, + ) + + try: + if sys.stdin.readline().strip() != "inspect": + raise ValueError("fixture did not receive an inspect request") + if not dispatch_started.wait(timeout=5): + raise TimeoutError("accepted turn was not dispatched") + turn = store.turn_for_client(session_id, "recoverable-request") + if turn is None: + raise AssertionError("accepted turn is missing") + queued_events = [ + event + for event in store.events_after( + session_id, + str(turn["turn_id"]), + None, + ) + if event.get("kind") == "turn.queued" + ] + user_messages = [ + message + for message in store.messages(session_id) + if message.get("role") == "user" + ] + mismatch_turns = _turn_count(store, mismatched_session_id) + mismatch_messages = store.messages(mismatched_session_id) + print( + json.dumps( + { + "fault_injected": failed, + "turn_count": _turn_count(store, session_id), + "user_message_count": len(user_messages), + "queued_event_count": len(queued_events), + "dispatch_count": dispatch_count, + "acceptance_capsule_present": "_acceptance" in turn, + "mismatched_turn_count": mismatch_turns, + "mismatched_message_count": len(mismatch_messages), + "mismatched_session_status": ( + store.load_session(mismatched_session_id) or {} + ).get("status"), + } + ), + flush=True, + ) + finally: + dispatch_release.set() + server.shutdown() + server.server_close() + server_thread.join(timeout=5) + runtime.close() + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts index c4262ae8d3..6e3c69bd69 100644 --- a/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts +++ b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts @@ -1,4 +1,9 @@ import assert from "node:assert/strict"; +import {spawn} from "node:child_process"; +import {once} from "node:events"; +import {existsSync} from "node:fs"; +import {resolve} from "node:path"; +import {createInterface} from "node:readline"; import { ChatApiError, @@ -15,6 +20,192 @@ const acceptedResponse = { }; const originalFetch = globalThis.fetch; +type FetchTrace = { + body: string; + responseBody: string; + status: number; +}; + +function parseObject(line: string): Record { + const value: unknown = JSON.parse(line); + if (value === null || typeof value !== "object" || Array.isArray(value)) { + throw new TypeError("expected a JSON object"); + } + return value; +} + +function requiredString( + value: Record, + key: string, +): string { + const field = value[key]; + if (typeof field !== "string") { + throw new TypeError(`${key} must be a string`); + } + return field; +} + +async function readRequiredLine( + lines: AsyncIterableIterator, +): Promise { + const line = await lines.next(); + if (line.done) { + throw new Error("acceptance HTTP fixture exited before returning a result"); + } + return line.value; +} + +function fetchInputUrl(input: string | URL | Request): string { + if (typeof input === "string") return input; + if (input instanceof URL) return input.toString(); + return input.url; +} + +async function expectChatApiError( + operation: () => Promise, + expectedStatus: number, +): Promise { + try { + await operation(); + } catch (error) { + assert.ok(error instanceof ChatApiError); + assert.equal(error.payload.http_status, expectedStatus); + return error; + } + throw new Error(`expected ChatApiError with HTTP ${expectedStatus}`); +} + +async function runHttpRecoveryScenario( + scenario: "before_transcript" | "after_queued", +): Promise { + const dashboardRoot = process.cwd(); + const repositoryRoot = resolve(dashboardRoot, "../../.."); + const fixturePath = resolve( + dashboardRoot, + "smoke/chat-turn-acceptance-http-fixture.py", + ); + const repositoryPython = resolve(repositoryRoot, ".venv/bin/python"); + const python = process.env.LOOPX_PYTHON + ?? (existsSync(repositoryPython) ? repositoryPython : "python3"); + const child = spawn( + python, + ["-u", fixturePath, scenario], + { + cwd: repositoryRoot, + stdio: ["pipe", "pipe", "pipe"], + }, + ); + const exited = once(child, "exit"); + let stderr = ""; + child.stderr.setEncoding("utf8"); + child.stderr.on("data", (chunk) => { + stderr += String(chunk); + }); + const output = createInterface({input: child.stdout}); + const lines = output[Symbol.asyncIterator](); + const traces: FetchTrace[] = []; + + try { + const fixture = parseObject(await readRequiredLine(lines)); + const origin = requiredString(fixture, "origin"); + const sessionId = requiredString(fixture, "session_id"); + const mismatchedSessionId = requiredString( + fixture, + "mismatched_session_id", + ); + globalThis.fetch = async (input, init) => { + const response = await originalFetch( + new URL(fetchInputUrl(input), `${origin}/`), + init, + ); + traces.push({ + body: String(init?.body ?? ""), + responseBody: await response.clone().text(), + status: response.status, + }); + return response; + }; + + const accepted = await acceptChatTurn( + sessionId, + "recover this request", + "recoverable-request", + ); + assert.equal(accepted.created, false); + assert.deepEqual( + traces.map(({status}) => status), + [503, 202], + ); + assert.deepEqual( + traces.map(({body}) => parseObject(body).client_turn_id), + ["recoverable-request", "recoverable-request"], + ); + const unavailable = parseObject(traces[0].responseBody); + assert.equal( + unavailable.error_code, + "chat_turn_acceptance_unavailable", + ); + assert.equal(unavailable.turn_replay_safe, true); + assert.equal( + traces[0].responseBody.includes("private"), + false, + ); + + traces.length = 0; + await expectChatApiError( + () => acceptChatTurn(sessionId, "", "invalid-request"), + 400, + ); + assert.deepEqual(traces.map(({status}) => status), [400]); + + traces.length = 0; + await expectChatApiError( + () => acceptChatTurn( + sessionId, + "conflicting request", + "conflicting-request", + ), + 409, + ); + assert.deepEqual(traces.map(({status}) => status), [409]); + + traces.length = 0; + const mismatch = await expectChatApiError( + () => acceptChatTurn( + mismatchedSessionId, + "must not write", + "home-mismatch-request", + ), + 424, + ); + assert.equal(mismatch.payload.error_code, "codex_home_mismatch"); + assert.deepEqual(traces.map(({status}) => status), [424]); + + child.stdin.end("inspect\n"); + const summary = parseObject(await readRequiredLine(lines)); + assert.deepEqual(summary, { + fault_injected: true, + turn_count: 1, + user_message_count: 1, + queued_event_count: 1, + dispatch_count: 1, + acceptance_capsule_present: false, + mismatched_turn_count: 0, + mismatched_message_count: 0, + mismatched_session_status: "ready", + }); + const [exitCode] = await exited; + assert.equal(exitCode, 0, stderr); + } finally { + globalThis.fetch = originalFetch; + output.close(); + if (child.exitCode === null) { + child.kill(); + await exited; + } + } +} + try { const bodies: string[] = []; globalThis.fetch = async (_input, init) => { @@ -150,6 +341,9 @@ try { ), ); assert.equal(conflictCalls, 1); + + await runHttpRecoveryScenario("before_transcript"); + await runHttpRecoveryScenario("after_queued"); } finally { globalThis.fetch = originalFetch; } diff --git a/tests/test_python_ci_workflow.py b/tests/test_python_ci_workflow.py index a036debc6b..0e5504a947 100644 --- a/tests/test_python_ci_workflow.py +++ b/tests/test_python_ci_workflow.py @@ -209,8 +209,13 @@ def test_presentation_exemption_retains_real_frontend_checks_and_force_full() -> assert "name: chat-bundle-${{ github.sha }}" in job producer = WORKFLOW.split(" chat-bundle:\n", 1)[1].split(" kernel-static-checks:\n", 1)[0] assert "npm run smoke:personal-workspace-packaged" in producer + assert "npm run smoke:chat-turn-acceptance-retry" in producer assert "npm run smoke:chat-upgrade" in producer - assert producer.index("npm run smoke:personal-workspace-packaged") < producer.index("actions/upload-artifact") + assert ( + producer.index("npm run smoke:personal-workspace-packaged") + < producer.index("npm run smoke:chat-turn-acceptance-retry") + < producer.index("actions/upload-artifact") + ) assert "scripts/chat_bundle.py verify --source" in job assert "status --short --untracked-files=all -- loopx/web/chat" not in job assert "continue-on-error" not in job From 52cb4cccca874b159172a6f16db94a5719585eb7 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 27 Sep 2026 20:54:44 +0800 Subject: [PATCH 5/6] fix(test): reuse shared Python discovery Signed-off-by: duanjialing.777 --- .../smoke/chat-turn-acceptance-retry-smoke.ts | 6 ++---- .../control_plane_ts/test_python_runtime.test.ts | 15 +++++++++++++-- 2 files changed, 15 insertions(+), 6 deletions(-) diff --git a/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts index 6e3c69bd69..0d80bd54ae 100644 --- a/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts +++ b/apps/presentation/dashboard/smoke/chat-turn-acceptance-retry-smoke.ts @@ -1,10 +1,10 @@ import assert from "node:assert/strict"; import {spawn} from "node:child_process"; import {once} from "node:events"; -import {existsSync} from "node:fs"; import {resolve} from "node:path"; import {createInterface} from "node:readline"; +import {resolveTestPython} from "../../../../scripts/test-python.mjs"; import { ChatApiError, acceptChatTurn, @@ -84,9 +84,7 @@ async function runHttpRecoveryScenario( dashboardRoot, "smoke/chat-turn-acceptance-http-fixture.py", ); - const repositoryPython = resolve(repositoryRoot, ".venv/bin/python"); - const python = process.env.LOOPX_PYTHON - ?? (existsSync(repositoryPython) ? repositoryPython : "python3"); + const python = resolveTestPython({repoRoot: repositoryRoot}); const child = spawn( python, ["-u", fixturePath, scenario], diff --git a/tests/control_plane_ts/test_python_runtime.test.ts b/tests/control_plane_ts/test_python_runtime.test.ts index 4048bf153d..1806bade06 100644 --- a/tests/control_plane_ts/test_python_runtime.test.ts +++ b/tests/control_plane_ts/test_python_runtime.test.ts @@ -30,7 +30,10 @@ test("an invalid explicit test Python never falls back silently", t => { /LOOPX_TEST_PYTHON does not resolve to Python 3\.11\+/); const current = resolveTestPython(); assert.equal(resolveTestPython({ env: { - ...process.env, LOOPX_TEST_PYTHON: current, LOOPX_PYTHON_BIN: missing, + ...process.env, + LOOPX_TEST_PYTHON: current, + LOOPX_PYTHON_BIN: missing, + LOOPX_PYTHON: missing, } }), current, "the test override takes precedence over legacy browser overrides"); }); @@ -90,6 +93,7 @@ test("test and browser smokes may not introduce bare python or python3 subproces .map(match => new RegExp(`\\b${match[1]}\\s*\\(\\s*["']python3?["']`)) .some(pattern => pattern.test(source)); const fallback = /(?:\?\?|\|\|)\s*["']python3?["']/; + const conditionalFallback = /\?\s*[^:;\n]+\s*:\s*["']python3?["']/; const assigned = /\b(?:const|let)\s+\w+\s*=\s*["']python3?["']/; const bare = JSON.stringify("python3"); const barePython = JSON.stringify("python"); @@ -104,6 +108,12 @@ test("test and browser smokes may not introduce bare python or python3 subproces assert.ok(fallback.test(`process.env.LOOPX_TEST_PYTHON ?? ${barePython}`)); assert.ok(fallback.test(`process.env.NEW_TEST_PYTHON || ${bare}`)); assert.ok(fallback.test(`process.env.NEW_TEST_PYTHON || ${barePython}`)); + assert.ok(conditionalFallback.test( + `process.env.LOOPX_TEST_PYTHON ?? (existsSync(repositoryPython) ? repositoryPython : ${bare})`, + )); + assert.ok(conditionalFallback.test( + `process.env.LOOPX_TEST_PYTHON ?? (existsSync(repositoryPython) ? repositoryPython : ${barePython})`, + )); assert.ok(assigned.test(`const PYTHON = ${bare}`)); assert.ok(assigned.test(`const PYTHON = ${barePython}`)); assert.ok(assigned.test(`const testInterpreter = ${bare}`)); @@ -123,7 +133,8 @@ test("test and browser smokes may not introduce bare python or python3 subproces else if (/\.(?:cjs|js|mjs|mts|ts)$/.test(entry.name)) { const source = readFileSync(join(root, path), "utf8"); if (direct.test(source) || promisified.test(source) || aliasedLaunch(source) - || fallback.test(source) || assigned.test(source)) offenders.push(path); + || fallback.test(source) || conditionalFallback.test(source) + || assigned.test(source)) offenders.push(path); } } } From 275a2a8c9be82d477475c224a27e8760adf1c4ca Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 28 Sep 2026 02:40:29 +0800 Subject: [PATCH 6/6] fix(chat): keep acceptance bridge within module budget Signed-off-by: duanjialing.777 --- loopx/chat_store.py | 2 +- loopx/{ => control_plane}/chat_turn_acceptance.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) rename loopx/{ => control_plane}/chat_turn_acceptance.py (99%) diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 3ad939aa75..4b5a650d5c 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -19,7 +19,7 @@ ) from .chat_event_cache import ChatEventCache from .chat_ingress import ChatIngressStore -from .chat_turn_acceptance import ( +from .control_plane.chat_turn_acceptance import ( CHAT_TURN_ACCEPTANCE_CAPSULE_SCHEMA, AcceptedManagedTurn, plan_managed_turn_acceptance, diff --git a/loopx/chat_turn_acceptance.py b/loopx/control_plane/chat_turn_acceptance.py similarity index 99% rename from loopx/chat_turn_acceptance.py rename to loopx/control_plane/chat_turn_acceptance.py index e15df8b7cc..cfdacbd02d 100644 --- a/loopx/chat_turn_acceptance.py +++ b/loopx/control_plane/chat_turn_acceptance.py @@ -7,7 +7,7 @@ import json from typing import Any -from .control_plane.effect_runtime import effect_runtime_result +from .effect_runtime import effect_runtime_result CHAT_TURN_ACCEPTANCE_REQUEST_SCHEMA = "loopx_chat_turn_acceptance_request_v0"