From 0aec8a7b38b1feba2fd18144c466c23814bfe43c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 03:09:50 +0800 Subject: [PATCH] fix(chat): settle accepted turns when runtime preparation fails Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../capable-manager-semantic-handoff-v0.md | 23 +++ .../rfcs/loopx-overall-roadmap-v0.md | 2 +- loopx/chat_runtime.py | 58 +++++-- loopx/extensions/lark/manager_context.py | 2 + .../test_lark_goal_topic_runtime.py | 2 + tests/test_chat_queue_preparation.py | 159 ++++++++++++++++++ 6 files changed, 233 insertions(+), 13 deletions(-) create mode 100644 tests/test_chat_queue_preparation.py diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md index 32cb861e6e..be63744b68 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md @@ -595,6 +595,29 @@ Target an ingress receipt within two seconds on a healthy local service, indepen Use existing service recovery and receipt pumps. No manager-specific business automation for each kind of request. Expose configuration and failures through the existing CLI, capability settings and manager conversation. Troubleshooting distinguishes model failure, tool/policy denial, state conflict, unreachable receiver and transport formatting/delivery failure. +**Accepted queue preparation failures (S1/S10, A12/A22/A23):** an accepted +request owns a terminal outcome even before an adapter starts. A missing runtime +asset, invalid workspace or failed session restoration must settle the affected +queued Turn through the shared Chat lifecycle and release its claim. Waiting for +an answer must observe that durable failure promptly, rather than wait for the +model timeout while leaving runnable work behind. Keep the original typed +provider failure where available; unexpected local preparation errors use +`runtime_unavailable`, with private diagnostics retained locally. Cancellation +and an already terminal result win over a late preparation error. Restoring the +runtime must not replay a failed request; the same ingress identity returns the +same failure, while a fresh explicit request can run after repair. + +The bounded Python queue repair uses the existing store's fenced failure and +claim-release operations for all queue callers; Lark only translates the typed +outcome. It does not create a separate manager scheduler or new TS authority. +The TS turn-driver migration must preserve this pre-dispatch failure matrix +alongside accepted-request recovery. Validate with a removed-release fixture, +multiple queued requests, a stop race, same-identity redelivery and a fresh +request after recovery. These qualify the preparation boundary, not successful +owner selection, receiver adoption or the complete A24 journey. Operational +recovery must also verify the service's actual installed release: a healthy HTTP +listener alone does not prove its lazy-loaded runtime assets still exist. + ## 11. Normative delivery plan Implement coherent end-to-end slices, not one PR per incidental field. The manager engineering owner maintains canonical Todos and a private incident-to-acceptance map; PRs link this RFC milestone and acceptance IDs. Public progress updates contain only safe results. Milestone completion requires current deployment evidence, not merged PR count. diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 2692441143..2a5ecbca60 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -45,7 +45,7 @@ P0 blocks correctness or continuity in the current user journey. P1 enables repe | **S7 Budget, scheduling and fleet scale · P0 observation/P1–P2 expansion** | Quota/scheduler and partial usage aggregates exist; full provider cost, distributed reservations and hundred-Agent concurrency need evidence | Separate configured budget, admission, consumption and estimates; unknown is not zero and replay cannot double-charge. R7 pagination/bounded summaries and [complete-history transport](typescript-control-plane-migration-v0.md), including refresh/replay/single-debit evidence beyond the RPC limit; provider/host limits, fairness, backpressure, event wake and isolation; report registration/activity/throughput and cost per accepted outcome separately | | **S8 Capabilities, extensions and domain integration · P1/P2** | Capability catalog, extension lifecycle, hooks, engineering/research/content/office capabilities and computer-use contracts exist | First exercise the shared control plane with existing issue-fix/PR-review and material/research callers. Every provider has readiness/version/permissions/default-off/uninstall/rollback/isolation and real-entry evidence. New domain effects start with one simulated operation, not a marketplace or workflow DSL | | **S9 Identity, authority, privacy and trust · continuous P0/P1–P2 remote** | Public/private scope, capability gates, fencing and confirmation contracts belong to existing owners | R1/R3 cover sender/audience/artifact scope and stale authority; R6 authenticates tenant/Goal/actor/host, rotation/revocation and least privilege. Qualify credential custody, untrusted tool/document inputs, dependency supply chain, audit retention/deletion and vulnerability response through real paths; roles/messages/memory mint no write authority | -| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery | +| **S10 Reliability, diagnostics and operations · P0/P1** | Recovery/canary, read-only diagnostics prototype and DSH event adapter exist; C0/C1, overhead and full operations qualification are open | Failure classification→observable state→recovery drill→regression prevention; process/storage/network/delivery failures and data growth. Accepted Chat requests must settle even when runtime preparation fails before dispatch; qualify missing runtime assets, stop races and recovery without replay under [the shared conversation operational contract](capable-manager-semantic-handoff-v0.md#10-operational-contract). Use [bounded repair lookup and targeted diagnostics](../../../skills/loopx-self-repair/references/targeted-diagnostics.md) to reduce redundant reads above the provider boundary; measure backend-specific cold/warm reads, writes and lock waits separately. Freeze SLO/RPO/RTO/capacity/retention boundaries and measure before qualification. Runbooks include upgrade, restore, stop and human takeover; test counts do not prove recovery | | **S11 Evaluation and scientific research · continuous P1/P2 research** | Benchmark toolkit, Explore, long-horizon portfolio and ten frontier-science tracks have designs/partial implementations | Pin native/passive/governed arms, model/harness/budget/task split and evaluator; report native scores, cost, failures, attention and uncertainty. Prioritize sequential evidence, continuation and stride; memory, formal kernel, curriculum/evolution, active experiments and multiscale state follow T01–T10 gates without automatic production treatment | | **S12 Release, developer experience and community governance · P0 hygiene/P1** | Install/source validation, registration, DCO/PR, test layers, contributor routes and bilingual docs exist | Qualify first work and upgrade/rollback from clean machines/release artifacts; host/OS support follows the release contract. Reduce localization/test/review effort for useful changes; preserve exact-head evidence, fixtures, compatibility, maintainer routing and contributor credit; retire duplicate protocols/stale evidence | | **S13 Adoption, ecosystem and sustainability · P1 discovery/P2 pilots** | Public adoption loop, showcases, licensing/governance and observer-first product contract exist; paid PMF is unproven | Gather independent first/repeat usage and exit reasons; reproducible cases and pilots with fixed budgets/acceptance/rollback. Retain reusable adapters/delivery guides. Account for model/compute/storage/support and maintenance costs; only repeated demand justifies commercial hosting/support/distribution decisions, with no invented SLA or open-source-term change | diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 59befe2fb7..60997503d1 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -4,6 +4,7 @@ from dataclasses import dataclass import json +import logging import os from pathlib import Path import threading @@ -1114,12 +1115,16 @@ def _drain_session_queue( with self.lock: attached = self.adapters.get(session_id) if attached is None or not attached.healthcheck(): - self._ensure_adapter( - session, - work_dir=work_dir, - objective=objective, - ) - continue + try: + self._ensure_adapter( + session, + work_dir=work_dir, + objective=objective, + ) + except Exception as exc: + self._fail_queue_preparation(session_id, active_turn_id, exc) + else: + continue active = self.store.load_turn(session_id, active_turn_id) if active and active.get("status") not in TERMINAL_TURN_STATES: self.closed.wait(0.05) @@ -1132,11 +1137,17 @@ def _drain_session_queue( refreshed = self.store.load_session(session_id) if refreshed is None: return - adapter = self._ensure_adapter( - refreshed, - work_dir=work_dir, - objective=objective, - ) + preparation_error = None + try: + adapter = self._ensure_adapter( + refreshed, + work_dir=work_dir, + objective=objective, + ) + except Exception as exc: + # An accepted queued request still owns an outcome when + # workspace preparation or adapter restoration fails. + preparation_error = exc with self.lock: self.session_queue_wakeups.discard(session_id) turn = self.store.claim_next_queued_turn(session_id) @@ -1150,6 +1161,9 @@ def _drain_session_queue( self.session_queue_threads.pop(session_id, None) return turn_id = str(turn["turn_id"]) + if preparation_error is not None: + self._fail_queue_preparation(session_id, turn_id, preparation_error) + continue with self.lock: self.turn_done_events[(session_id, turn_id)] = threading.Event() self._run_turn( @@ -1167,6 +1181,25 @@ def _drain_session_queue( self.session_queue_threads.pop(session_id, None) self.session_queue_wakeups.discard(session_id) + def _fail_queue_preparation( + self, session_id: str, turn_id: str, error: Exception, + ) -> None: + 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"}, + ) + def _run_turn( self, *, @@ -1390,12 +1423,13 @@ def _fail_turn( *, status: str, gate: dict[str, Any] | None = None, + expected_statuses: set[str] | None = None, ) -> None: completed = utc_now() failed = self.store.update_turn( session_id, turn_id, - expected_statuses={"starting", "running"}, + expected_statuses=expected_statuses if expected_statuses is not None else {"starting", "running"}, status=status, error_code=error_code, error=message, diff --git a/loopx/extensions/lark/manager_context.py b/loopx/extensions/lark/manager_context.py index 3149753c85..906fb4f53d 100644 --- a/loopx/extensions/lark/manager_context.py +++ b/loopx/extensions/lark/manager_context.py @@ -51,6 +51,8 @@ def manager_failure_reply(error: Exception) -> tuple[str, str]: "host_gate": "执行器未能完成本次调用,请检查管家执行器的具体错误", "upstream_invalid_request": "上游拒绝了请求参数,请检查所选模型与 Codex CLI、账户的兼容性", "server_restarted": "管家服务在处理过程中重启", + "runtime_unavailable": "本地运行环境初始化失败,请检查 LoopX 安装与运行配置,修复后重试", + "resume_failed": "原 Agent 会话恢复失败,请检查执行器状态,修复后重试", "manager_channel_executor_rebind_required": ( "管家的执行器已由本机设置更改,需要重新应用一次管家连接" ), diff --git a/tests/extensions/test_lark_goal_topic_runtime.py b/tests/extensions/test_lark_goal_topic_runtime.py index 0c3d1da673..f4ec57d3a9 100644 --- a/tests/extensions/test_lark_goal_topic_runtime.py +++ b/tests/extensions/test_lark_goal_topic_runtime.py @@ -2570,6 +2570,8 @@ def answer(_route, _text): @pytest.mark.parametrize("error_code,label", [ ("cyber_policy", "安全策略拦截"), ("rate_limit_exceeded", "请求频率限制"), + ("runtime_unavailable", "本地运行环境初始化失败"), + ("resume_failed", "原 Agent 会话恢复失败"), ("private-upstream-detail", "管家处理失败"), ]) @pytest.mark.parametrize("reply_ok", [True, False]) diff --git a/tests/test_chat_queue_preparation.py b/tests/test_chat_queue_preparation.py new file mode 100644 index 0000000000..a970d6ae29 --- /dev/null +++ b/tests/test_chat_queue_preparation.py @@ -0,0 +1,159 @@ +"""Accepted requests must settle even when the runtime cannot be prepared.""" + +import threading + +import pytest + +from loopx.chat_agent import CodexChatAgentError +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_store import ChatSessionStore +from loopx.extensions.lark.goal_topic_runtime import ( + LarkGoalTopicTurnFailed, answer_lark_goal_topic, +) +from loopx.extensions.lark.manager_context import manager_failure_reply + + +def session_runtime(tmp_path, *, channel="goal.public-research"): + store = ChatSessionStore(tmp_path) + session = store.create_session( + goal_id="public-research", agent_id="codex", channel_id=channel, + executor_endpoint_id="codex", adapter_kind="codex_app_server", + upstream_thread_id="public-thread", upstream_mode="chat", + ) + return store, str(session["session_id"]), ChatRuntimeController( + store=store, codex_bin="unused-test-codex", + ) + + +def drain(runtime, session_id, tmp_path): + runtime.resume_session_queue(session_id=session_id, work_dir=tmp_path, objective="Research") + # A worker that already retired has completed its queue cleanup. + with runtime.lock: + worker = runtime.session_queue_threads.get(session_id) + if worker: + worker.join(3) + assert not worker.is_alive() + + +@pytest.mark.parametrize("claimed", [False, True]) +def test_missing_runtime_asset_settles_accepted_work_without_dispatch( + tmp_path, monkeypatch, claimed, +): + store, sid, runtime = session_runtime(tmp_path) + turns = [store.create_queued_turn(sid, client_turn_id=f"request-{i}", message="Research cash flow")[0] + for i in range(2)] + if claimed: + assert store.claim_next_queued_turn(sid) + + def missing_asset(*args, **kwargs): + (tmp_path / "removed-release" / "SKILL.md").read_text() + + monkeypatch.setattr(runtime, "_ensure_adapter", missing_asset) + drain(runtime, sid, tmp_path) + for turn in turns: + result = runtime.wait_for_turn(session_id=sid, turn_id=turn["turn_id"], timeout_sec=0.1) + assert result["status"] == "failed" + assert result["error_code"] == "runtime_unavailable" + assert result.get("started_at") is None + assert result.get("upstream_turn_id") is None + assert str(tmp_path) not in result["error"] + events = store.events_after(sid, turn["turn_id"], None) + assert [e["kind"] for e in events].count("turn.failed") == 1 + assert not store.queued_turns(sid) + assert store.load_session(sid)["active_turn_id"] is None + + class Adapter: + upstream_thread_id = "public-thread" + messages = [] + + def healthcheck(self): + return True + + def start_turn(self, message, sink): + self.messages.append(message) + return {"message": "Recovered"} + + adapter = Adapter() + monkeypatch.setattr(runtime, "_ensure_adapter", lambda *a, **kw: adapter) + retry, created = runtime.enqueue_turn( + session_id=sid, client_turn_id="request-0", message="Research cash flow", + work_dir=tmp_path, objective="Research", + ) + assert not created + assert retry["status"] == "failed" + new, _ = runtime.enqueue_turn( + session_id=sid, client_turn_id="explicit-new-request", message="New request", + work_dir=tmp_path, objective="Research", + ) + assert runtime.wait_for_turn(session_id=sid, turn_id=new["turn_id"], timeout_sec=3)["status"] == "completed" + assert adapter.messages == ["New request"] + + +def test_stop_wins_over_late_preparation_failure(tmp_path, monkeypatch): + store, sid, runtime = session_runtime(tmp_path) + entered, release = threading.Event(), threading.Event() + + def preparation(*args, **kwargs): + entered.set() + assert release.wait(3) + raise OSError("runtime removed") + + monkeypatch.setattr(runtime, "_ensure_adapter", preparation) + turn, _ = runtime.enqueue_turn( + session_id=sid, client_turn_id="stop-race", message="Research", + work_dir=tmp_path, objective="Research", + ) + assert entered.wait(3) + assert runtime.interrupt_turn(session_id=sid, turn_id=turn["turn_id"])["status"] == "interrupted" + release.set() + drain(runtime, sid, tmp_path) + assert store.load_turn(sid, turn["turn_id"])["status"] == "interrupted" + assert not any(e["kind"] == "turn.failed" for e in store.events_after(sid, turn["turn_id"], None)) + + +def test_removed_manager_release_fails_before_starting_provider(tmp_path, monkeypatch): + import loopx.chat_manager as manager + + store, sid, runtime = session_runtime(tmp_path, channel="manager") + monkeypatch.setattr(manager, "__file__", str(tmp_path / "removed-release" / "chat_manager.py")) + starts = [] + monkeypatch.setattr(runtime, "_start_adapter", lambda **kw: starts.append(kw)) + turn, _ = runtime.enqueue_turn( + session_id=sid, client_turn_id="missing-manager-skill", message="Research cash flow", + work_dir=tmp_path, objective="Research", + ) + result = runtime.wait_for_turn(session_id=sid, turn_id=turn["turn_id"], timeout_sec=3) + assert result["status"] == "failed" + assert result["error_code"] == "runtime_unavailable" + assert result.get("started_at") is None + assert starts == [] + 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"), +]) +def test_lark_wait_receives_typed_failure_without_model_or_timeout(tmp_path, monkeypatch, error, code): + store, sid, runtime = session_runtime(tmp_path) + + def preparation(*args, **kwargs): + raise error + + monkeypatch.setattr(runtime, "_ensure_adapter", preparation) + real_wait = runtime.wait_for_turn + monkeypatch.setattr(runtime, "wait_for_turn", lambda **kw: real_wait(**kw, timeout_sec=3)) + with pytest.raises(LarkGoalTopicTurnFailed) as failure: + answer_lark_goal_topic( + route={"goal_id": "public-research", "session_id": sid, + "ingress_mode": "session_queue", "message_id": "om_public_question", + "topic_root_message_id": "om_public_root"}, + text="Research Microsoft's cash flow", work_dir=tmp_path, + objective="Research", runtime_controller=runtime, + ) + assert failure.value.error_code == code + reply_code, reply = manager_failure_reply(failure.value) + assert reply_code == code + assert "修复后重试" in reply + assert "private installation path" not in reply + assert not store.queued_turns(sid)