From 70e1ee7df10cb9cba411ddf62efb8fb4b91336cb Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 9 Sep 2026 12:03:54 +0800 Subject: [PATCH 1/8] fix(chat): preserve attached session ownership Signed-off-by: duanjialing.777 --- loopx/chat_runtime.py | 39 ++++++++++++--- tests/test_attached_session_broker.py | 65 ++++++++++++++++++++++++ tests/test_chat_session_active_turn.py | 68 ++++++++++++++++++++++++++ 3 files changed, 164 insertions(+), 8 deletions(-) diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 6b37d05804..8ef4b84963 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -409,10 +409,6 @@ def open_session( agent_goal_id: str | None = None, ) -> tuple[dict[str, Any], bool]: capability = next((item for item in self.capabilities() if item["agent_id"] == agent_id), None) - if capability is None: - raise ValueError(f"unknown Agent endpoint: {agent_id}") - if not capability["available"]: - raise ValueError(f"Agent endpoint is unavailable: {agent_id}") if mode not in {"resume_latest", "new"}: raise ValueError("mode must be resume_latest or new") selected_channel = channel_id or f"goal.{goal_id}" @@ -421,15 +417,22 @@ def open_session( with self.lock: route_lock = self.session_open_locks.setdefault(route_key, threading.Lock()) with route_lock: + latest = None if mode == "resume_latest": latest = self.store.latest_session( goal_id=None if selected_channel == "manager" else goal_id, agent_id=agent_id, channel_id=selected_channel, ) - if latest is not None: - self._ensure_adapter(latest, work_dir=work_dir, objective=objective) + if latest is not None and latest.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: return latest, True + if capability is None: + raise ValueError(f"unknown Agent endpoint: {agent_id}") + if not capability["available"]: + raise ValueError(f"Agent endpoint is unavailable: {agent_id}") + if latest is not None: + self._ensure_adapter(latest, work_dir=work_dir, objective=objective) + return latest, True adapter = self._start_adapter( agent_id=agent_id, work_dir=work_dir, @@ -974,8 +977,13 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None: self._fail_turn(session_id, turn_id, exc.error_code, str(exc), status="failed", gate=exc.gate) if not adapter.healthcheck(): with self.lock: - self.adapters.pop(session_id, None) - self.store.update_session(session_id, status="stale", last_error_code="transport_disconnected") + if self.adapters.get(session_id) is adapter: + self.adapters.pop(session_id, None) + self.store.update_session( + session_id, + status="stale", + last_error_code="transport_disconnected", + ) except Exception as exc: # noqa: BLE001 - preserve compact runtime failure. event_buffer.close() if consume_interrupted(): @@ -1030,6 +1038,21 @@ def interrupt_turn(self, *, session_id: str, turn_id: str) -> dict[str, Any]: 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.", + }, + ) interrupting = self.store.update_turn( session_id, turn_id, diff --git a/tests/test_attached_session_broker.py b/tests/test_attached_session_broker.py index 29b429d2cf..3012d5d4dd 100644 --- a/tests/test_attached_session_broker.py +++ b/tests/test_attached_session_broker.py @@ -13,6 +13,7 @@ claim_attached_agent_turn, complete_attached_agent_turn, ) +from loopx.chat_agent import CodexChatAgentError from loopx.chat_runtime import ChatRuntimeController from loopx.chat_server import ChatRequestHandler from loopx.chat_store import CHAT_SESSION_MODE_ATTACHED, ChatSessionStore @@ -362,6 +363,70 @@ def test_web_and_lark_share_ordered_attached_session_without_spawning( ] +def test_resume_latest_reuses_attached_session_without_local_adapter( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session = _bind(store)["session"] + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + monkeypatch.setattr( + runtime, + "_start_adapter", + lambda **_kwargs: pytest.fail("attached Session must not start a local adapter"), + ) + + resumed, reused = runtime.open_session( + goal_id=GOAL_ID, + agent_id=AGENT_ID, + work_dir=tmp_path, + objective="sample objective", + mode="resume_latest", + ) + + assert reused is True + assert resumed["session_id"] == session["session_id"] + assert runtime.adapters == {} + + +def test_attached_claimed_turn_interrupt_fails_closed(tmp_path: Path) -> None: + store = ChatSessionStore(tmp_path) + session_id = str(_bind(store)["session"]["session_id"]) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + turn, _created = store.create_queued_turn( + session_id, + client_turn_id="claimed-interrupt", + message="keep host ownership", + ) + claim_attached_agent_turn( + store=store, + session_id=session_id, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="claimed-interrupt", + ) + + with pytest.raises(CodexChatAgentError) as raised: + runtime.interrupt_turn(session_id=session_id, turn_id=str(turn["turn_id"])) + + assert raised.value.error_code == "attached_session_interrupt_unavailable" + current_turn = store.load_turn(session_id, str(turn["turn_id"])) + current_session = store.load_session(session_id) + assert current_turn is not None and current_turn["status"] == "running" + assert current_session is not None and current_session["status"] == "busy" + assert current_session["active_turn_id"] == turn["turn_id"] + assert ( + claim_attached_agent_turn( + store=store, + session_id=session_id, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="next-claim", + )["claimed"] + is False + ) + + def test_attached_completion_uses_canonical_response_and_terminal_events( tmp_path: Path, ) -> None: diff --git a/tests/test_chat_session_active_turn.py b/tests/test_chat_session_active_turn.py index 938759c1fa..e6416101bd 100644 --- a/tests/test_chat_session_active_turn.py +++ b/tests/test_chat_session_active_turn.py @@ -6,6 +6,7 @@ import pytest import loopx.chat_store as chat_store +from loopx.chat_agent import CodexChatAgentError from loopx.chat_runtime import ChatRuntimeController from loopx.chat_store import ( SESSION_QUEUE_MAX_PENDING, @@ -56,6 +57,19 @@ def close_session(self) -> None: self.closed = True +class _FailingChatAdapter(_HealthyChatAdapter): + def start_turn(self, message: str, event_sink) -> dict[str, object]: + del message, event_sink + raise CodexChatAgentError( + "transport disconnected", + error_code="transport_disconnected", + gate={"kind": "host_tool_gate"}, + ) + + def healthcheck(self) -> bool: + return False + + def _slow_new_turn_writes(monkeypatch) -> set[str]: original_atomic_write = chat_store._atomic_write_json turn_write_threads: set[str] = set() @@ -1222,6 +1236,60 @@ def test_managed_close_rejects_pending_queue_without_stranding_it( assert store.load_turn(session_id, str(queued["turn_id"]))["status"] == "queued" # type: ignore[index] +def test_failed_old_turn_does_not_remove_replacement_adapter( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session = 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 = str(session["session_id"]) + first_turn, _created = store.create_turn( + session_id, + client_turn_id="failed-old-turn", + message="fail on the old adapter", + ) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + old_adapter = _FailingChatAdapter() + replacement_adapter = _HealthyChatAdapter() + runtime.adapters[session_id] = old_adapter + original_fail_turn = runtime._fail_turn + replacement_turn: dict[str, object] = {} + + def fail_then_replace(*args: object, **kwargs: object) -> None: + original_fail_turn(*args, **kwargs) # type: ignore[arg-type] + with runtime.lock: + runtime.adapters[session_id] = replacement_adapter # type: ignore[assignment] + replacement_turn.update( + store.create_turn( + session_id, + client_turn_id="replacement-turn", + message="continue on the replacement adapter", + )[0] + ) + + monkeypatch.setattr(runtime, "_fail_turn", fail_then_replace) + + runtime._run_turn( + session_id=session_id, + turn_id=str(first_turn["turn_id"]), + message="fail on the old adapter", + attachments=[], + adapter=old_adapter, + ) + + current = store.load_session(session_id) + assert runtime.adapters[session_id] is replacement_adapter + assert current is not None and current["status"] == "busy" + assert current["active_turn_id"] == replacement_turn["turn_id"] + + @pytest.mark.parametrize("terminal_status", sorted(TERMINAL_TURN_STATES)) def test_managed_close_clears_terminal_active_turn_left_by_crash( tmp_path: Path, From e91885654df8e6be3e96d26f3c3222f19e8c8e80 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Thu, 10 Sep 2026 13:13:11 +0800 Subject: [PATCH 2/8] fix(coordination): stabilize fence read errors Signed-off-by: duanjialing.777 --- loopx/control_plane/coordination/legacy_writer_fence.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/loopx/control_plane/coordination/legacy_writer_fence.ts b/loopx/control_plane/coordination/legacy_writer_fence.ts index 8a8f4507e1..81af31fbf1 100644 --- a/loopx/control_plane/coordination/legacy_writer_fence.ts +++ b/loopx/control_plane/coordination/legacy_writer_fence.ts @@ -142,10 +142,12 @@ export async function loadLegacyCoordinationWriterFence( return { status: "loaded", fence }; } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return { status: "missing" }; + const path = (error as NodeJS.ErrnoException).path; + const reason = error instanceof Error ? error.message : "legacy writer fence read failed"; return { status: "failed", reason_code: "legacy_writer_fence_read_failed", - reason: error instanceof Error ? error.message : "legacy writer fence read failed", + reason: typeof path === "string" ? reason.replace(` '${path}'`, "") : reason, }; } } From ff8c0de475865986e99d2a404b072a2372d1a045 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 11:20:16 +0800 Subject: [PATCH 3/8] test(chat): cover attached interrupt endpoint error code Signed-off-by: duanjialing.777 --- tests/test_attached_session_broker.py | 48 +++++++++++++++++++++++++++ 1 file changed, 48 insertions(+) diff --git a/tests/test_attached_session_broker.py b/tests/test_attached_session_broker.py index 3012d5d4dd..64bf65de9d 100644 --- a/tests/test_attached_session_broker.py +++ b/tests/test_attached_session_broker.py @@ -427,6 +427,54 @@ def test_attached_claimed_turn_interrupt_fails_closed(tmp_path: Path) -> None: ) +def test_attached_interrupt_endpoint_preserves_typed_failure(tmp_path: Path) -> None: + store = ChatSessionStore(tmp_path) + session_id = str(_bind(store)["session"]["session_id"]) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + turn, _created = store.create_queued_turn( + session_id, + client_turn_id="endpoint-interrupt", + message="keep host ownership", + ) + turn_id = str(turn["turn_id"]) + claim_attached_agent_turn( + store=store, + session_id=session_id, + host_surface=HOST_SURFACE, + host_session_id=HOST_SESSION_ID, + claim_id="endpoint-interrupt", + ) + responses: list[dict[str, object]] = [] + + class Handler: + server = SimpleNamespace(runtime_controller=runtime) + + def _send_error(self, message: str, **kwargs: object) -> None: + responses.append({"error": message, **kwargs}) + + def _send_json(self, *_args: object, **_kwargs: object) -> None: + raise AssertionError("interrupt should fail closed") + + ChatRequestHandler._interrupt_turn(Handler(), session_id, turn_id) # type: ignore[arg-type] + + assert responses == [ + { + "error": "The attached host does not expose interrupt control to LoopX Chat.", + "status": 424, + "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.", + }, + } + ] + assert store.load_turn(session_id, turn_id)["status"] == "running" # type: ignore[index] + session = store.load_session(session_id) + assert session is not None and session["status"] == "busy" + assert session["active_turn_id"] == turn_id + + def test_attached_completion_uses_canonical_response_and_terminal_events( tmp_path: Path, ) -> None: From 4246bc1294f3995f37f026661bb04153e28c6b83 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 11:22:41 +0800 Subject: [PATCH 4/8] fix(chat): preserve interrupt error codes Signed-off-by: duanjialing.777 --- loopx/chat_server.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/chat_server.py b/loopx/chat_server.py index d389fca8d1..1dc55e6cf9 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -831,7 +831,7 @@ def _interrupt_turn(self, session_id: str, turn_id: str) -> None: self._send_error("chat turn was not found", status=404) return except CodexChatAgentError as exc: - self._send_error(str(exc), status=424, gate=exc.gate) + self._send_error(str(exc), status=424, error_code=exc.error_code, gate=exc.gate) return self._send_json( { From 5eb8d4635f399b2a09995db12880cae3fa736d76 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 22:24:43 +0800 Subject: [PATCH 5/8] test(dashboard): disambiguate recovery pending state Signed-off-by: duanjialing.777 --- examples/personal-workspace-browser-smoke.mjs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/examples/personal-workspace-browser-smoke.mjs b/examples/personal-workspace-browser-smoke.mjs index 15d1f6b6c6..8ffd88336d 100644 --- a/examples/personal-workspace-browser-smoke.mjs +++ b/examples/personal-workspace-browser-smoke.mjs @@ -2892,14 +2892,14 @@ async function main() { await page.locator(".personal-goal-link").first().click(); await goalNavigation.getByRole("button", { name: "Chat" }).click(); await page.getByText("保持运行,用于验证刷新恢复。").waitFor({ state: "visible", timeout: 10_000 }); - await page.getByText("正在整理…").waitFor({ state: "hidden", timeout: 10_000 }); + await page.getByText("正在整理…", { exact: true }).waitFor({ state: "hidden", timeout: 10_000 }); const recovered = page.__loopxRuntime.sessions.get(recoveryTurn.sessionId); if (recovered?.active_turn_id !== null && recovered?.active_turn_id !== recoveryTurn.turnId) { throw new Error("Recovered Session points at a different active Turn"); } pass(6, "Reload restored visible Goal history and resumed the active Turn SSE stream."); } catch (error) { - fail(6, "Reload did not restore the active Goal conversation and reconnect its active Turn within 10 seconds."); + fail(6, `Reload did not restore the active Goal conversation and reconnect its active Turn within 10 seconds: ${error.message}`); await page.screenshot({ path: resolve(outputDir, "refresh-recovery-failed.png"), fullPage: true, animations: "disabled" }); observations.push(`Refresh recovery failure: ${error.message}`); } From 84a76e125b9b8681bd4ea164ee459c965237f879 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 23:08:54 +0800 Subject: [PATCH 6/8] test(dashboard): replay completed mock turn after reload Signed-off-by: duanjialing.777 (cherry picked from commit 31209395b98bc14f84f5b4a82d246e6ea2e4a085) Signed-off-by: duanjialing.777 --- examples/personal-workspace-browser-smoke.mjs | 22 ++++++++++--------- 1 file changed, 12 insertions(+), 10 deletions(-) diff --git a/examples/personal-workspace-browser-smoke.mjs b/examples/personal-workspace-browser-smoke.mjs index 9366ddf261..ee9350e3f1 100644 --- a/examples/personal-workspace-browser-smoke.mjs +++ b/examples/personal-workspace-browser-smoke.mjs @@ -847,18 +847,19 @@ async function installApi(page, { goalSubagentConfigurationEnabled = true } = {} const answer = "已沿用当前 Goal 与 Agent Session。接下来会先核对状态,再继续推进。"; await new Promise((resolveWait) => setTimeout(resolveWait, /(中断控制|刷新恢复)/u.test(turnMessages.get(turnId) ?? "") ? 5000 : 1200)); const activeSession = sessions.get(sessionId); + const visible = messages.get(sessionId) ?? []; + const completed = visible.some((message) => message.message_id === `${turnId}-assistant`); + const event = (id, kind, payload) => `id: ${id}\nevent: ${kind}\ndata: ${JSON.stringify({ event_id: id, sequence: Number(id), kind, created_at: "2026-08-13T01:00:02Z", payload })}\n\n`; if (!activeSession || activeSession.active_turn_id !== turnId) { - await route.fulfill({ contentType: "text/event-stream", body: "", status: 200 }); + await route.fulfill({ contentType: "text/event-stream", body: completed ? event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], gate: null } }) : "", status: 200 }); return; } - const visible = messages.get(sessionId) ?? []; - if (!visible.some((message) => message.message_id === `${turnId}-assistant`)) { + if (!completed) { visible.push({ message_id: `${turnId}-assistant`, turn_id: turnId, role: "assistant", text: answer, created_at: "2026-08-13T01:00:02Z" }); } messages.set(sessionId, visible); - const event = (id, kind, payload) => `id: ${id}\nevent: ${kind}\ndata: ${JSON.stringify({ event_id: id, sequence: Number(id), kind, created_at: "2026-08-13T01:00:02Z", payload })}\n\n`; - await route.fulfill({ contentType: "text/event-stream", body: event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], gate: null } }), status: 200 }); sessions.set(sessionId, { ...activeSession, active_turn_id: null, status: "ready", updated_at: "2026-08-13T01:00:02Z" }); + await route.fulfill({ contentType: "text/event-stream", body: event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], gate: null } }), status: 200 }); return; } if (url.pathname === "/api/chat/goals/contexts") { @@ -1090,20 +1091,21 @@ async function installApi(page, { goalSubagentConfigurationEnabled = true } = {} : "已沿用当前 Goal 与 Agent Session。接下来会先核对状态,再继续推进。"; await new Promise((resolveWait) => setTimeout(resolveWait, /(中断控制|刷新恢复)/u.test(operatorMessage) ? 5000 : 1200)); const activeSession = sessions.get(sessionId); + const visible = messages.get(sessionId) ?? []; + const completed = visible.some((message) => message.message_id === `${turnId}-assistant`); + const event = (id, kind, payload) => `id: ${id}\nevent: ${kind}\ndata: ${JSON.stringify({ event_id: id, sequence: Number(id), kind, created_at: "2026-08-13T01:00:02Z", payload })}\n\n`; if (!activeSession || activeSession.active_turn_id !== turnId) { - await route.fulfill({ contentType: "text/event-stream", body: "", status: 200 }); + await route.fulfill({ contentType: "text/event-stream", body: completed ? event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], protected_action: protectedAction, gate: null } }) : "", status: 200 }); return; } if (sessionId && messages.has(sessionId)) { - const visible = messages.get(sessionId); - if (!visible.some((message) => message.message_id === `${turnId}-assistant`)) { + if (!completed) { visible.push({ message_id: `${turnId}-assistant`, turn_id: turnId, role: "assistant", text: answer, created_at: "2026-08-13T01:00:02Z" }); } } - const event = (id, kind, payload) => `id: ${id}\nevent: ${kind}\ndata: ${JSON.stringify({ event_id: id, sequence: Number(id), kind, created_at: "2026-08-13T01:00:02Z", payload })}\n\n`; - await route.fulfill({ contentType: "text/event-stream", body: event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], protected_action: protectedAction, gate: null } }), status: 200 }); const current = sessions.get(sessionId); if (current?.active_turn_id === turnId) sessions.set(sessionId, { ...current, active_turn_id: null, status: "ready", updated_at: "2026-08-13T01:00:02Z" }); + await route.fulfill({ contentType: "text/event-stream", body: event("1", "assistant.delta", { text: answer }) + event("2", "turn.completed", { response: { schema_version: "loopx_chat_agent_response_v0", message: answer, proposals: [], protected_action: protectedAction, gate: null } }), status: 200 }); }); await page.route("**/api/actions?**", async (route) => { const url = new URL(route.request().url()); From 9908571a6a8513bbb592e3a68732561d1e10fd6d Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 23:35:02 +0800 Subject: [PATCH 7/8] test(dashboard): target exact recovery pending state Signed-off-by: duanjialing.777 (cherry picked from commit 8e2262a63063b4684fdd86c75c741d6d6ad4a9c4) Signed-off-by: duanjialing.777 --- examples/personal-workspace-browser-smoke.mjs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/examples/personal-workspace-browser-smoke.mjs b/examples/personal-workspace-browser-smoke.mjs index 4490c0de68..c814a67619 100644 --- a/examples/personal-workspace-browser-smoke.mjs +++ b/examples/personal-workspace-browser-smoke.mjs @@ -2918,7 +2918,7 @@ async function main() { await page.locator(".personal-goal-link").first().click(); await goalNavigation.getByRole("button", { name: "Chat" }).click(); await page.getByText("保持运行,用于验证刷新恢复。").waitFor({ state: "visible", timeout: 10_000 }); - await page.getByText("正在整理…").waitFor({ state: "hidden", timeout: 10_000 }); + await page.getByText("正在整理…", { exact: true }).waitFor({ state: "hidden", timeout: 10_000 }); const recovered = page.__loopxRuntime.sessions.get(recoveryTurn.sessionId); if (recovered?.active_turn_id !== null && recovered?.active_turn_id !== recoveryTurn.turnId) { throw new Error("Recovered Session points at a different active Turn"); From 596716c8dbab7e2fc117938f4662fa2f2411d3ec Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Fri, 11 Sep 2026 23:44:28 +0800 Subject: [PATCH 8/8] fix(todo): type effect decision results Signed-off-by: duanjialing.777 (cherry picked from commit f4f2b9bbcc51f28cf4a3dc216ba19a4efb3fb461) Signed-off-by: duanjialing.777 --- loopx/control_plane/todos/decision_scope.py | 30 ++++++++++++++------- 1 file changed, 20 insertions(+), 10 deletions(-) diff --git a/loopx/control_plane/todos/decision_scope.py b/loopx/control_plane/todos/decision_scope.py index 07cf551777..ebc763a9fb 100644 --- a/loopx/control_plane/todos/decision_scope.py +++ b/loopx/control_plane/todos/decision_scope.py @@ -1,7 +1,7 @@ """Legacy input codec for the single typed decision-dependency rule owner.""" from __future__ import annotations -from typing import Any +from typing import Any, cast from ..effect_runtime import effect_runtime_result from .contract import ( @@ -90,8 +90,10 @@ def standing_decision_authority_for_agent(authority: dict[str, Any] | None, *, agent_id: str | None) -> dict[str, Any] | None: if _authority(authority) is None: return None - return _evaluate("standing", authority=_authority(authority), - agent_id=normalize_todo_claimed_by(agent_id)) + return cast(dict[str, Any] | None, _evaluate( + "standing", authority=_authority(authority), + agent_id=normalize_todo_claimed_by(agent_id), + )) def build_required_decision_scope_consistency( @@ -101,14 +103,14 @@ def build_required_decision_scope_consistency( user_source_items: list[dict[str, Any]] | None = None, standing_decision_authority: dict[str, Any] | None = None, ) -> dict[str, Any]: - return _evaluate("consistency", + return cast(dict[str, Any], _evaluate("consistency", agent_items=_source(agent_todo_summary, agent_source_items, _AGENT_SUMMARY_ITEM_KEYS), user_items=_source(user_todo_summary, user_source_items, _USER_SUMMARY_ITEM_KEYS), agent_id=normalize_todo_claimed_by(agent_id), registered_agents=sorted({value for raw in registered_agent_ids or [] if (value := normalize_todo_claimed_by(raw))}), standing_authority=_authority(standing_decision_authority), - ) + )) def build_required_decision_scope_repair_hint( @@ -193,26 +195,34 @@ def decision_scope_covers(gate_scope: Any, required_scope: Any) -> bool: required = normalize_todo_decision_scope(required_scope) if not gate or not required: return False - return _evaluate("covers", gate_scope=gate, required_scope=required) + return cast(bool, _evaluate("covers", gate_scope=gate, required_scope=required)) def decision_scope_gate_relation(gate: dict[str, Any], agent_item: dict[str, Any]) -> dict[str, Any] | None: - return _evaluate("scope_relation", gate=_facts(gate), item=_facts(agent_item)) + return cast(dict[str, Any] | None, _evaluate( + "scope_relation", gate=_facts(gate), item=_facts(agent_item), + )) def exact_todo_gate_relation(gate: dict[str, Any], agent_item: dict[str, Any]) -> dict[str, Any] | None: - return _evaluate("exact_relation", gate=_facts(gate), item=_facts(agent_item)) + return cast(dict[str, Any] | None, _evaluate( + "exact_relation", gate=_facts(gate), item=_facts(agent_item), + )) def todo_gate_relation(gate: dict[str, Any], agent_item: dict[str, Any]) -> dict[str, Any] | None: - return _evaluate("relation", gate=_facts(gate), item=_facts(agent_item)) + return cast(dict[str, Any] | None, _evaluate( + "relation", gate=_facts(gate), item=_facts(agent_item), + )) def todo_gate_relations(gates: list[dict[str, Any]], items: list[dict[str, Any]]) -> list[list[dict[str, Any] | None]]: """Evaluate a consumer's candidate set in one RPC, retaining positional identity.""" if not gates or not items: return [[] for _ in gates] - return _evaluate("relations", gates=[_facts(gate) for gate in gates], items=[_facts(item) for item in items]) + return cast(list[list[dict[str, Any] | None]], _evaluate( + "relations", gates=[_facts(gate) for gate in gates], items=[_facts(item) for item in items], + )) def todo_gate_relation_blocks_agent(relation: dict[str, Any] | None) -> bool: