Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
127 changes: 93 additions & 34 deletions loopx/control_plane/collaboration/peers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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

Expand All @@ -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. "
Expand Down
139 changes: 138 additions & 1 deletion tests/test_collaboration_goal_instance.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import hashlib
import json
import subprocess
import sys
Expand Down Expand Up @@ -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",
Expand All @@ -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)


Expand Down Expand Up @@ -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

Expand Down
Loading