From 025be07e3a17c312d5eb57d306ef2887c60a70f0 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sun, 4 Oct 2026 16:17:43 +0800 Subject: [PATCH] fix(collaboration): commit inbox reads after response validation read_inbox recorded request and peer-return deliveries before response enrichment completed. Enrichment then re-resolved a reusable Goal alias, so recreation could read a replacement workspace while leaving success receipts behind.\n\nCapture the admitted Goal snapshot for readiness checks, then re-enter that Goal lifetime before committing receipts. Regression tests cover recreation and enrichment failures. Signed-off-by: duanjialing.777 --- loopx/control_plane/collaboration/peers.py | 127 ++++++++++++++----- tests/test_collaboration_goal_instance.py | 139 ++++++++++++++++++++- 2 files changed, 231 insertions(+), 35 deletions(-) diff --git a/loopx/control_plane/collaboration/peers.py b/loopx/control_plane/collaboration/peers.py index 0a8ca74aac..009cda2abd 100644 --- a/loopx/control_plane/collaboration/peers.py +++ b/loopx/control_plane/collaboration/peers.py @@ -227,9 +227,9 @@ def request( } -def returns(root, goal_id, agent_id, *, mark_read=False, scope=None): - """Re-offer results until the requester explicitly acknowledges consumption.""" +def _collect_returns(root, goal_id, agent_id, *, scope=None): items = [] + deliveries = [] folder = ( _root(root) / "peer-operations" @@ -328,32 +328,49 @@ def returns(root, goal_id, agent_id, *, mark_read=False, scope=None): ) if len(items) > 20: break - if mark_read: - state = path.with_name(path.stem + ".delivery.json") - with _request_lock( - root, - row["request_id"], - scope, - path.with_suffix(".lock"), - ): - if not state.exists(): - _write( - state, - { - "status": "delivered", - **({"result_key": path.stem} if "result_key" in reply else {}), - "delivered_at": _now(), - "kind": "requester_cli_read", - **( - {"goal_ref": row["goal_ref"]} - if "goal_ref" in row - else {} - ), - }, - ) + deliveries.append((path, row, reply)) if len(items) > 20: break - return {"items": items[:20], "has_more": len(items) > 20} + return {"items": items[:20], "has_more": len(items) > 20}, deliveries + + +def _record_return_reads(root, deliveries, *, scope=None): + for path, row, reply in deliveries: + state = path.with_name(path.stem + ".delivery.json") + with _request_lock( + root, + row["request_id"], + scope, + path.with_suffix(".lock"), + ): + if not state.exists(): + _write( + state, + { + "status": "delivered", + **({"result_key": path.stem} if "result_key" in reply else {}), + "delivered_at": _now(), + "kind": "requester_cli_read", + **( + {"goal_ref": row["goal_ref"]} + if "goal_ref" in row + else {} + ), + }, + ) + + +def returns(root, goal_id, agent_id, *, mark_read=False, scope=None): + """Re-offer results until the requester explicitly acknowledges consumption.""" + result, deliveries = _collect_returns( + root, + goal_id, + agent_id, + scope=scope, + ) + if mark_read: + _record_return_reads(root, deliveries, scope=scope) + return result def consume_return( @@ -493,16 +510,15 @@ def _request_id(value): return value -def input_readiness( +def _input_readiness_for_goal( registry, goal_id, + goal, brief, *, workspace=None, configured_workspace: bool = False, ): - """Check local input versions, without fetching or claiming agent comprehension.""" - goal = _goal(registry, goal_id) goal_workspace = Path(goal["repo"]).resolve() selected = goal_workspace if workspace is not None and Path(workspace).resolve() != goal_workspace: @@ -562,6 +578,25 @@ def input_readiness( return result +def input_readiness( + registry, + goal_id, + brief, + *, + workspace=None, + configured_workspace: bool = False, +): + """Check local input versions, without fetching or claiming agent comprehension.""" + return _input_readiness_for_goal( + registry, + goal_id, + _goal(registry, goal_id), + brief, + workspace=workspace, + configured_workspace=configured_workspace, + ) + + def read_inbox( root, registry, @@ -590,16 +625,20 @@ def read_inbox( operation_cursor=operation_cursor, scope=goal_scope, ) - peer_returns = returns( + peer_returns, return_deliveries = _collect_returns( root, goal_id, agent_id, - mark_read=True, scope=goal_scope, ) if peer_returns["items"]: result["peer_returns"] = peer_returns - record_read(root, result["items"], scope=goal_scope) + else: + result.pop("peer_returns", None) + observed_goal = dict(goal_scope.goal) + observed_goal_ref = ( + dict(goal_scope.caller_goal_ref or {}) if goal_scope.exact else None + ) from .links import receiver_followthrough @@ -608,9 +647,29 @@ def read_inbox( # in Goal lifetime admission. for item in result["items"]: if item.get("brief"): - item["input_readiness"] = input_readiness( - registry, goal_id, item["brief"], workspace=workspace + item["input_readiness"] = _input_readiness_for_goal( + registry, + goal_id, + observed_goal, + item["brief"], + workspace=workspace, ) + if result["items"] or return_deliveries: + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(agent_id,), + caller_goal_ref=observed_goal_ref, + ) as goal_scope: + if result["items"]: + record_read(root, result["items"], scope=goal_scope) + else: + decide_collaboration_lifecycle( + goal_scope, + operation="read_record", + record=return_deliveries[0][1], + ) + _record_return_reads(root, return_deliveries, scope=goal_scope) result["followthrough"] = ( "Use each request's receiver_followthrough to reconcile it with actual Core work. " "Record an explicit assessment even when continuing other work; a read is not a decision. " diff --git a/tests/test_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py index 09fe75b12a..55b25232c5 100644 --- a/tests/test_collaboration_goal_instance.py +++ b/tests/test_collaboration_goal_instance.py @@ -1,3 +1,4 @@ +import hashlib import json import subprocess import sys @@ -95,7 +96,7 @@ def _http_snapshot(root, registry, session_id): server.server_close() -def _recreate(registry: Path) -> None: +def _recreate(registry: Path, *, repo: Path | None = None) -> None: with exclusive_cross_runtime_file_lock( guard_path(registry, "delivery"), operation="test_recreate_goal", @@ -106,6 +107,8 @@ def _recreate(registry: Path) -> None: ) as transaction: payload = transaction.payload_copy() payload["goals"][0]["goal_instance_id"] = INSTANCE_B + if repo is not None: + payload["goals"][0]["repo"] = str(repo) transaction.commit(payload) @@ -293,6 +296,140 @@ def test_recreated_goal_cannot_observe_or_mutate_prior_instance_requests( ) +def test_inbox_read_rejects_recreation_before_response_commit( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + from loopx.control_plane.collaboration import links + + registry = _create_source_registry(tmp_path) + input_path = tmp_path / "input.txt" + input_path.write_text("instance A", encoding="utf-8") + brief = { + **_brief("Read only the originating instance"), + "inputs": [ + { + "ref": input_path.name, + "description": "Instance-bound input", + "sha256": hashlib.sha256(input_path.read_bytes()).hexdigest(), + } + ], + } + receipt = request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + "read-before-recreation", + brief, + ) + replacement = tmp_path / "replacement" + replacement.mkdir() + (replacement / input_path.name).write_text("instance B", encoding="utf-8") + render_followthrough = links.receiver_followthrough + + def recreate_before_input_read(root, registry_path, items): + rendered = render_followthrough(root, registry_path, items) + _recreate(registry, repo=replacement) + return rendered + + monkeypatch.setattr(links, "receiver_followthrough", recreate_before_input_read) + + with pytest.raises(ValueError, match="historical_mutation_forbidden"): + read_inbox( + tmp_path, + registry, + "delivery", + "reviewer", + caller_goal_ref={ + "goal_id": "delivery", + "goal_instance_id": INSTANCE_A, + }, + ) + + assert not (_root(tmp_path) / "reads" / f"{receipt['request_id']}.json").exists() + + +@pytest.mark.parametrize("interruption", ["enrichment_failure", "goal_recreation"]) +def test_inbox_read_does_not_commit_peer_return_before_response_succeeds( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + interruption: str, +) -> None: + from loopx.control_plane.collaboration import links + from loopx.control_plane.collaboration.peers import return_result + + registry = _create_source_registry(tmp_path) + receipt = request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + "return-before-response", + _brief("Return only after a successful read"), + ) + goal_ref = { + "goal_id": "delivery", + "goal_instance_id": INSTANCE_A, + } + read_inbox( + tmp_path, + registry, + "delivery", + "reviewer", + caller_goal_ref=goal_ref, + ) + acknowledge( + tmp_path, + "delivery", + "reviewer", + receipt["request_id"], + "adopt", + "Review the current instance.", + registry=registry, + caller_goal_ref=goal_ref, + ) + return_result( + tmp_path, + "delivery", + "reviewer", + receipt["request_id"], + "Review complete.", + registry=registry, + caller_goal_ref=goal_ref, + ) + + def interrupt_response(*_args): + if interruption == "goal_recreation": + _recreate(registry) + return None + raise RuntimeError("response enrichment failed") + + monkeypatch.setattr(links, "receiver_followthrough", interrupt_response) + + expected_error = ValueError if interruption == "goal_recreation" else RuntimeError + expected_message = ( + "historical_mutation_forbidden" + if interruption == "goal_recreation" + else "response enrichment failed" + ) + with pytest.raises(expected_error, match=expected_message): + read_inbox( + tmp_path, + registry, + "delivery", + "builder", + caller_goal_ref=goal_ref, + ) + + delivery = ( + _root(tmp_path) / "replies" / receipt["request_id"] / "conclusion.delivery.json" + ) + assert not delivery.exists() + + def test_historical_request_does_not_borrow_work_from_recreated_goal(tmp_path, monkeypatch): from loopx.control_plane.collaboration import links