From 9e0925a44697907aebef4f1bdf04e35cfe1392b6 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 10:19:39 +0800 Subject: [PATCH] fix(collaboration): prepare return route before publishing requests Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../capable-manager-semantic-handoff-v0.md | 12 ++ .../capabilities/manager_context/__init__.py | 8 +- .../capabilities/manager_context/roundtrip.py | 3 + tests/test_context_handoff_publication.py | 164 ++++++++++++++++++ 4 files changed, 185 insertions(+), 2 deletions(-) create mode 100644 tests/test_context_handoff_publication.py diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md index 688711d86..9bc87c0ca 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md @@ -626,6 +626,18 @@ 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. +**Request publication and return preparation (S1/S3, A8/A9/A24):** persist +and verify the trusted original-conversation route before publishing an Inbox +entry. An independent receiver may consume the entry before the sender receives +its acknowledgement. A failed route write/readback must leave no visible work; +a prepared route without an entry is inert and an exact retry can finish it. +Keep the existing request identity and conflict checks. This matches the peer +request's route-before-entry ordering; it changes Python filesystem IO ordering, +not the shared TypeScript request or return-state authority. Exercise App, +Goal Chat and external-audience ingress, immediate receiver results, interrupted +publication and original-transcript return after restart. This qualification +does not establish native worker selection, execution or the full A24 journey. + ## 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/loopx/capabilities/manager_context/__init__.py b/loopx/capabilities/manager_context/__init__.py index de14582e8..b183c02cb 100644 --- a/loopx/capabilities/manager_context/__init__.py +++ b/loopx/capabilities/manager_context/__init__.py @@ -160,13 +160,17 @@ def deliver( exists = path.exists() if exists and {k: v for k, v in _read(path).items() if k not in {"delivered_at", "source_channel"}} != value: raise ValueError("context request identity conflict") + # The entry rename publishes work to independently scheduled receivers. + # Persist and verify its exact return route first, as peer requests do. + # A route without an entry is inert and recoverable by the same retry; + # an entry without a route can execute but cannot return its result. + from .roundtrip import register + register(runtime_root, value, session, turn) if not exists: from .tracking import _now _write(path, value | {"delivered_at": _now(), "source_channel": session.get("channel_id")}) if {k: v for k, v in _read(path).items() if k not in {"delivered_at", "source_channel"}} != value: raise ValueError("context delivery readback failed") - from .roundtrip import register - register(runtime_root, value, session, turn) return { "request_id": request_id, "status": "delivered", diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index 7a4e77994..21ca14c8b 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -140,6 +140,9 @@ def register(root, row, session, turn): raise ValueError("context return route conflict") else: _write(path, value | {"registered_at": _now()}) + saved = _read(path) + if any(saved.get(k) != v for k, v in value.items()): + raise ValueError("context return route readback failed") def _route(root, row): diff --git a/tests/test_context_handoff_publication.py b/tests/test_context_handoff_publication.py new file mode 100644 index 000000000..f3759a22c --- /dev/null +++ b/tests/test_context_handoff_publication.py @@ -0,0 +1,164 @@ +"""A receiver must never observe delegated work without its exact return route.""" + +import json + +import pytest + +import loopx.capabilities.manager_context as delivery +import loopx.capabilities.manager_context.roundtrip as roundtrip +from loopx.chat_store import ChatSessionStore + + +@pytest.fixture(params=["manager", "goal.research", "manager.external.test"]) +def conversation(tmp_path, request): + root = tmp_path + registry = root / "registry.json" + registry.write_text(json.dumps({"goals": [{ + "id": "research", "repo": str(root), + "coordination": {"registered_agents": ["worker"]}, + }]})) + channel = request.param + external = channel.startswith("manager.external.") + store = ChatSessionStore(root) + session = store.create_session( + goal_id="research" if channel == "goal.research" else "loopx-manager", + agent_id="codex", adapter_kind="codex_app_server", + upstream_thread_id="fixture", channel_id=channel, + ) + turn, _ = store.create_turn( + session["session_id"], client_turn_id="owner-request", + message="Ask the research owner for a Chinese draft. Do not publish it.", + origin="lark" if external else "web", + ) + target = {"goal_id": "research", "agent_id": "worker"} + if external: + delivery._write(delivery._root(root) / "policy.json", { + "schema_version": delivery.POLICY_SCHEMA, + "sources": {channel: {"sender_ids": ["owner"], "targets": [target]}}, + }) + delivery.register_ingress( + root, session_id=session["session_id"], + client_turn_id=turn["client_turn_id"], channel=channel, + sender_id="owner", message=turn["message"], source_id="lark:request", + ) + return root, registry, store, session, turn, target + + +@pytest.mark.parametrize("failure", ["route_write", "route_readback", "entry_write"]) +def test_failed_preparation_never_exposes_work_and_exact_retry_recovers( + conversation, monkeypatch, failure, +): + root, registry, store, session, turn, target = conversation + original_read = roundtrip._read + + def fail_route_write(path, value): + raise OSError("return storage unavailable") + + def fail_route_readback(path): + raise OSError("return readback unavailable") + + def fail_entry_write(path, value): + raise OSError("inbox storage unavailable") + + with monkeypatch.context() as patch: + if failure == "route_write": + patch.setattr(roundtrip, "_write", fail_route_write) + elif failure == "route_readback": + patch.setattr(roundtrip, "_read", fail_route_readback) + else: + patch.setattr(delivery, "_write", fail_entry_write) + with pytest.raises(OSError): + delivery.deliver(root, registry, session=session, turn=turn, request=target) + assert delivery.pending(root, **target)["items"] == [] + + receipt = delivery.deliver(root, registry, session=session, turn=turn, request=target) + assert not receipt["replayed"] + assert len(delivery.pending(root, **target)["items"]) == 1 + assert delivery.deliver(root, registry, session=session, turn=turn, request=target) == { + **receipt, "replayed": True, + } + route = original_read(delivery._root(root) / "roundtrips" / (receipt["request_id"] + ".json")) + assert route["session_id"] == session["session_id"] + assert route["client_turn_id"] == turn["client_turn_id"] + assert route["channel_id"] == session["channel_id"] + + +def test_receiver_can_report_at_publication_before_sender_gets_receipt( + conversation, monkeypatch, +): + root, registry, store, session, turn, target = conversation + original_write = delivery._write + returned_text = "已完成中文草稿,保留原约束,尚未发布。" + + def receive_immediately(path, value): + original_write(path, value) + # The atomic entry rename is the publication boundary. An independently + # scheduled receiver can read here, before deliver() returns to Chat. + entry, = delivery.pending(root, **target)["items"] + delivery.acknowledge(root, **target, request_id=entry["request_id"], + decision="adopt", reason="Keep the no-publication constraint") + roundtrip.report(root, "research", "worker", entry["request_id"], + "conclusion", returned_text) + + with monkeypatch.context() as patch: + patch.setattr(delivery, "_write", receive_immediately) + receipt = delivery.deliver(root, registry, session=session, turn=turn, request=target) + store.update_turn(session["session_id"], turn["turn_id"], status="completed", + response={"message": "Delivered", "context_handoff_receipt": receipt}) + store.finalize_managed_turn_completion(session["session_id"], turn["turn_id"]) + sends = [] + + def sender(*args, **kwargs): + sends.append((args, kwargs)) + return {"sent": True} + + # App readback survives restarting the store and needs no second model Turn. + if not session["channel_id"].startswith("manager.external."): + for _ in range(2): + roundtrip.drain(root, registry, ChatSessionStore(root), sender) + returned = [row for row in store.messages(session["session_id"]) + if row.get("origin") == "manager_followup"] + assert len(returned) == 1 and returned_text in returned[0]["text"] + assert returned[0]["turn_id"] == turn["turn_id"] + assert sends == [] + else: + assert roundtrip.reply_status(root, receipt)[0]["status"] == "queued" + + +def test_lost_acknowledgement_leaves_a_returnable_request_and_replays_once( + conversation, monkeypatch, +): + root, registry, _, session, turn, target = conversation + original_write = delivery._write + + def lose_ack(path, value): + original_write(path, value) + raise OSError("publication acknowledgement lost") + + with monkeypatch.context() as patch: + patch.setattr(delivery, "_write", lose_ack) + with pytest.raises(OSError): + delivery.deliver(root, registry, session=session, turn=turn, request=target) + entry, = delivery.pending(root, **target)["items"] + route = roundtrip._route(root, entry) + assert route["session_id"] == session["session_id"] + receipt = delivery.deliver(root, registry, session=session, turn=turn, request=target) + assert receipt["replayed"] and receipt["request_id"] == entry["request_id"] + assert len(delivery.pending(root, **target)["items"]) == 1 + + +def test_conflicting_return_route_cannot_publish_work(conversation, monkeypatch): + root, registry, _, session, turn, target = conversation + original_write = roundtrip._write + + def wrong_destination(path, value): + original_write(path, {**value, "session_id": "unrelated-conversation"}) + + with monkeypatch.context() as patch: + patch.setattr(roundtrip, "_write", wrong_destination) + with pytest.raises(ValueError, match="return route readback failed"): + delivery.deliver(root, registry, session=session, turn=turn, request=target) + assert delivery.pending(root, **target)["items"] == [] + with pytest.raises(ValueError, match="return route conflict"): + delivery.deliver(root, registry, session=session, turn=turn, request=target) + assert delivery.pending(root, **target)["items"] == []