From 2b98053379b61beefc5307d7e66aa8f3b6333820 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:57:10 +0800 Subject: [PATCH 1/2] fix(collaboration): recover return verification without prose classifiers Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../app-conversation-and-async-inbox-v0.md | 8 +++ loopx/capabilities/manager_context/README.md | 7 +++ .../capabilities/manager_context/roundtrip.py | 54 +++++++------------ .../collaboration/return_delivery.ts | 8 +++ .../manager_return_delivery.test.ts | 23 ++++++++ tests/test_manager_context_roundtrip.py | 45 ++++++++++------ 6 files changed, 94 insertions(+), 51 deletions(-) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index a08fbe6bcd..2953509ba8 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -220,6 +220,14 @@ Todo/lease, model admission and artifact acceptance keep their existing owners. Dispatch events wake the existing driver within admission; polling repairs gaps. An inbox is not permission to start another automation. +The return-verification slice uses the existing TS classification owner for +both adapter results and typed resolution failures. Exception text is diagnostic, +not route/authority evidence: a transient read failure retains its locator and +backoff, then reconciles the original result without another send. Explicit +revocation, lost routing and missing initial receipts remain terminal. This +qualifies the persisted recovery boundary, not live provider availability or +the complete GQ09 journey. + Before each extraction report base/head real-call latency, boundary crossings, bytes, owners deleted/retained and compatibility callers. Product delivery must not wait for full Python retirement. Python may retain IO; TS owns migrated diff --git a/loopx/capabilities/manager_context/README.md b/loopx/capabilities/manager_context/README.md index 352eeff08c..cad60d244b 100644 --- a/loopx/capabilities/manager_context/README.md +++ b/loopx/capabilities/manager_context/README.md @@ -312,6 +312,13 @@ readback; `manager-context` remains the sole result/delivery writer. The typed attempt validation and verification classification; Python retains file-lock, persistence and adapter orchestration only. +Unclassified readback exceptions retain the saved attempt and retry with backoff; +their wording never establishes revoked authority or a missing route. Adapters +must raise `ReturnResolutionBlocked` with an exact typed resolution reason for +those permanent failures. This replaces the old exception-substring fallback. +A successful later readback updates the same App transcript receipt without +another model turn or external send. Actual live grant checks still precede it. + New handoffs persist their exact original return route. Legacy requests remain queryable; a receiver can explicitly report one only when its exact persisted Chat receipt uniquely recovers the route. Historical timestamps stay unknown. diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index 21ca14c8bd..b4c5676968 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -44,18 +44,6 @@ "delivery_state_unreadable", } -# Bounded reasons why one return cannot be resolved at all. Adapters raise -# ``ReturnResolutionBlocked`` with one of these instead of relying on their prose -# being re-parsed, so delivery state never depends on message substrings. -RETURN_RESOLUTION_REASONS = frozenset( - { - "return_authorization_unavailable", - "original_route_unavailable", - "initial_delivery_receipt_unavailable", - } -) - - class ReturnResolutionBlocked(ValueError, RuntimeError): """A return cannot be resolved; ``reason`` is the provider-neutral code. @@ -64,7 +52,10 @@ class ReturnResolutionBlocked(ValueError, RuntimeError): """ def __init__(self, reason, message): - if reason not in RETURN_RESOLUTION_REASONS: + decision = _verification_decision({ + "verification_performed": False, "reply_verified": False, "blocker": reason, + }) + if decision["status"] != "explicit_unverified": raise ValueError("unsupported return resolution reason") self.reason = reason super().__init__(message) @@ -105,23 +96,12 @@ def _verification_decision(outcome): def _verification_exception_error(exc): - reason = getattr(exc, "reason", None) - if reason in RETURN_RESOLUTION_REASONS: - return reason - # Compatibility fallback for adapter text that still arrives as prose. New - # adapter failures must raise ReturnResolutionBlocked with a typed reason. - message = str(exc) - if ( - "authorization" in message - or "authorized" in message - or "authority" in message - ): - return "return_authorization_unavailable" - if any(token in message for token in ("conversation", "route", "binding", "target")): - return "original_route_unavailable" - if "initial reply" in message or "initial_receipt" in message: - return "initial_delivery_receipt_unavailable" - return None + decision = _verification_decision({ + "verification_performed": False, + "reply_verified": False, + "blocker": exc.reason if isinstance(exc, ReturnResolutionBlocked) else None, + }) + return decision["error"] if decision["status"] == "explicit_unverified" else None def register(root, row, session, turn): @@ -170,16 +150,18 @@ def _route(root, row): if receipt.get("request_id") == row["request_id"]: matches.append((session, turn)) if len(matches) != 1: - raise ValueError("original Chat return route unavailable or ambiguous") + raise ReturnResolutionBlocked( + "original_route_unavailable", "original Chat return route unavailable or ambiguous" + ) register(root, row, *matches[0]) value = _read(path) if value.get("kind") == "peer" and (row.get("source_kind") != "peer" or value.get("source_agent_id") != row.get("source_agent_id")): - raise ValueError("peer return route identity mismatch") + raise ReturnResolutionBlocked("original_route_unavailable", "peer return route identity mismatch") if any( value.get(k) != row.get(k) for k in ("request_id", "goal_id", "agent_id", "source_id") ): - raise ValueError("context return route identity mismatch") + raise ReturnResolutionBlocked("original_route_unavailable", "context return route identity mismatch") return value @@ -338,16 +320,16 @@ def drain(root, registry, store, external_sender, *, now=None, cancelled=lambda: or session.get("channel_id") != route["channel_id"] or not turn ): - raise ValueError("original_conversation_unavailable") + raise ReturnResolutionBlocked("original_route_unavailable", "original_conversation_unavailable") grant = authority(root, registry, session, turn) target = {k: row[k] for k in ("goal_id", "agent_id")} if ( target not in grant["targets"] or grant.get("source_id") != row["source_id"] ): - raise ValueError("return_authorization_unavailable") + raise ReturnResolutionBlocked("return_authorization_unavailable", "return_authorization_unavailable") if turn.get("status") != "completed": - raise ValueError("initial_receipt_not_completed") + raise ReturnResolutionBlocked("initial_delivery_receipt_unavailable", "initial_receipt_not_completed") if ( path.stem == "decision" and (path.parent / "conclusion.json").exists() diff --git a/loopx/control_plane/collaboration/return_delivery.ts b/loopx/control_plane/collaboration/return_delivery.ts index 1efb99a6a3..7165579f60 100644 --- a/loopx/control_plane/collaboration/return_delivery.ts +++ b/loopx/control_plane/collaboration/return_delivery.ts @@ -96,6 +96,14 @@ export function classifyManagerReturnVerification(value: unknown): JsonObject { }; } if (!performed) { + // A transport outage is not proof that the original route or grant was + // revoked. Only exact, adapter-declared resolution reasons stop readback. + const blocker = outcome.blocker; + if (blocker === "return_authorization_unavailable" + || blocker === "original_route_unavailable" + || blocker === "initial_delivery_receipt_unavailable") { + return { status: "explicit_unverified", error: blocker, verification: null }; + } return { status: "verification_required", error: "provider_verification_unavailable", diff --git a/tests/control_plane_ts/manager_return_delivery.test.ts b/tests/control_plane_ts/manager_return_delivery.test.ts index 4b2cefd2b4..d4c52e894a 100644 --- a/tests/control_plane_ts/manager_return_delivery.test.ts +++ b/tests/control_plane_ts/manager_return_delivery.test.ts @@ -74,3 +74,26 @@ test("classifies verification without exposing provider prose", () => { }, ); }); + +test("only typed resolution blockers stop an unavailable verification", () => { + for (const blocker of [ + "return_authorization_unavailable", + "original_route_unavailable", + "initial_delivery_receipt_unavailable", + ]) { + assert.deepEqual(classifyManagerReturnVerification({ + verification_performed: false, reply_verified: false, blocker, + }), { status: "explicit_unverified", error: blocker, verification: null }); + } + for (const blocker of [null, "route lookup temporarily unavailable", + "authorization service read timed out", "initial reply read interrupted", + "original_route_unavailable: timeout", { reason: "original_route_unavailable" }]) { + assert.deepEqual(classifyManagerReturnVerification({ + verification_performed: false, reply_verified: false, blocker, + }), { + status: "verification_required", + error: "provider_verification_unavailable", + verification: null, + }); + } +}); diff --git a/tests/test_manager_context_roundtrip.py b/tests/test_manager_context_roundtrip.py index 6abda05053..04a9e8d9a7 100644 --- a/tests/test_manager_context_roundtrip.py +++ b/tests/test_manager_context_roundtrip.py @@ -643,18 +643,15 @@ def accepted_attempt(): @pytest.mark.parametrize( - "failure", + "reason", [ - lambda: ReturnResolutionBlocked( - "return_authorization_unavailable", "context return authority revoked" - ), - lambda: ValueError("context return authority revoked"), - lambda: ValueError("manager connection no longer authorized"), + "return_authorization_unavailable", + "original_route_unavailable", + "initial_delivery_receipt_unavailable", ], - ids=["typed-reason", "prose-revoked", "prose-unauthorized"], ) -def test_return_authority_revocation_terminalizes_without_repeating_readback( - flow, failure +def test_typed_return_resolution_terminalizes_without_repeating_readback( + flow, reason ): root, registry, store, create = flow _, _, receipt = create(True) @@ -671,14 +668,14 @@ def send_with_attempt(self, route, session, turn, text, record_attempt): def verify(self, *_args): self.verify_calls += 1 - raise failure() + raise ReturnResolutionBlocked(reason, "Adapter resolution blocked") transport = Transport() drain(root, registry, store, transport) drain(root, registry, ChatSessionStore(root), transport) state = reply_status(root, receipt)[0] assert state["status"] == "explicit_unverified" - assert state["error"] == "return_authorization_unavailable" + assert state["error"] == reason assert transport.verify_calls == 1 # A terminal reason is never re-read, not even after the backoff window. drain( @@ -691,23 +688,34 @@ def verify(self, *_args): assert transport.verify_calls == 1 -def test_unclassified_verification_failure_backs_off_instead_of_hot_looping(flow): +@pytest.mark.parametrize("message", [ + "provider readback transport failed", + "route lookup temporarily unavailable", + "authorization service read timed out", + "initial reply read interrupted", +]) +def test_unclassified_verification_failure_recovers_without_resending(flow, message): root, registry, store, create = flow - _, _, receipt = create(True) + session, _, receipt = create(True) rid = receipt["request_id"] acknowledge(root, "research", "worker", rid, "adopt", "Checked") report(root, "research", "worker", rid, "conclusion", "Bounded result.") class Transport: verify_calls = 0 + send_calls = 0 def send_with_attempt(self, route, session, turn, text, record_attempt): + self.send_calls += 1 record_attempt(accepted_attempt()) - return {"external_write_performed": True, "reply_verified": False} + # A recorded provider write resumes readback regardless of wording. + raise RuntimeError(message) def verify(self, *_args): self.verify_calls += 1 - raise RuntimeError("provider readback transport failed") + if self.verify_calls == 1: + raise RuntimeError(message) + return {"verification_performed": True, "reply_verified": True} transport = Transport() drain(root, registry, store, transport) @@ -726,6 +734,13 @@ def verify(self, *_args): now=datetime.now(timezone.utc) + timedelta(days=2), ) assert transport.verify_calls == 2 + assert transport.send_calls == 1 + assert reply_status(root, receipt)[0]["status"] == "delivered" + snapshot = project_chat_session_snapshot(root, ChatSessionStore(root), session["session_id"]) + replies = [m for m in snapshot["messages"] if m.get("origin") == "manager_followup"] + assert len(replies) == 1 + assert replies[0]["return_delivery"]["status"] == "delivered" + assert message not in json.dumps(snapshot) def test_public_delivery_projection_normalizes_unknown_private_state(flow): From 93bf4cb04e551c1e0fbcafd9a6043427c485db26 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 20:58:35 +0800 Subject: [PATCH 2/2] test(collaboration): qualify verification recovery after Goal recreation Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../app-conversation-and-async-inbox-v0.md | 3 ++ .../project_registry_io_manifest_v1.json | 4 +- tests/test_collaboration_goal_instance.py | 38 ++++++++++++++++++- 3 files changed, 42 insertions(+), 3 deletions(-) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 52d263f787..47b8f12652 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -225,6 +225,9 @@ both adapter results and typed resolution failures. Exception text is diagnostic not route/authority evidence: a transient read failure retains its locator and backoff, then reconciles the original result without another send. Explicit revocation, lost routing and missing initial receipts remain terminal. This +also holds after a `source_session_v1` Goal is recreated: a crash-persisted +attempt recovers on the original GoalRef and conversation, transient verification +backs off without resending, and typed terminal blockers stay stopped. This qualifies the persisted recovery boundary, not live provider availability or the complete GQ09 journey. diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 21dbd26142..e45ac09333 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -215,7 +215,7 @@ }, { "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", - "line": 883, + "line": 865, "column": 22, "kind": "codec_read", "api": "load_project_registry", @@ -447,7 +447,7 @@ }, { "site": "loopx/cli.py::.main::codec_read:load_project_registry#1", - "line": 837, + "line": 846, "column": 17, "kind": "codec_read", "api": "load_project_registry", diff --git a/tests/test_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py index 1b50b58512..103fc72b3b 100644 --- a/tests/test_collaboration_goal_instance.py +++ b/tests/test_collaboration_goal_instance.py @@ -21,6 +21,7 @@ ) from loopx.capabilities.manager_context.roundtrip import ( EXACT_DELIVERY_ADMISSION_SECONDS, + ReturnResolutionBlocked, drain, report, ) @@ -513,8 +514,17 @@ def __call__(self, *_args): } +@pytest.mark.parametrize("verification_failure, terminal", [ + (None, False), + ("route lookup temporarily unavailable", False), + ("authorization service read timed out", False), + ("initial reply read interrupted", False), + ("original_route_unavailable", True), + ("return_authorization_unavailable", True), + ("initial_delivery_receipt_unavailable", True), +]) def test_exact_external_return_verifies_after_recreation_without_resend( - tmp_path: Path, + tmp_path: Path, verification_failure: str | None, terminal: bool, ) -> None: registry = _create_source_registry(tmp_path) store, _, receipt = _external_manager_request(tmp_path, registry) @@ -546,6 +556,10 @@ def send_with_attempt( def verify(self, *_args): self.verify_calls += 1 + if self.verify_calls == 1 and verification_failure is not None: + if terminal: + raise ReturnResolutionBlocked(verification_failure, "Resolution blocked") + raise RuntimeError(verification_failure) return { "ok": True, "verification_performed": True, @@ -590,6 +604,28 @@ def verify(self, *_args): ) assert transport.send_calls == 1 assert transport.verify_calls == 1 + state_path = ( + _root(tmp_path) / "replies" / receipt["request_id"] / "conclusion.delivery.json" + ) + if verification_failure is not None: + failed = json.loads(state_path.read_text(encoding="utf-8")) + assert failed["goal_ref"] == receipt["goal_ref"] + assert failed["status"] == ("explicit_unverified" if terminal else "retry_pending") + if terminal: + assert failed["error"] == verification_failure + else: + assert failed["attempt"] == first["attempt"] + # Another immediate pump must honor backoff, even after recreation. + drain(tmp_path, registry, ChatSessionStore(tmp_path), transport, + now=admitted_at + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS + 1)) + assert transport.verify_calls == 1 + drain(tmp_path, registry, ChatSessionStore(tmp_path), transport, + now=admitted_at + timedelta(days=1)) + assert transport.send_calls == 1 + assert transport.verify_calls == (1 if terminal else 2) + if terminal: + assert json.loads(state_path.read_text(encoding="utf-8")) == failed + return state = json.loads( ( _root(tmp_path)