diff --git a/pyproject.toml b/pyproject.toml
index 37be35eef9..72a346a7dc 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -51,6 +51,7 @@ include = ["loopx*"]
"loopx" = ["web/chat/*.html", "web/chat/assets/*", "web/chat/manifest.webmanifest", "web/chat/pwa/*"]
"loopx.control_plane" = ["*.json", "*.ts", "presentation/*.ts"]
"loopx.control_plane.agents" = ["*.ts"]
+"loopx.control_plane.collaboration" = ["*.ts"]
"loopx.control_plane.coordination" = ["*.ts", "*.json"]
"loopx.control_plane.goals" = ["*.ts"]
"loopx.control_plane.quota" = ["*.ts"]
diff --git a/tests/control_plane_ts/manager_return_delivery.test.ts b/tests/control_plane_ts/manager_return_delivery.test.ts
new file mode 100644
index 0000000000..1f6e7d6235
--- /dev/null
+++ b/tests/control_plane_ts/manager_return_delivery.test.ts
@@ -0,0 +1,65 @@
+import assert from "node:assert/strict";
+import test from "node:test";
+
+import {
+ classifyManagerReturnVerification,
+ normalizeManagerReturnDeliveryAttempt,
+} from "../../loopx/control_plane/collaboration/return_delivery.ts";
+
+const attempt = {
+ schema_version: "manager_return_delivery_attempt_v0",
+ provider: "lark",
+ message_ref: "om_provider_reply",
+ intent_digest: `sha256:${"a".repeat(64)}`,
+ provider_receipt: `sha256:${"b".repeat(64)}`,
+};
+
+test("normalizes the exact provider-neutral delivery attempt", () => {
+ assert.deepEqual(normalizeManagerReturnDeliveryAttempt(attempt), attempt);
+ assert.throws(
+ () => normalizeManagerReturnDeliveryAttempt({ ...attempt, private_payload: "no" }),
+ /unsupported or missing fields/,
+ );
+ assert.throws(
+ () => normalizeManagerReturnDeliveryAttempt({ ...attempt, message_ref: "bad ref" }),
+ /message_ref is invalid/,
+ );
+});
+
+test("classifies verification without exposing provider prose", () => {
+ assert.deepEqual(
+ classifyManagerReturnVerification({
+ verification_performed: true,
+ reply_verified: true,
+ }),
+ {
+ status: "delivered",
+ error: null,
+ verification: "reconciled_after_restart",
+ },
+ );
+ assert.deepEqual(
+ classifyManagerReturnVerification({
+ verification_performed: false,
+ reply_verified: false,
+ blocker: "private provider outage detail",
+ }),
+ {
+ status: "verification_required",
+ error: "provider_verification_unavailable",
+ verification: null,
+ },
+ );
+ assert.deepEqual(
+ classifyManagerReturnVerification({
+ verification_performed: true,
+ reply_verified: false,
+ blocker: "private provider mismatch detail",
+ }),
+ {
+ status: "explicit_unverified",
+ error: "provider_delivery_mismatch",
+ verification: null,
+ },
+ );
+});
diff --git a/tests/extensions/test_lark_inbox_reactions.py b/tests/extensions/test_lark_inbox_reactions.py
index 23c0e05035..1d448fee0f 100644
--- a/tests/extensions/test_lark_inbox_reactions.py
+++ b/tests/extensions/test_lark_inbox_reactions.py
@@ -22,6 +22,7 @@
from loopx.extensions.lark.inbox_reply import (
reply_lark_event_inbox,
send_lark_inbox_message,
+ verify_lark_inbox_reply,
)
@@ -927,6 +928,105 @@ def test_multiline_readback_must_preserve_line_structure(tmp_path: Path) -> None
assert result["reply_verified"] is False
+def test_unverified_reply_records_private_locator_and_read_only_recovery(
+ tmp_path: Path,
+) -> None:
+ config, _, project = _fixture(tmp_path, lifecycle=False)
+ attempts: list[dict[str, str]] = []
+ first_runner = ReplyRunner(matching_readback=False)
+
+ sent = reply_lark_event_inbox(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text="处理完成",
+ execute=True,
+ runner=first_runner,
+ delivery_attempt_recorder=attempts.append,
+ )
+
+ assert sent["status"] == "sent_unverified"
+ assert attempts == [
+ {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_reply_fixture",
+ "intent_digest": attempts[0]["intent_digest"],
+ "provider_receipt": sent["idempotency_key"],
+ }
+ ]
+ assert attempts[0]["intent_digest"].startswith("sha256:")
+ assert "message_ref" not in sent
+
+ recovery_runner = ReplyRunner()
+ recovered = verify_lark_inbox_reply(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text="处理完成",
+ attempt=attempts[0],
+ runner=recovery_runner,
+ )
+
+ assert recovered["reply_verified"] is True
+ assert recovered["verification_performed"] is True
+ assert not any(
+ "+messages-send" in call or "+messages-reply" in call
+ for call in recovery_runner.calls
+ )
+
+
+def test_read_only_recovery_rejects_changed_intent_without_provider_call(
+ tmp_path: Path,
+) -> None:
+ config, _, project = _fixture(tmp_path, lifecycle=False)
+ runner = ReplyRunner()
+ result = verify_lark_inbox_reply(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text="different result",
+ attempt={
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_reply_fixture",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ },
+ runner=runner,
+ )
+
+ assert result["reply_verified"] is False
+ assert result["verification_performed"] is True
+ assert result["blocker"] == "provider_delivery_intent_conflict"
+ assert runner.calls == []
+
+
+def test_reply_fails_closed_after_send_when_locator_persistence_fails(
+ tmp_path: Path,
+) -> None:
+ config, _, project = _fixture(tmp_path, lifecycle=False)
+ runner = ReplyRunner()
+
+ result = reply_lark_event_inbox(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text="处理完成",
+ execute=True,
+ runner=runner,
+ delivery_attempt_recorder=lambda _attempt: (_ for _ in ()).throw(
+ OSError("synthetic private persistence failure")
+ ),
+ )
+
+ assert result["external_write_performed"] is True
+ assert result["reply_verified"] is False
+ assert result["blocker"] == "lark_inbox_reply_delivery_attempt_not_persisted"
+ assert sum("+messages-reply" in call and "--dry-run" not in call for call in runner.calls) == 1
+ assert not any("+messages-mget" in call for call in runner.calls)
+
+
def test_verified_reply_accepts_provider_token_or_rendered_mention_name(
tmp_path: Path,
) -> None:
diff --git a/tests/extensions/test_lark_manager_returns.py b/tests/extensions/test_lark_manager_returns.py
index 398ebb11d8..316707588f 100644
--- a/tests/extensions/test_lark_manager_returns.py
+++ b/tests/extensions/test_lark_manager_returns.py
@@ -155,8 +155,9 @@ def invoke():
runner=transport,
)
- with pytest.raises(ValueError, match="initial reply"):
+ with pytest.raises(ValueError, match="initial reply") as pending:
invoke()
+ assert pending.value.reason == "initial_delivery_receipt_unavailable"
assert not calls
acknowledge_lark_event_inbox(
project=root,
@@ -170,8 +171,9 @@ def invoke():
reaction_id="reaction_Get", emoji_type="Get",
)
if revoke_before_send:
- with pytest.raises(ValueError, match="revoked"):
+ with pytest.raises(ValueError, match="revoked") as revoked:
invoke()
+ assert revoked.value.reason == "return_authorization_unavailable"
assert not sent
else:
result = invoke()
diff --git a/tests/extensions/test_lark_markdown_reply.py b/tests/extensions/test_lark_markdown_reply.py
index 791b199467..69f378c6e7 100644
--- a/tests/extensions/test_lark_markdown_reply.py
+++ b/tests/extensions/test_lark_markdown_reply.py
@@ -7,7 +7,10 @@
lark_markdown_readback_matches, normalize_lark_outbound_text,
safe_lark_plain_text_fallback,
)
-from loopx.extensions.lark.inbox_reply import reply_lark_event_inbox
+from loopx.extensions.lark.inbox_reply import (
+ reply_lark_event_inbox,
+ verify_lark_inbox_reply,
+)
from test_lark_inbox_reactions import ReplyRunner, _fixture
TEXT = "进展\n\n- **结果**\n - 证据\n\n```python\nif ok:\n done()\n```"
@@ -108,6 +111,60 @@ def runner(args):
assert sent[0][sent[0].index("--text") + 1] == text.strip()
+def test_large_post_fallback_can_be_verified_without_resending(tmp_path):
+ config, _, project = _fixture(tmp_path, lifecycle=False)
+ text = "- 完整内容\n" * 2000
+ fallback = ReplyRunner(readback_text="different text")
+ attempts = []
+
+ def runner(args):
+ if "--content" in args:
+ return {
+ "returncode": 0,
+ "stdout": json.dumps(
+ {
+ "api": [
+ {
+ "body": {
+ "msg_type": "post",
+ "content": args[args.index("--content") + 1],
+ }
+ }
+ ]
+ }
+ ),
+ }
+ return fallback(args)
+
+ first = reply_lark_event_inbox(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text=text,
+ content_format="markdown",
+ execute=True,
+ runner=runner,
+ delivery_attempt_recorder=attempts.append,
+ )
+ assert first["status"] == "sent_unverified"
+ assert first["content_format"] == "text"
+
+ recovery = ReplyRunner(readback_text=text.strip())
+ verified = verify_lark_inbox_reply(
+ project=project,
+ config_path=config,
+ message_id="om_reaction_fixture",
+ text=text,
+ attempt=attempts[0],
+ runner=recovery,
+ )
+ assert verified["reply_verified"] is True
+ assert not any(
+ "+messages-send" in call or "+messages-reply" in call
+ for call in recovery.calls
+ )
+
+
def test_markdown_request_with_mentions_keeps_verified_text_transport(tmp_path):
config, _, project = _fixture(tmp_path, lifecycle=False)
text = '
Reviewer please review'
diff --git a/tests/test_manager_context_roundtrip.py b/tests/test_manager_context_roundtrip.py
index bb749a4629..c01054df62 100644
--- a/tests/test_manager_context_roundtrip.py
+++ b/tests/test_manager_context_roundtrip.py
@@ -18,7 +18,11 @@
)
from loopx.capabilities.manager_context.roundtrip import (
ReturnService,
+ ReturnResolutionBlocked,
+ _hash,
drain,
+ project_chat_return_deliveries,
+ project_chat_session_snapshot,
report,
reply_status,
)
@@ -321,7 +325,422 @@ def ambiguous(*args):
now=datetime.now(timezone.utc) + timedelta(days=2),
)
assert len(calls) == 1
- assert reply_status(root, r)[0]["status"] == "verification_required"
+ assert reply_status(root, r)[0] == {
+ "phase": "conclusion",
+ "status": "explicit_unverified",
+ "created_at": reply_status(root, r)[0]["created_at"],
+ "delivered_at": None,
+ "error": "provider_locator_unavailable",
+ }
+
+
+def test_known_provider_locator_is_verified_after_restart_without_resend(flow):
+ root, registry, store, create = flow
+ session, _, receipt = create(True)
+ rid = receipt["request_id"]
+ acknowledge(root, "research", "worker", rid, "adopt", "Checked")
+ report(
+ root,
+ "research",
+ "worker",
+ rid,
+ "conclusion",
+ "Processed with a recorded validation result.",
+ )
+
+ class Transport:
+ def __init__(self):
+ self.send_calls = 0
+ self.verify_calls = 0
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ self.send_calls += 1
+ record_attempt(
+ {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_provider_reply",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ }
+ )
+ raise RuntimeError("synthetic process interruption after provider send")
+
+ def verify(self, route, session, turn, text, attempt):
+ self.verify_calls += 1
+ assert attempt["message_ref"] == "om_provider_reply"
+ return {
+ "ok": True,
+ "verification_performed": True,
+ "reply_verified": True,
+ }
+
+ transport = Transport()
+ drain(root, registry, store, transport)
+ first = reply_status(root, receipt)[0]
+ assert first["status"] == "verification_required"
+ assert "message_ref" not in first and "intent_digest" not in first
+ first_projection = project_chat_return_deliveries(
+ root, session["session_id"], store.messages(session["session_id"])
+ )
+ assert next(
+ row for row in first_projection if row.get("origin") == "manager_followup"
+ )["return_delivery"]["status"] == "verification_required"
+ assert "om_provider_reply" not in json.dumps(first_projection)
+
+ drain(root, registry, ChatSessionStore(root), transport)
+ recovered = reply_status(root, receipt)[0]
+ assert recovered["status"] == "delivered"
+ assert recovered["verification"] == "reconciled_after_restart"
+ assert transport.send_calls == 1
+ assert transport.verify_calls == 1
+
+ projected = project_chat_return_deliveries(
+ root, session["session_id"], store.messages(session["session_id"])
+ )
+ returned = [row for row in projected if row.get("origin") == "manager_followup"]
+ assert returned[0]["return_delivery"] == {
+ "schema_version": "manager_return_delivery_status_v0",
+ **recovered,
+ }
+
+
+def test_chat_snapshot_keeps_current_session_delivery_after_unrelated_route_limit(flow):
+ root, _, store, _ = flow
+ session = store.create_session(
+ goal_id="loopx-manager",
+ agent_id="codex",
+ adapter_kind="codex_app_server",
+ upstream_thread_id="projection-limit",
+ channel_id="manager",
+ )
+ turn, _ = store.create_turn(
+ session["session_id"],
+ client_turn_id="projection-limit",
+ message="Check delivery projection completeness",
+ origin="web",
+ )
+
+ expected = {}
+ for request_id, status in (
+ ("e" * 64, "delivered"),
+ ("f" * 64, "verification_required"),
+ ):
+ _write(
+ _root(root) / "roundtrips" / f"{request_id}.json",
+ {"request_id": request_id, "session_id": session["session_id"]},
+ )
+ _write(
+ _root(root) / "replies" / request_id / "conclusion.json",
+ {
+ "request_id": request_id,
+ "phase": "conclusion",
+ "created_at": "2026-09-14T00:00:00+00:00",
+ },
+ )
+ delivery = {"status": status}
+ if status == "verification_required":
+ delivery["error"] = "provider_delivery_unverified"
+ _write(
+ _root(root) / "replies" / request_id / "conclusion.delivery.json",
+ delivery,
+ )
+ message_id = "handoff." + _hash([request_id, "conclusion"])
+ store.append_message(
+ session["session_id"],
+ role="agent",
+ text=f"Projected {status}",
+ turn_id=turn["turn_id"],
+ origin="manager_followup",
+ message_id=message_id,
+ )
+ expected[message_id] = status
+
+ for index in range(2000):
+ request_id = f"{index:064x}"
+ _write(
+ _root(root) / "roundtrips" / f"{request_id}.json",
+ {"request_id": request_id, "session_id": f"unrelated-{index}"},
+ )
+
+ snapshot = project_chat_session_snapshot(root, store, session["session_id"])
+ returned = {
+ row["message_id"]: row["return_delivery"]["status"]
+ for row in snapshot["messages"]
+ if row.get("message_id") in expected
+ }
+ assert returned == expected
+
+
+def test_provider_verification_outage_retries_read_only_without_resend(flow):
+ root, registry, store, create = flow
+ _, _, receipt = create(True)
+ rid = receipt["request_id"]
+ acknowledge(root, "research", "worker", rid, "adopt", "Checked")
+ report(root, "research", "worker", rid, "conclusion", "Bounded result.")
+
+ class Transport:
+ send_calls = 0
+ verify_calls = 0
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ self.send_calls += 1
+ record_attempt(
+ {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_provider_reply",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ }
+ )
+ return {"external_write_performed": True, "reply_verified": False}
+
+ def verify(self, *_args):
+ self.verify_calls += 1
+ return {
+ "ok": False,
+ "verification_performed": False,
+ "reply_verified": False,
+ "blocker": "provider_verification_unavailable",
+ }
+
+ transport = Transport()
+ now = datetime.now(timezone.utc)
+ drain(root, registry, store, transport, now=now)
+ drain(root, registry, ChatSessionStore(root), transport, now=now)
+ state = reply_status(root, receipt)[0]
+ assert state["status"] == "verification_required"
+ assert state["error"] == "provider_verification_unavailable"
+ assert transport.send_calls == 1 and transport.verify_calls == 1
+
+
+def test_provider_verification_mismatch_is_terminal_and_public_safe(flow):
+ root, registry, store, create = flow
+ _, _, receipt = create(True)
+ rid = receipt["request_id"]
+ acknowledge(root, "research", "worker", rid, "adopt", "Checked")
+ report(root, "research", "worker", rid, "conclusion", "Bounded result.")
+
+ class Transport:
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ record_attempt(
+ {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_private_provider_reply",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ }
+ )
+ return {"external_write_performed": True, "reply_verified": False}
+
+ def verify(self, *_args):
+ return {
+ "ok": False,
+ "verification_performed": True,
+ "reply_verified": False,
+ "blocker": "private provider mismatch detail",
+ }
+
+ 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"] == "provider_delivery_mismatch"
+ assert "private" not in str(state)
+
+
+def test_provider_verification_stops_after_return_authority_revocation(flow):
+ root, registry, store, create = flow
+ _, _, 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
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ record_attempt(
+ {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_provider_reply",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ }
+ )
+ return {"external_write_performed": True, "reply_verified": False}
+
+ def verify(self, *_args):
+ self.verify_calls += 1
+ return {"verification_performed": True, "reply_verified": True}
+
+ transport = Transport()
+ drain(root, registry, store, transport)
+ _write(
+ _root(root) / "policy.json", {"schema_version": POLICY_SCHEMA, "sources": {}}
+ )
+ drain(root, registry, ChatSessionStore(root), transport)
+
+ state = reply_status(root, receipt)[0]
+ assert state["status"] == "explicit_unverified"
+ assert state["error"] == "return_authorization_unavailable"
+ assert transport.verify_calls == 0
+
+
+def test_legacy_verification_required_without_locator_is_never_resent(flow):
+ """A record from before locators were persisted must not be guessed or resent."""
+ root, registry, store, create = flow
+ _, _, receipt = create(True)
+ rid = receipt["request_id"]
+ acknowledge(root, "research", "worker", rid, "adopt", "Checked")
+ report(root, "research", "worker", rid, "conclusion", "Legacy ambiguous result.")
+ _write(
+ _root(root) / "replies" / rid / "conclusion.delivery.json",
+ {"status": "verification_required", "error": "provider_delivery_unverified"},
+ )
+
+ class Transport:
+ send_calls = 0
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ self.send_calls += 1
+ return {"reply_verified": True}
+
+ def verify(self, *_args):
+ raise AssertionError("a locatorless return has nothing to verify")
+
+ 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"] == "provider_locator_unavailable"
+ # Terminal: no later pump may resend the already published conclusion.
+ drain(
+ root,
+ registry,
+ ChatSessionStore(root),
+ transport,
+ now=datetime.now(timezone.utc) + timedelta(days=2),
+ )
+ assert transport.send_calls == 0
+
+
+def accepted_attempt():
+ return {
+ "schema_version": "manager_return_delivery_attempt_v0",
+ "provider": "lark",
+ "message_ref": "om_provider_reply",
+ "intent_digest": "sha256:" + "a" * 64,
+ "provider_receipt": "sha256:" + "b" * 64,
+ }
+
+
+@pytest.mark.parametrize(
+ "failure",
+ [
+ lambda: ReturnResolutionBlocked(
+ "return_authorization_unavailable", "context return authority revoked"
+ ),
+ lambda: ValueError("context return authority revoked"),
+ lambda: ValueError("manager connection no longer authorized"),
+ ],
+ ids=["typed-reason", "prose-revoked", "prose-unauthorized"],
+)
+def test_return_authority_revocation_terminalizes_without_repeating_readback(
+ flow, failure
+):
+ root, registry, store, create = flow
+ _, _, 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
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ record_attempt(accepted_attempt())
+ return {"external_write_performed": True, "reply_verified": False}
+
+ def verify(self, *_args):
+ self.verify_calls += 1
+ raise failure()
+
+ 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 transport.verify_calls == 1
+ # A terminal reason is never re-read, not even after the backoff window.
+ drain(
+ root,
+ registry,
+ ChatSessionStore(root),
+ transport,
+ now=datetime.now(timezone.utc) + timedelta(days=2),
+ )
+ assert transport.verify_calls == 1
+
+
+def test_unclassified_verification_failure_backs_off_instead_of_hot_looping(flow):
+ root, registry, store, create = flow
+ _, _, 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
+
+ def send_with_attempt(self, route, session, turn, text, record_attempt):
+ record_attempt(accepted_attempt())
+ return {"external_write_performed": True, "reply_verified": False}
+
+ def verify(self, *_args):
+ self.verify_calls += 1
+ raise RuntimeError("provider readback transport failed")
+
+ transport = Transport()
+ drain(root, registry, store, transport)
+ drain(root, registry, ChatSessionStore(root), transport)
+ assert reply_status(root, receipt)[0]["status"] == "verification_required"
+ assert transport.verify_calls == 1
+ # The kept locator stays retryable, but the next pump must respect the
+ # backoff instead of re-running the provider readback every few seconds.
+ drain(root, registry, ChatSessionStore(root), transport)
+ assert transport.verify_calls == 1
+ drain(
+ root,
+ registry,
+ ChatSessionStore(root),
+ transport,
+ now=datetime.now(timezone.utc) + timedelta(days=2),
+ )
+ assert transport.verify_calls == 2
+
+
+def test_public_delivery_projection_normalizes_unknown_private_state(flow):
+ root, _, _, create = flow
+ _, _, receipt = create(True)
+ rid = receipt["request_id"]
+ acknowledge(root, "research", "worker", rid, "adopt", "Checked")
+ report(root, "research", "worker", rid, "conclusion", "Bounded result.")
+ _write(
+ _root(root) / "replies" / rid / "conclusion.delivery.json",
+ {"status": "verification_required", "error": "private provider detail"},
+ )
+
+ state = reply_status(root, receipt)[0]
+ assert state["status"] == "explicit_unverified"
+ assert state["error"] == "delivery_state_unreadable"
+ assert "private" not in str(state)
def test_background_service_delivers_without_another_agent_or_query(flow):