diff --git a/loopx/extensions/lark/docs/realtime-conversation-readiness.md b/loopx/extensions/lark/docs/realtime-conversation-readiness.md index 8f7c50b4ea..0b422a9a28 100644 --- a/loopx/extensions/lark/docs/realtime-conversation-readiness.md +++ b/loopx/extensions/lark/docs/realtime-conversation-readiness.md @@ -80,6 +80,15 @@ approval store, scheduler or model-execution authority. Missing evidence is not permission to restart a request in a fresh model thread, elevate host policy or acknowledge work as completed. +Private DM result reconciliation uses the existing Chat store's terminal-state +owner, including `timed_out`, for ordinary replies and commission results. +Timeouts return an explicit failure and the original conversation's recovery +controls. Provider readback still gates delivery: restarting the transport or +redelivering the source reconciles the saved attempt without another send or +model execution. Synthetic host/provider checks cover idle/hard DM timeouts +and an injected commission timeout; a live-provider recovery drill remains +part of the switch gate below. + ## Qualification before switching The acceptance owner is the steward RFC's diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index 50df77b0ea..829fabca37 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -11,7 +11,7 @@ from typing import Any from ...capabilities.native_chat.external_conversations import ChatExternalConversations -from ...chat_store import _atomic_write_json, _read_json +from ...chat_store import TERMINAL_TURN_STATES, _atomic_write_json, _read_json from ...file_lock import exclusive_file_lock from .conversation_identity import identity_ref, lark_private_source from .event_inbox import acknowledge_lark_event_inbox, ingest_lark_event_inbox @@ -261,10 +261,11 @@ def reconcile(self) -> int: # independent of terminal execution and reply delivery. self._deliver(path, record, "admission", "已持久受理到原 Agent 会话;等待原宿主领取。/status 查看持久队列,/project 返回普通项目对话。实时停止暂不支持,请在原宿主处理。" if native.get("agent_target") else "已持久受理;若已有执行,本条会排队。可发送 /status、/stop 或 /new。") turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) - if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: + if not turn or turn["status"] not in TERMINAL_TURN_STATES: continue response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else ( "本次执行已停止。" if turn["status"] == "interrupted" else + "本次执行超时,原会话已保留;请发送 /status 查看状态后再决定是否重试。" if turn["status"] == "timed_out" else "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") else: response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") @@ -274,12 +275,14 @@ def reconcile(self) -> int: resources = record.get("commission_resources") or {} if resources.get("session_id") and resources.get("turn_id"): first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) - if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: + if not first_turn or first_turn["status"] not in TERMINAL_TURN_STATES: record["status"] = "commission_running" _atomic_write_json(path, record) continue result_text = str((first_turn.get("response") or {}).get("message") or "") - result_text = "委托执行结果:\n" + result_text if first_turn["status"] == "completed" else "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。" + result_text = ("委托执行结果:\n" + result_text if first_turn["status"] == "completed" else + "委托首轮执行超时;原 Goal 和回执已保留,请查看状态后决定恢复。" if first_turn["status"] == "timed_out" else + "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。") proposal_id = native.get("proposal_id") if proposal_id: result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。" diff --git a/tests/test_lark_private_timeouts.py b/tests/test_lark_private_timeouts.py new file mode 100644 index 0000000000..1c08b8b175 --- /dev/null +++ b/tests/test_lark_private_timeouts.py @@ -0,0 +1,97 @@ +"""Real native timeout settlement with synthetic host and provider transports.""" +import json +import threading + +import pytest +from test_chat_ordinary_project import ordinary # noqa: F401 +from test_lark_private_conversations import connect +from test_native_steward_private import steward # noqa: F401 + +from loopx.chat_codex_goal import CodexGoalDriver +from loopx.extensions.lark.private_conversations import LarkPrivateConversations + + +@pytest.mark.parametrize("timeout_kind", ["idle", "hard"]) +def test_private_timeout_returns_once_after_delivery_readback_recovery(ordinary, timeout_kind): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + runtime.idle_timeout_sec = .1 if timeout_kind == "idle" else 30 + runtime.hard_timeout_sec = .1 if timeout_kind == "hard" else 30 + try: + event = provider.event("notes-app", "timeout", "wait for interrupt") + assert transport.admit("notes-app", event)["status"] == "durably_accepted" + row = transport.core.pending()[0] + sid, tid = row["session_id"], row["turn_id"] + terminal = runtime.wait_for_turn(session_id=sid, turn_id=tid, timeout_sec=10) + assert terminal["status"] == "timed_out" + assert terminal["error_code"] == f"{timeout_kind}_timeout" + original = store.load_session(sid)["upstream_thread_id"] + + provider.verify_replies = False + assert transport.reconcile() == 0 + assert any(profile == "notes-app" and "本次执行超时" in text for profile, text in provider.writes) + assert not transport.core.read_request(row["request_ref"]).get("delivery_verified") + count = len(provider.writes) + restarted = LarkPrivateConversations(controller=runtime, runtime_root=store.root.parent, + runner=provider, cli_bin="lark-cli") + assert restarted.reconcile() == 0 + assert len(provider.writes) == count + provider.verify_replies = True + assert restarted.reconcile() == 1 + assert restarted.core.read_request(row["request_ref"])["delivery_verified"] is True + assert restarted.health()[row["binding_id"]] == {"pending_count": 0, "recovery_count": 0} + assert restarted.admit("notes-app", {**event, "event_id": "redelivery"})["status"] == "durably_accepted" + restarted.reconcile() + assert len(provider.writes) == count + assert len(store.list_sessions()) == 1 + assert store.load_session(sid)["upstream_thread_id"] == original + assert store.load_session(sid)["goal_id"] is None + assert store.load_turn(sid, tid)["status"] == "timed_out" + requests = [json.loads(line) for line in ordinary[4].read_text().splitlines()] + assert len([r for r in requests if r.get("method") == "turn/start"]) == 1 + finally: + runtime.close() + + +def test_commission_timeout_returns_original_resources_without_restarting_work(steward, monkeypatch): # noqa: F811 + store, runtime, provider, transport, binding, capture, _ = steward + started, release = threading.Event(), threading.Event() + # Inject a typed host timeout after Goal activation. Core persists the + # terminal Turn; this does not qualify the real native Goal deadline. + def timeout(self, emit): + started.set() + assert release.wait(10), "timeout fixture was not released" + raise self.session._timeout_error("hard_timeout", "Synthetic host timeout") + monkeypatch.setattr(CodexGoalDriver, "_observe", timeout) + event = provider.event("steward-app", "delegate", "/delegate --tokens 12000 wait for interrupt") + assert transport.admit("steward-app", event)["status"] == "command_recorded" + transport.reconcile() + proposal = transport.core.actions.store.list()[0] + confirm = provider.event("steward-app", "confirm", "/confirm " + proposal["proposal_id"]) + assert transport.admit("steward-app", confirm)["status"] == "command_recorded" + transport.reconcile() + resources = transport.core.actions.load(proposal["proposal_id"])["receipt"]["resource_ids"] + assert started.wait(10) + provider.verify_replies = False + release.set() + turn = runtime.wait_for_turn(session_id=resources["session_id"], turn_id=resources["turn_id"], timeout_sec=10) + assert turn["status"] == "timed_out" and turn["error_code"] == "hard_timeout" + transport.reconcile() + assert any(profile == "steward-app" and "委托首轮执行超时" in text and + f"/resume-commission {proposal['proposal_id']}" in text for profile, text in provider.writes) + count = len(provider.writes) + before = capture.read_text() + restarted = LarkPrivateConversations(controller=runtime, runtime_root=store.root.parent, + runner=provider, cli_bin="lark-cli") + assert restarted.reconcile() == 0 + assert len(provider.writes) == count + provider.verify_replies = True + assert restarted.reconcile() == 1 + assert restarted.health()[binding["binding_id"]] == {"pending_count": 0, "recovery_count": 0} + restarted.admit("steward-app", {**confirm, "event_id": "redelivery"}) + restarted.reconcile() + assert len(provider.writes) == count + assert transport.core.actions.load(proposal["proposal_id"])["receipt"]["resource_ids"] == resources + assert len(json.loads(runtime.registry_path.read_text())["goals"]) == 1 + assert store.load_turn(resources["session_id"], resources["turn_id"])["status"] == "timed_out" + assert capture.read_text() == before + assert not any(profile == "notes-app" for profile, _ in provider.writes)