From 61fb7c74fd89fec413ee977bc96d4cbb3e98532b Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sat, 26 Sep 2026 21:41:23 +0800 Subject: [PATCH 1/4] feat(collaboration): bind inbox work to goal instances Signed-off-by: duanjialing.777 --- .../capabilities/manager_context/__init__.py | 189 +++-- .../capabilities/manager_context/roundtrip.py | 594 ++++++++++++++- .../capabilities/manager_context/tracking.py | 90 ++- loopx/cli.py | 13 +- loopx/cli_commands/manager_inbox.py | 14 +- loopx/collaboration_mcp.py | 137 +++- .../collaboration/goal_instance_lifecycle.ts | 250 +++++++ .../collaboration/goal_instance_scope.py | 191 +++++ loopx/control_plane/collaboration/inbox.py | 220 +++++- loopx/control_plane/collaboration/peers.py | 455 +++++++++--- .../control_plane/effect_runtime_handlers.ts | 5 + ...laboration_goal_instance_lifecycle.test.ts | 180 +++++ tests/test_collaboration_goal_instance.py | 685 ++++++++++++++++++ 13 files changed, 2782 insertions(+), 241 deletions(-) create mode 100644 loopx/control_plane/collaboration/goal_instance_lifecycle.ts create mode 100644 loopx/control_plane/collaboration/goal_instance_scope.py create mode 100644 tests/control_plane_ts/collaboration_goal_instance_lifecycle.test.ts create mode 100644 tests/test_collaboration_goal_instance.py diff --git a/loopx/capabilities/manager_context/__init__.py b/loopx/capabilities/manager_context/__init__.py index de14582e89..c3ebcb60fe 100644 --- a/loopx/capabilities/manager_context/__init__.py +++ b/loopx/capabilities/manager_context/__init__.py @@ -8,15 +8,27 @@ from ...agent_registry import registered_agent_ids_for_goal from ...file_lock import exclusive_file_lock -from ...history import load_registry from ...control_plane.collaboration import conversation_scope +from ...control_plane.collaboration.goal_instance_scope import ( + collaboration_goal_scope, + decide_collaboration_lifecycle, +) from ...control_plane.goals.activation import goal_is_stopped +from ...control_plane.projects.registry_codec import load_project_registry # Retained imports are the shipped manager-context API; the shared owner is neutral. from ...control_plane.collaboration.inbox import ( - ENTRY_SCHEMA as ENTRY_SCHEMA, _hash as _hash, _read as _read, - _root as _root, _write as _write, normalize_request as normalize_request, - pending as pending, acknowledge as acknowledge, + ENTRY_SCHEMA as ENTRY_SCHEMA, + EXACT_ENTRY_SCHEMA, + _hash as _hash, + _read as _read, + _request_lock, + _root as _root, + _target, + _write as _write, + acknowledge as acknowledge, + normalize_request as normalize_request, + pending as pending, ) POLICY_SCHEMA = "loopx_manager_context_policy_v1" @@ -71,7 +83,7 @@ def authority( if registry_path is None: return {"mode": "unavailable", "targets": []} try: - registry = load_registry(registry_path) + registry = load_project_registry(registry_path) if not isinstance(registry, dict): raise ValueError("invalid registry") except (OSError, ValueError, TypeError): @@ -132,51 +144,118 @@ def deliver( runtime_root: Path, registry_path: Path, *, session: dict, turn: dict, request: dict ) -> dict: request = normalize_request(request) - grant = authority(runtime_root, registry_path, session, turn) target = {key: request[key] for key in ("goal_id", "agent_id")} - if target not in grant["targets"]: - raise ValueError("context recipient is not authorized or registered") - content = str(turn.get("message") or "") - if session.get("channel_id", "").startswith("manager.external."): - ingress = _read( + with collaboration_goal_scope( + registry_path, + goal_id=request["goal_id"], + agents=(request["agent_id"],), + require_active=True, + ) as goal_scope: + decide_collaboration_lifecycle( + goal_scope, + operation="request_create", + ) + grant = authority(runtime_root, registry_path, session, turn) + if target not in grant["targets"]: + raise ValueError("context recipient is not authorized or registered") + content = str(turn.get("message") or "") + if session.get("channel_id", "").startswith("manager.external."): + ingress = _read( + _root(runtime_root) + / "ingress" + / (_hash([session["session_id"], turn["client_turn_id"]]) + ".json") + ) + content = str(ingress["source_message"]) + if not content.strip() or len(content) > 20_000: + raise ValueError("invalid context content") + exact_target = _target( + request["goal_id"], + request["agent_id"], + goal_scope, + ) + request_id = _hash([grant["source_id"], exact_target]) + value = { + "schema_version": ( + EXACT_ENTRY_SCHEMA if goal_scope.exact else ENTRY_SCHEMA + ), + "request_id": request_id, + **request, + **goal_scope.record_identity(), + "source_id": grant["source_id"], + "message": content, + "instruction": INSTRUCTION, + } + path = ( _root(runtime_root) - / "ingress" - / (_hash([session["session_id"], turn["client_turn_id"]]) + ".json") + / "entries" + / _hash(exact_target) + / (request_id + ".json") ) - content = str(ingress["source_message"]) - if not content.strip() or len(content) > 20_000: - raise ValueError("invalid context content") - request_id = _hash([grant["source_id"], target]) - value = { - "schema_version": ENTRY_SCHEMA, - "request_id": request_id, - **request, - "source_id": grant["source_id"], - "message": content, - "instruction": INSTRUCTION, - } - path = _root(runtime_root) / "entries" / _hash(target) / (request_id + ".json") - with exclusive_file_lock(path.with_suffix(".lock")): - 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") - 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", - "replayed": exists, - "goal_id": request["goal_id"], - "agent_id": request["agent_id"], - "priority_changed": False, - "todo_created": False, - "execution_interrupted": False, - } + + def persist_entry() -> bool: + exists = path.exists() + if ( + exists + and { + key: item + for key, item in _read(path).items() + if key not in {"delivered_at", "source_channel"} + } + != value + ): + raise ValueError("context request identity conflict") + if not exists: + from .tracking import _now + + _write( + path, + value + | { + "delivered_at": _now(), + "source_channel": session.get("channel_id"), + }, + ) + if ( + { + key: item + for key, item in _read(path).items() + if key not in {"delivered_at", "source_channel"} + } + != value + ): + raise ValueError("context delivery readback failed") + return exists + + from .roundtrip import _register_unlocked, register + + if goal_scope.exact: + with _request_lock( + runtime_root, + request_id, + goal_scope, + path.with_suffix(".lock"), + ): + exists = persist_entry() + _register_unlocked(runtime_root, value, session, turn) + else: + with exclusive_file_lock(path.with_suffix(".lock")): + exists = persist_entry() + register(runtime_root, value, session, turn) + return { + "request_id": request_id, + "status": "delivered", + "replayed": exists, + "goal_id": request["goal_id"], + "agent_id": request["agent_id"], + **( + {"goal_ref": dict(goal_scope.caller_goal_ref or {})} + if goal_scope.exact + else {} + ), + "priority_changed": False, + "todo_created": False, + "execution_interrupted": False, + } def turn_start_hook( @@ -189,7 +268,17 @@ def turn_start_hook( def produce(): try: - inbox = pending(runtime_root, goal_id, agent_id) + with collaboration_goal_scope( + registry_path, + goal_id=goal_id, + agents=(agent_id,), + ) as goal_scope: + inbox = pending( + runtime_root, + goal_id, + agent_id, + scope=goal_scope, + ) count = len(inbox["items"]) + len(inbox.get("peer_returns", {}).get("items", [])) status, error = ("observed" if count else "empty"), None except (OSError, ValueError): @@ -272,7 +361,7 @@ def configure_evidence_scope(runtime_root: Path, registry_path: Path, *, channel """Local operator grants only selected Goal summaries to an exact audience.""" if not re.fullmatch(r"manager\.external\.[a-f0-9]{24}", channel): raise ValueError("an exact external manager channel is required") - registry = load_registry(registry_path) + registry = load_project_registry(registry_path) available = {g.get("id") for g in registry.get("goals", []) if isinstance(g, dict)} if any(g not in available for g in goal_ids): raise ValueError("every read Goal must be registered") @@ -314,7 +403,7 @@ def is_target(item: dict) -> bool: return item.get("goal_id") == goal_id and item.get("agent_id") == agent_id if grant: - registry = load_registry(registry_path) + registry = load_project_registry(registry_path) goal = next( (g for g in registry.get("goals", []) if isinstance(g, dict) and g.get("id") == goal_id), None, diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index fdbb1be71d..6c1da6ab65 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -10,18 +10,31 @@ import re import threading from datetime import datetime, timezone, timedelta +from uuid import uuid4 from . import _root, _read, _write, _hash, authority from .tracking import _entry, _now from ...file_lock import exclusive_file_lock from ...control_plane.collaboration import conversation_scope +from ...control_plane.collaboration.goal_instance_scope import ( + collaboration_goal_scope, + decide_collaboration_lifecycle, +) +from ...control_plane.projects.registry_codec import ( + SOURCE_SESSION_PROFILE_ID, + load_project_registry, +) from ...presentation.public_safety import scan_public_boundary_text from ...control_plane.effect_runtime import EffectRuntimeRejected, effect_runtime_result -from ...control_plane.collaboration.inbox import needs_conclusion as needs_conclusion +from ...control_plane.collaboration.inbox import ( + _request_lock, + needs_conclusion as needs_conclusion, +) PHASES = ("decision", "conclusion") DELIVERY_STATUSES = { + "admitted", "queued", "retry_pending", "verification_required", @@ -29,6 +42,7 @@ "superseded", "explicit_unverified", } +EXACT_DELIVERY_ADMISSION_SECONDS = 60 DELIVERY_ERRORS = { "provider_delivery_unverified", "provider_locator_unavailable", @@ -124,27 +138,43 @@ def _verification_exception_error(exc): return None -def register(root, row, session, turn): - """Called by trusted Chat delivery, never with model-authored routing.""" - value = {k: row[k] for k in ("request_id", "goal_id", "agent_id", "source_id")} +def _register_unlocked(root, row, session, turn): + value = { + key: row[key] + for key in ("request_id", "goal_id", "agent_id", "source_id", "goal_ref") + if key in row + } value.update( session_id=session["session_id"], client_turn_id=turn["client_turn_id"], channel_id=session["channel_id"], ) path = _root(root) / "roundtrips" / (row["request_id"] + ".json") - with exclusive_file_lock(path.with_suffix(".lock")): - if path.exists(): - old = _read(path) - if any(old.get(k) != v for k, v in value.items()): - raise ValueError("context return route conflict") - else: - _write(path, value | {"registered_at": _now()}) + if path.exists(): + old = _read(path) + if any(old.get(key) != value for key, value in value.items()): + raise ValueError("context return route conflict") + else: + _write(path, value | {"registered_at": _now()}) + + +def register(root, row, session, turn, *, scope=None): + """Called by trusted Chat delivery, never with model-authored routing.""" + path = _root(root) / "roundtrips" / (row["request_id"] + ".json") + with _request_lock( + root, + row["request_id"], + scope, + path.with_suffix(".lock"), + ): + _register_unlocked(root, row, session, turn) -def _route(root, row): +def _route(root, row, *, scope=None): path = _root(root) / "roundtrips" / (row["request_id"] + ".json") if not path.exists(): + if scope is not None and scope.exact: + raise ValueError("exact context return route unavailable") # Legacy opt-in only when the receiver actually publishes a reply. Recover # exact trusted Chat receipts, never infer a destination from Goal alone. from ...chat_store import ChatSessionStore @@ -174,23 +204,69 @@ def _route(root, row): raise ValueError("peer return route identity mismatch") if any( value.get(k) != row.get(k) - for k in ("request_id", "goal_id", "agent_id", "source_id") + for k in ("request_id", "goal_id", "agent_id", "source_id", "goal_ref") + if k in row ): raise ValueError("context return route identity mismatch") return value -def report(root, goal_id, agent_id, request_id, phase, text): +def report( + root, + goal_id, + agent_id, + request_id, + phase, + text, + *, + registry=None, + caller_goal_ref=None, + scope=None, +): """Chat audience adapter; the shared Inbox owns result validation/persistence.""" - row = _entry(root, goal_id, agent_id, request_id) - route = _route(root, row) + if scope is None and registry is not None: + from ...control_plane.collaboration.goal_instance_scope import ( + collaboration_goal_scope, + ) + + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + return report( + root, + goal_id, + agent_id, + request_id, + phase, + text, + scope=goal_scope, + ) + row = _entry(root, goal_id, agent_id, request_id, scope=scope) + route = _route(root, row, scope=scope) if route.get("kind") == "peer" and phase != "conclusion": raise ValueError("peer replies require a conclusion; report a concrete result or blocker") if (route["channel_id"] != "peer" and not conversation_scope(route)["private_conversation"] and not scan_public_boundary_text(text)["ok"]): raise ValueError("reply contains private boundary material; write an audience-safe conclusion") from ...control_plane.collaboration.inbox import record_result - return {**record_result(root, row, phase, text), "status": "queued_for_requester" if route.get("kind") == "peer" else "queued_for_original_conversation"} + return { + **record_result( + root, + row, + phase, + text, + scope=scope, + route=route, + ), + "status": ( + "queued_for_requester" + if route.get("kind") == "peer" + else "queued_for_original_conversation" + ), + } def reply_status(root, row): @@ -289,9 +365,493 @@ def project_chat_session_snapshot(root, store, session_id): return snapshot +def _initial_delivery_proved(row, route, turn): + receipt = (turn.get("response") or {}).get("context_handoff_receipt") or {} + return ( + turn.get("status") == "completed" + and receipt.get("request_id") == row["request_id"] + and receipt.get("goal_id") == row["goal_id"] + and receipt.get("agent_id") == row["agent_id"] + and receipt.get("goal_ref") == row["goal_ref"] + and route.get("goal_ref") == row["goal_ref"] + ) + + +def _exact_return_scope(registry, reply): + return collaboration_goal_scope( + registry, + goal_id=reply["goal_id"], + agents=(), + caller_goal_ref=reply["goal_ref"], + ) + + +def _exact_return_context(root, registry, store, path, state_path, now): + reply = _read(path) + if not isinstance(reply.get("goal_ref"), dict): + raise ValueError("exact return reply is missing goal_ref") + with _exact_return_scope(registry, reply) as scope: + row = _entry( + root, + reply["goal_id"], + reply["agent_id"], + reply["request_id"], + scope=scope, + ) + if ( + path.parent.name != row["request_id"] + or reply.get("phase") != path.stem + or any( + reply.get(key) != row.get(key) + for key in ("source_id", "goal_ref") + ) + ): + raise ValueError("return_reply_identity_mismatch") + route = _route(root, row, scope=scope) + if route.get("kind") == "peer": + return None + session = store.load_session(route["session_id"]) + turn = store.turn_for_client(route["session_id"], route["client_turn_id"]) + if ( + not session + or session.get("status") == "closed" + or session.get("channel_id") != route["channel_id"] + or not turn + ): + raise ValueError("original_conversation_unavailable") + grant = authority(root, registry, session, turn) + target = {key: row[key] for key in ("goal_id", "agent_id")} + if ( + target not in grant["targets"] + or grant.get("source_id") != row["source_id"] + ): + raise ValueError("return_authorization_unavailable") + initial_delivery_proved = _initial_delivery_proved(row, route, turn) + decide_collaboration_lifecycle( + scope, + operation="original_return_admit", + record=row, + route=route, + initial_delivery_proved=initial_delivery_proved, + ) + with _request_lock( + root, + row["request_id"], + scope, + path.with_suffix(".lock"), + ): + state = _read(state_path) if state_path.exists() else {} + if state.get("status") in { + "delivered", + "superseded", + "explicit_unverified", + }: + return None + admission = state.get("admission") + if ( + isinstance(admission, dict) + and isinstance(admission.get("expires_at"), str) + and now.isoformat() < admission["expires_at"] + ): + return None + if ( + state.get("retry_at") + and state.get("status") != "admitted" + and now.isoformat() < state["retry_at"] + ): + return None + if path.stem == "decision" and ( + path.parent / "conclusion.json" + ).exists(): + _write( + state_path, + { + "status": "superseded", + "reason": "conclusion_ready", + "goal_ref": row["goal_ref"], + }, + ) + return None + token = uuid4().hex + prior_status = ( + admission.get("prior_status") + if isinstance(admission, dict) + else state.get("status", "queued") + ) + admitted = { + **state, + "status": "admitted", + "goal_ref": row["goal_ref"], + "admission": { + "token": token, + "prior_status": prior_status, + "admitted_at": now.isoformat(), + "expires_at": ( + now + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS) + ).isoformat(), + }, + } + admitted.pop("retry_at", None) + _write(state_path, admitted) + return { + "reply": reply, + "row": row, + "route": route, + "session": session, + "turn": turn, + "store": store, + "state_path": state_path, + "state": admitted, + "token": token, + } + + +def _write_exact_return_state( + root, + registry, + context, + value, + *, + preserve_admission=False, +): + reply = context["reply"] + with _exact_return_scope(registry, reply) as scope: + row = _entry( + root, + reply["goal_id"], + reply["agent_id"], + reply["request_id"], + scope=scope, + ) + route = _route(root, row, scope=scope) + turn = context["store"].turn_for_client( + route["session_id"], + route["client_turn_id"], + ) + decide_collaboration_lifecycle( + scope, + operation="original_return_settle", + record=row, + route=route, + initial_delivery_proved=( + isinstance(turn, dict) + and _initial_delivery_proved(row, route, turn) + ), + ) + with _request_lock( + root, + row["request_id"], + scope, + context["state_path"].with_suffix(".lock"), + ): + current = _read(context["state_path"]) + admission = current.get("admission") + if ( + not isinstance(admission, dict) + or admission.get("token") != context["token"] + ): + raise ValueError("exact return delivery admission changed") + result = {**value, "goal_ref": row["goal_ref"]} + if preserve_admission: + result["admission"] = admission + else: + result.pop("admission", None) + _write(context["state_path"], result) + + +def _retry_state(state, now, *, error): + attempts = int(state.get("attempts", 0)) + 1 + result = { + key: value + for key, value in state.items() + if key not in {"admission", "delivered_at", "message_id", "verification"} + } + result.update( + status="retry_pending", + attempts=attempts, + error=error, + retry_at=( + now + timedelta(seconds=min(300, 5 * 2 ** min(attempts, 6))) + ).isoformat(), + ) + return result + + +def _drain_exact(root, registry, store, external_sender, *, now, cancelled): + processed = 0 + for path in sorted((_root(root) / "replies").glob("*/*.json")): + if cancelled(): + break + if path.stem not in PHASES: + continue + state_path = path.with_name(path.stem + ".delivery.json") + try: + context = _exact_return_context( + root, + registry, + store, + path, + state_path, + now, + ) + except (OSError, ValueError, KeyError, TypeError, RuntimeError): + logging.getLogger(__name__).warning( + "Exact manager return admission unavailable" + ) + continue + if context is None: + continue + row = context["row"] + route = context["route"] + session = context["session"] + turn = context["turn"] + reply = context["reply"] + state = context["state"] + prefix = "处理结论" if path.stem == "conclusion" else "处理进展" + text = ( + f"{prefix} · {row['agent_id']} · 委托 {row['request_id'][:8]}" + f"\n\n{reply['text']}" + ) + mid = "handoff." + _hash([row["request_id"], path.stem]) + try: + if cancelled(): + return processed + store.append_message( + route["session_id"], + role="agent", + text=text, + turn_id=turn["turn_id"], + origin="manager_followup", + message_id=mid, + ) + if cancelled(): + return processed + if conversation_scope(session)["private_conversation"]: + _write_exact_return_state( + root, + registry, + context, + { + "status": "delivered", + "delivered_at": now.isoformat(), + "message_id": mid, + }, + ) + processed += 1 + continue + + prior_status = state["admission"]["prior_status"] + if prior_status == "verification_required" or state.get("attempt") is not None: + attempt = state.get("attempt") + if attempt is None or _attempt_locator(attempt) is None: + _write_exact_return_state( + root, + registry, + context, + { + **({"attempt": attempt} if attempt is not None else {}), + "status": "explicit_unverified", + "error": "provider_locator_unavailable", + }, + ) + processed += 1 + continue + normalized_attempt = _delivery_attempt(attempt) + verifier = getattr(external_sender, "verify", None) + if not callable(verifier): + _write_exact_return_state( + root, + registry, + context, + { + "status": "explicit_unverified", + "error": "provider_verifier_unavailable", + }, + ) + processed += 1 + continue + try: + decision = _verification_decision( + verifier( + route, + session, + turn, + text, + normalized_attempt, + ) + ) + except ( + OSError, + ValueError, + KeyError, + TypeError, + RuntimeError, + ) as exc: + error = _verification_exception_error(exc) + next_state = ( + { + "status": "explicit_unverified", + "error": error, + } + if error + else _retry_state( + {"attempt": normalized_attempt, **state}, + now, + error="provider_verification_unavailable", + ) + ) + _write_exact_return_state( + root, + registry, + context, + next_state, + ) + processed += 1 + continue + if decision["status"] == "delivered": + next_state = { + "status": "delivered", + "delivered_at": now.isoformat(), + "message_id": mid, + "provider_receipt": normalized_attempt["provider_receipt"], + "reply_verified": True, + "verification": decision["verification"], + } + elif decision["status"] == "explicit_unverified": + next_state = { + "status": "explicit_unverified", + "error": decision["error"], + } + else: + next_state = _retry_state( + {"attempt": normalized_attempt, **state}, + now, + error=decision["error"], + ) + _write_exact_return_state( + root, + registry, + context, + next_state, + ) + processed += 1 + continue + + def record_attempt(value): + nonlocal state + attempt = _delivery_attempt(value) + existing = state.get("attempt") + if existing is not None and _delivery_attempt(existing) != attempt: + raise ValueError("manager return delivery attempt conflict") + state = {**state, "attempt": attempt} + _write_exact_return_state( + root, + registry, + context, + state, + preserve_admission=True, + ) + + sender = getattr(external_sender, "send_with_attempt", None) + sent = ( + sender(route, session, turn, text, record_attempt) + if callable(sender) + else external_sender(route, session, turn, text) + ) + if sent.get("reply_verified") is not True: + if sent.get("external_write_performed") is True: + attempt = state.get("attempt") + next_state = ( + { + "status": "verification_required", + "error": "provider_delivery_unverified", + "attempt": attempt, + } + if attempt is not None and _attempt_locator(attempt) is not None + else { + **({"attempt": attempt} if attempt is not None else {}), + "status": "explicit_unverified", + "error": "provider_locator_unavailable", + } + ) + _write_exact_return_state( + root, + registry, + context, + next_state, + ) + processed += 1 + continue + raise ValueError("return_transport_unavailable") + _write_exact_return_state( + root, + registry, + context, + { + "status": "delivered", + "delivered_at": now.isoformat(), + "message_id": mid, + "provider_receipt": sent.get("idempotency_key"), + "reply_verified": True, + }, + ) + except (OSError, ValueError, KeyError, TypeError, RuntimeError) as exc: + attempt = state.get("attempt") + error = ( + _verification_exception_error(exc) + if attempt is not None + else None + ) + if error: + next_state = { + "status": "explicit_unverified", + "error": error, + } + elif attempt is not None and _attempt_locator(attempt) is None: + next_state = { + "status": "explicit_unverified", + "error": "provider_locator_unavailable", + "attempt": attempt, + } + else: + next_state = _retry_state( + state, + now, + error="original_route_or_return_delivery_unavailable", + ) + try: + _write_exact_return_state( + root, + registry, + context, + next_state, + ) + except (OSError, ValueError, KeyError, TypeError, RuntimeError): + logging.getLogger(__name__).warning( + "Exact manager return settlement unavailable" + ) + processed += 1 + if processed >= 20: + break + return processed + + def drain(root, registry, store, external_sender, *, now=None, cancelled=lambda: False): """Restart-safe return delivery; transport retries never rerun the worker/model.""" now = now or datetime.now(timezone.utc) + try: + profile_id = load_project_registry(registry).get("profile_id") + except (OSError, ValueError, TypeError): + profile_id = None + if profile_id == SOURCE_SESSION_PROFILE_ID: + return _drain_exact( + root, + registry, + store, + external_sender, + now=now, + cancelled=cancelled, + ) processed = 0 for path in sorted((_root(root) / "replies").glob("*/*.json")): if cancelled(): diff --git a/loopx/capabilities/manager_context/tracking.py b/loopx/capabilities/manager_context/tracking.py index ef282cabd1..29a1ff862a 100644 --- a/loopx/capabilities/manager_context/tracking.py +++ b/loopx/capabilities/manager_context/tracking.py @@ -9,9 +9,16 @@ from . import _read, _root, _write from ...control_plane.collaboration.inbox import ( - _now as _now, _entry as _entry, _receipt as _receipt, record_read as record_read, + _entry as _entry, + _now as _now, + _receipt as _receipt, + _request_lock, + record_read as record_read, +) +from ...control_plane.collaboration.goal_instance_scope import ( + collaboration_goal_scope, + decide_collaboration_lifecycle, ) -from ...file_lock import exclusive_file_lock from ...todos import list_goal_todos from ...chat_manager_details import _text @@ -25,8 +32,18 @@ def _core_todos(registry_path, root, goal_id): return {r["todo_id"]: r for r in result.get("todos", []) if r.get("todo_id")} -def link(root, registry_path, goal_id, agent_id, request_id, todo_ids, evidence_ids): - _entry(root, goal_id, agent_id, request_id) +def link( + root, + registry_path, + goal_id, + agent_id, + request_id, + todo_ids, + evidence_ids, + *, + caller_goal_ref=None, + scope=None, +): if not todo_ids and not evidence_ids: raise ValueError("at least one Core Todo or evidence reference required") if len(todo_ids) > 16 or len(evidence_ids) > 16: @@ -41,12 +58,35 @@ def link(root, registry_path, goal_id, agent_id, request_id, todo_ids, evidence_ row = rows.get(tid, {}) if row.get("claimed_by") != agent_id and row.get("bound_agent") != agent_id: raise ValueError("linked Todo must belong to the receiving Agent") + if scope is None: + with collaboration_goal_scope( + registry_path, + goal_id=goal_id, + agents=(agent_id,), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + return link( + root, + registry_path, + goal_id, + agent_id, + request_id, + todo_ids, + evidence_ids, + scope=goal_scope, + ) + row = _entry(root, goal_id, agent_id, request_id, scope=scope) + decide_collaboration_lifecycle( + scope, + operation="artifact_link", + record=row, + ) path = _root(root) / "links" / (request_id + ".json") - with exclusive_file_lock(path.with_suffix(".lock")): + with _request_lock(root, request_id, scope, path.with_suffix(".lock")): old, error = _receipt( root, "links", - {"request_id": request_id, "goal_id": goal_id, "agent_id": agent_id}, + row, ) if error: raise ValueError(error) @@ -54,13 +94,12 @@ def link(root, registry_path, goal_id, agent_id, request_id, todo_ids, evidence_ refs = sorted(set(old.get("evidence_ids", [])) | set(evidence_ids)) if len(tids) > 16 or len(refs) > 16: raise ValueError("too many context links") - value = dict( - request_id=request_id, - goal_id=goal_id, - agent_id=agent_id, - todo_ids=tids, - evidence_ids=refs, - ) + value = { + key: row[key] + for key in ("request_id", "goal_id", "agent_id", "goal_ref") + if key in row + } + value.update(todo_ids=tids, evidence_ids=refs) if any(old.get(k) != v for k, v in value.items()): _write(path, value | {"updated_at": _now()}) return {"ok": True, **value} @@ -125,7 +164,24 @@ def query( and legacy_channels.get(row.get("source_id")) == {channel_id} ): continue - _entry(root, row["goal_id"], row["agent_id"], row["request_id"]) + with collaboration_goal_scope( + registry_path, + goal_id=row["goal_id"], + agents=(), + caller_goal_ref=row.get("goal_ref"), + ) as goal_scope: + entry = _entry( + root, + row["goal_id"], + row["agent_id"], + row["request_id"], + scope=goal_scope, + ) + decide_collaboration_lifecycle( + goal_scope, + operation="history_inspect", + record=entry, + ) rows.append(row) except (OSError, ValueError, KeyError, TypeError): unreadable += 1 @@ -161,7 +217,11 @@ def query( else "core_todo_unavailable_or_not_found", } ) - item = {k: row[k] for k in ("request_id", "goal_id", "agent_id", "source_id")} + item = { + key: row[key] + for key in ("request_id", "goal_id", "agent_id", "source_id", "goal_ref") + if key in row + } if row.get("source_kind") == "peer": item.update(source_kind="peer", source_agent_id=row["source_agent_id"], parent_request_id=row.get("parent_request_id")) if owner_scope and row.get("brief"): diff --git a/loopx/cli.py b/loopx/cli.py index e1bc9700a4..22f132dfec 100644 --- a/loopx/cli.py +++ b/loopx/cli.py @@ -810,7 +810,18 @@ def main(argv: list[str] | None = None) -> int: ) if args.command == "manager-inbox": - return handle_manager_inbox(args, registry_path, effective_runtime_root(registry_path, args.runtime_root)) + from .control_plane.projects.registry_codec import load_project_registry + from .paths import resolve_runtime_root + + return handle_manager_inbox( + args, + registry_path, + resolve_runtime_root( + load_project_registry(registry_path), + args.runtime_root, + registry_path=registry_path, + ), + ) if args.command == "delegation": return handle_delegation(args, registry_path, effective_runtime_root(registry_path, args.runtime_root)) diff --git a/loopx/cli_commands/manager_inbox.py b/loopx/cli_commands/manager_inbox.py index 836e54ead7..3239f433ab 100644 --- a/loopx/cli_commands/manager_inbox.py +++ b/loopx/cli_commands/manager_inbox.py @@ -3,12 +3,12 @@ import json from pathlib import Path from ..agent_registry import registered_agent_ids_for_goal -from ..history import load_registry from ..capabilities.manager_context import ( acknowledge, configure_delivery_target, configure_evidence_scope, ) +from ..control_plane.projects.registry_codec import load_project_registry def register_manager_inbox(subparsers, add_format): @@ -90,7 +90,7 @@ def handle_manager_inbox(args, registry_path, runtime_root): ) print(json.dumps(result, ensure_ascii=False, indent=2)) return 0 - registry = load_registry(registry_path) + registry = load_project_registry(registry_path) goal = next( (g for g in registry.get("goals", []) if g.get("id") == args.goal_id), None ) @@ -108,7 +108,13 @@ def handle_manager_inbox(args, registry_path, runtime_root): args.peer_agent_id, args.operation_id, json.loads(raw), args.parent_request_id) elif args.manager_inbox_action == "acknowledge-return": from ..control_plane.collaboration.peers import consume_return - result = consume_return(runtime_root, args.goal_id, args.agent_id, args.request_id) + result = consume_return( + runtime_root, + args.goal_id, + args.agent_id, + args.request_id, + registry=registry_path, + ) elif args.manager_inbox_action == "read": from ..control_plane.collaboration.peers import read_inbox result = read_inbox(runtime_root, registry_path, args.goal_id, args.agent_id, @@ -126,6 +132,7 @@ def handle_manager_inbox(args, registry_path, runtime_root): args.request_id or "", args.phase, args.reply_text or "", + registry=registry_path, ) elif args.manager_inbox_action == "link": from ..capabilities.manager_context.tracking import link @@ -163,6 +170,7 @@ def handle_manager_inbox(args, registry_path, runtime_root): args.request_id or "", args.decision or "", args.reason or "", + registry=registry_path, ) except (OSError, ValueError) as exc: result = {"ok": False, "error": str(exc)} diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index c302dc9132..aee2bd5cca 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -40,6 +40,11 @@ from .control_plane.collaboration.inbox import _hash, _read, _write, _root, _receipt from .control_plane.collaboration.peers import return_result from .control_plane.collaboration.inbox import acknowledge, _entry, normalize_request +from .control_plane.collaboration.goal_instance_scope import ( + capture_collaboration_goal_ref, + collaboration_goal_scope, + decide_collaboration_lifecycle, +) from .control_plane.collaboration import delegation_results from .control_plane.collaboration.peers import ( _goal, @@ -188,7 +193,11 @@ def create_server( def register_collaboration_tools(server: FastMCP, root: Path, registry: Path, goal_id: str, agent_id: str, workspace: Path) -> None: - _goal(registry, goal_id, agent_id) + caller_goal_ref = capture_collaboration_goal_ref( + registry, + goal_id=goal_id, + agent_id=agent_id, + ) def check_scope(): # Revocation is read on every tool call, including a long-lived server. @@ -201,7 +210,15 @@ def read_context(cursor: str | None = None) -> dict: Follow next_cursor for later requests. Omit cursor to start a fresh scan. Pages are live; reading all pages does not complete outstanding work. """ - return read_inbox(root, registry, goal_id, agent_id, workspace=workspace, cursor=cursor) + return read_inbox( + root, + registry, + goal_id, + agent_id, + workspace=workspace, + cursor=cursor, + caller_goal_ref=caller_goal_ref, + ) @server.tool() def assess_request( @@ -211,7 +228,16 @@ def assess_request( ) -> dict: """Record your independent decision; this does not change task ownership or priority.""" check_scope() - return acknowledge(root, goal_id, agent_id, request_id, decision, reason) + return acknowledge( + root, + goal_id, + agent_id, + request_id, + decision, + reason, + registry=registry, + caller_goal_ref=caller_goal_ref, + ) @server.tool() def request_peer( @@ -237,28 +263,40 @@ def request_peer( operation_id, brief, parent_request_id, + caller_goal_ref=caller_goal_ref, ) @server.tool() def return_result(request_id: str, text: str) -> dict: """Save an evidence-backed conclusion or explicit blocker for the original requester.""" check_scope() - row = _entry(root, goal_id, agent_id, request_id) - if row.get("source_kind") == "peer": - from .control_plane.collaboration.peers import return_result as save_result - - return save_result(root, goal_id, agent_id, request_id, text) # The host adapter selects Chat/Lark transport; the shared collaboration # owner never depends on presentation or manager capabilities. from .capabilities.manager_context.roundtrip import report - return report(root, goal_id, agent_id, request_id, "conclusion", text) + return report( + root, + goal_id, + agent_id, + request_id, + "conclusion", + text, + registry=registry, + caller_goal_ref=caller_goal_ref, + ) @server.tool() def consume_peer_result(request_id: str) -> dict: """Acknowledge a peer result after reading and using/rejecting it; no work-state mutation.""" check_scope() - return consume_return(root, goal_id, agent_id, request_id) + return consume_return( + root, + goal_id, + agent_id, + request_id, + registry=registry, + caller_goal_ref=caller_goal_ref, + ) class Delegations: @@ -271,6 +309,11 @@ class Delegations: def __init__(self, root: Path, registry: Path, goal_id: str, agent_id: str, config: Path): self.root, self.registry = root.resolve(), registry.resolve() self.goal_id, self.agent_id, self.config = goal_id, agent_id, config.resolve() + self.goal_ref = capture_collaboration_goal_ref( + self.registry, + goal_id=self.goal_id, + agent_id=self.agent_id, + ) def binding(self, binding_id: str, *, require_active: bool = False) -> dict: _goal(self.registry, self.goal_id, self.agent_id, require_active=require_active) @@ -391,7 +434,8 @@ def start(self, binding_id: str, operation_id: str, brief: dict, if not exists: delegation_results.require_dependencies(self, binding, brief) delivered = request(self.root, self.registry, self.goal_id, self.agent_id, - binding["agent_id"], operation_id, brief, parent_request_id) + binding["agent_id"], operation_id, brief, parent_request_id, + caller_goal_ref=self.goal_ref) identity = {"binding": binding, "request_id": delivered["request_id"], "operation_id": operation_id} if exists: if _read(path).get("identity") != identity: @@ -719,20 +763,54 @@ def _record_turn_result( def _receiver_adopted(self, row: dict, binding: dict) -> bool: request_id = row["identity"]["request_id"] - decision, error = _receipt( - self.root, - "decisions", - _entry(self.root, self.goal_id, binding["agent_id"], request_id), - ) + with collaboration_goal_scope( + self.registry, + goal_id=self.goal_id, + agents=(), + caller_goal_ref=self.goal_ref, + ) as goal_scope: + entry = _entry( + self.root, + self.goal_id, + binding["agent_id"], + request_id, + scope=goal_scope, + ) + decide_collaboration_lifecycle( + goal_scope, + operation="history_inspect", + record=entry, + ) + decision, error = _receipt( + self.root, + "decisions", + entry, + ) return not error and bool(decision) and decision["decision"] == "adopt" def _delegation_bootstrap(self, row: dict, binding: dict) -> dict: request_id = row["identity"]["request_id"] + with collaboration_goal_scope( + self.registry, + goal_id=self.goal_id, + agents=(), + caller_goal_ref=self.goal_ref, + ) as goal_scope: + entry = _entry( + self.root, + self.goal_id, + binding["agent_id"], + request_id, + scope=goal_scope, + ) + decide_collaboration_lifecycle( + goal_scope, + operation="history_inspect", + record=entry, + ) return { "request_id": request_id, - "brief": _entry( - self.root, self.goal_id, binding["agent_id"], request_id - )["brief"], + "brief": entry["brief"], "instruction": ( "Use the loopx_delegation tools to read_context and call " "assess_request for this request before working. If you adopt " @@ -960,9 +1038,24 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: self._complete_delegated_todo(row, binding) row["artifacts"] = self._accepted(binding) if not (_root(self.root) / "replies" / request_id / "conclusion.json").exists(): - return_result(self.root, self.goal_id, binding["agent_id"], request_id, - json.dumps({"todo_id": binding["todo_id"], "status": "accepted", - "artifacts": [{k: v for k, v in item.items() if k != "text"} for item in row["artifacts"]]})) + return_result( + self.root, + self.goal_id, + binding["agent_id"], + request_id, + json.dumps( + { + "todo_id": binding["todo_id"], + "status": "accepted", + "artifacts": [ + {k: v for k, v in item.items() if k != "text"} + for item in row["artifacts"] + ], + } + ), + registry=self.registry, + caller_goal_ref=self.goal_ref, + ) self._observe(path, row, "accepted", canonical_done=True, acceptance_ready=True, artifacts_current=True) except (ValueError, KeyError, subprocess.TimeoutExpired, EffectRuntimeRemoteError) as exc: # Retain uncertain execution for explicit same-operation recovery. diff --git a/loopx/control_plane/collaboration/goal_instance_lifecycle.ts b/loopx/control_plane/collaboration/goal_instance_lifecycle.ts new file mode 100644 index 0000000000..3d62a7e9e3 --- /dev/null +++ b/loopx/control_plane/collaboration/goal_instance_lifecycle.ts @@ -0,0 +1,250 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { assertNever, jsonObject } from "../runtime_decode.ts"; +import { + parseExactGoalRef, + type ExactGoalRef, +} from "../goals/goal_instance_identity.ts"; + +const SOURCE_SESSION_PROFILE_ID = "source_session_v1"; + +const OPERATIONS = [ + "request_create", + "inbox_observe", + "read_record", + "receiver_decide", + "artifact_link", + "result_publish", + "peer_return_observe", + "peer_return_consume", + "original_return_admit", + "original_return_settle", + "history_inspect", +] as const; + +type CollaborationOperation = (typeof OPERATIONS)[number]; + +type WireGoalRef = Readonly<{ + goal_id: string; + goal_instance_id: string; +}>; + +type CollaborationLifecycleFacts = Readonly<{ + operation: CollaborationOperation; + callerGoalRef: ExactGoalRef; + currentGoalRef: ExactGoalRef; + recordGoalRef: ExactGoalRef | null; + routeGoalRef: ExactGoalRef | null; + initialDeliveryProved: boolean; +}>; + +export type CollaborationLifecycleDecision = + | Readonly<{ kind: "legacy" }> + | Readonly<{ + kind: "allow"; + mode: "current_instance" | "historical_result" | "historical_read"; + goal_ref: WireGoalRef; + }> + | Readonly<{ + kind: "omit"; + code: "different_goal_instance" | "legacy_unbound"; + }> + | Readonly<{ + kind: "reject"; + code: + | "unsupported_profile" + | "missing_exact_goal_ref" + | "stale_goal_instance" + | "record_instance_mismatch" + | "route_instance_mismatch" + | "historical_mutation_forbidden" + | "initial_delivery_unproved"; + }>; + +function operation(value: unknown): CollaborationOperation { + switch (value) { + case "request_create": + case "inbox_observe": + case "read_record": + case "receiver_decide": + case "artifact_link": + case "result_publish": + case "peer_return_observe": + case "peer_return_consume": + case "original_return_admit": + case "original_return_settle": + case "history_inspect": + return value; + default: + throw new EffectRuntimeRequestError( + "collaboration lifecycle operation is unsupported", + ); + } +} + +function optionalGoalRef(value: unknown): ExactGoalRef | null { + if (value === null) return null; + const parsed = parseExactGoalRef(value); + return parsed.kind === "parsed" ? parsed.value : null; +} + +function goalRefsEqual(left: ExactGoalRef, right: ExactGoalRef): boolean { + return left.goalId.value === right.goalId.value + && left.goalInstanceId.value === right.goalInstanceId.value; +} + +function wireGoalRef(value: ExactGoalRef): WireGoalRef { + return { + goal_id: value.goalId.value, + goal_instance_id: value.goalInstanceId.value, + }; +} + +function isObservation(operation: CollaborationOperation): boolean { + return operation === "inbox_observe" + || operation === "peer_return_observe"; +} + +function requiresRoute(operation: CollaborationOperation): boolean { + return operation === "result_publish" + || operation === "peer_return_observe" + || operation === "peer_return_consume" + || operation === "original_return_admit" + || operation === "original_return_settle"; +} + +function permitsHistorical(operation: CollaborationOperation): boolean { + return operation === "result_publish" + || operation === "original_return_admit" + || operation === "original_return_settle" + || operation === "history_inspect"; +} + +function lifecycleFacts( + value: JsonObject, + selectedOperation: CollaborationOperation, +): CollaborationLifecycleFacts | CollaborationLifecycleDecision { + const callerGoalRef = optionalGoalRef(value.caller_goal_ref); + const currentGoalRef = optionalGoalRef(value.current_goal_ref); + if (callerGoalRef === null || currentGoalRef === null) { + return { kind: "reject", code: "missing_exact_goal_ref" }; + } + const recordGoalRef = optionalGoalRef(value.record_goal_ref); + const routeGoalRef = optionalGoalRef(value.route_goal_ref); + return { + operation: selectedOperation, + callerGoalRef, + currentGoalRef, + recordGoalRef, + routeGoalRef, + initialDeliveryProved: value.initial_delivery_proved === true, + }; +} + +function decideExactLifecycle( + facts: CollaborationLifecycleFacts, +): CollaborationLifecycleDecision { + const { + operation, + callerGoalRef, + currentGoalRef, + recordGoalRef, + routeGoalRef, + } = facts; + if (operation === "request_create") { + return goalRefsEqual(callerGoalRef, currentGoalRef) + ? { + kind: "allow", + mode: "current_instance", + goal_ref: wireGoalRef(callerGoalRef), + } + : { kind: "reject", code: "stale_goal_instance" }; + } + + if (recordGoalRef === null) { + return isObservation(operation) + ? { kind: "omit", code: "legacy_unbound" } + : { kind: "reject", code: "missing_exact_goal_ref" }; + } + if (!goalRefsEqual(recordGoalRef, callerGoalRef)) { + return isObservation(operation) + ? { kind: "omit", code: "different_goal_instance" } + : { kind: "reject", code: "record_instance_mismatch" }; + } + if (requiresRoute(operation)) { + if (routeGoalRef === null || !goalRefsEqual(routeGoalRef, callerGoalRef)) { + return { kind: "reject", code: "route_instance_mismatch" }; + } + } + + const current = goalRefsEqual(callerGoalRef, currentGoalRef); + if (current) { + if ( + (operation === "original_return_admit" + || operation === "original_return_settle") + && !facts.initialDeliveryProved + ) { + return { kind: "reject", code: "initial_delivery_unproved" }; + } + return { + kind: "allow", + mode: "current_instance", + goal_ref: wireGoalRef(callerGoalRef), + }; + } + if (!permitsHistorical(operation)) { + return isObservation(operation) + ? { kind: "omit", code: "different_goal_instance" } + : { kind: "reject", code: "historical_mutation_forbidden" }; + } + if ( + (operation === "original_return_admit" + || operation === "original_return_settle") + && !facts.initialDeliveryProved + ) { + return { kind: "reject", code: "initial_delivery_unproved" }; + } + switch (operation) { + case "result_publish": + case "original_return_admit": + case "original_return_settle": + return { + kind: "allow", + mode: "historical_result", + goal_ref: wireGoalRef(callerGoalRef), + }; + case "history_inspect": + return { + kind: "allow", + mode: "historical_read", + goal_ref: wireGoalRef(callerGoalRef), + }; + case "inbox_observe": + case "read_record": + case "receiver_decide": + case "artifact_link": + case "peer_return_observe": + case "peer_return_consume": + return { kind: "reject", code: "historical_mutation_forbidden" }; + default: + return assertNever(operation, "unsupported collaboration operation"); + } +} + +export function decideCollaborationLifecycle( + value: unknown, +): CollaborationLifecycleDecision { + const raw = jsonObject(value); + if (!raw) { + throw new EffectRuntimeRequestError( + "collaboration lifecycle facts must be an object", + ); + } + const selectedOperation = operation(raw.operation); + if (raw.profile_id === null) return { kind: "legacy" }; + if (raw.profile_id !== SOURCE_SESSION_PROFILE_ID) { + return { kind: "reject", code: "unsupported_profile" }; + } + const facts = lifecycleFacts(raw, selectedOperation); + return "operation" in facts ? decideExactLifecycle(facts) : facts; +} diff --git a/loopx/control_plane/collaboration/goal_instance_scope.py b/loopx/control_plane/collaboration/goal_instance_scope.py new file mode 100644 index 0000000000..66581e212e --- /dev/null +++ b/loopx/control_plane/collaboration/goal_instance_scope.py @@ -0,0 +1,191 @@ +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager +from dataclasses import dataclass +from pathlib import Path +from typing import Any + +from ...agent_registry import registered_agent_ids_for_goal +from ...file_lock import exclusive_cross_runtime_file_lock +from ..effect_runtime import effect_runtime_result +from ..goals.source_session_registry_state import exact_goal_ref, guard_path +from ..projects.registry_codec import ( + SOURCE_SESSION_PROFILE_ID, + load_project_registry, +) + + +@dataclass(frozen=True, slots=True) +class CollaborationGoalScope: + registry_path: Path + goal_id: str + goal: dict[str, Any] + profile_id: str | None + caller_goal_ref: dict[str, str] | None + current_goal_ref: dict[str, str] | None + + @property + def exact(self) -> bool: + return self.profile_id == SOURCE_SESSION_PROFILE_ID + + def target(self, agent_id: str) -> dict[str, Any]: + if self.exact: + return { + "goal_ref": dict(self.caller_goal_ref or {}), + "agent_id": agent_id, + } + return {"goal_id": self.goal_id, "agent_id": agent_id} + + def record_identity(self) -> dict[str, Any]: + result: dict[str, Any] = {"goal_id": self.goal_id} + if self.exact: + result["goal_ref"] = dict(self.caller_goal_ref or {}) + return result + + +def _registered_goal( + registry: dict[str, Any], + *, + goal_id: str, + agents: tuple[str, ...], + require_active: bool, +) -> dict[str, Any]: + goal = next( + ( + candidate + for candidate in registry.get("goals", []) + if isinstance(candidate, dict) and candidate.get("id") == goal_id + ), + None, + ) + if goal is None or any( + agent not in registered_agent_ids_for_goal(goal) for agent in agents + ): + raise ValueError("collaboration requires registered Agents of the same Goal") + if require_active and goal.get("status") in {"stopped", "archived"}: + raise ValueError("collaboration Goal is stopped or archived") + return goal + + +def _source_goal_ref(goal_id: str, goal: dict[str, Any]) -> dict[str, str]: + instance_id = goal.get("goal_instance_id") + if not isinstance(instance_id, str): + raise ValueError("source-session Goal is missing goal_instance_id") + return exact_goal_ref(goal_id, instance_id) + + +def _caller_ref( + *, + goal_id: str, + requested: dict[str, str] | None, + current: dict[str, str], +) -> dict[str, str]: + if requested is None: + return dict(current) + return exact_goal_ref( + str(requested.get("goal_id") or goal_id), + str(requested.get("goal_instance_id") or ""), + ) + + +@contextmanager +def collaboration_goal_scope( + registry_path: Path, + *, + goal_id: str, + agents: tuple[str, ...], + caller_goal_ref: dict[str, str] | None = None, + require_active: bool = False, +) -> Iterator[CollaborationGoalScope]: + """Hold the alias lifetime guard while one collaboration operation commits.""" + + registry_path = Path(registry_path).expanduser().resolve() + with exclusive_cross_runtime_file_lock( + guard_path(registry_path, goal_id), + operation="collaboration_goal_lifetime", + ): + registry = load_project_registry(registry_path) + goal = _registered_goal( + registry, + goal_id=goal_id, + agents=agents, + require_active=require_active, + ) + profile = registry.get("profile_id") + if profile != SOURCE_SESSION_PROFILE_ID: + yield CollaborationGoalScope( + registry_path=registry_path, + goal_id=goal_id, + goal=goal, + profile_id=None, + caller_goal_ref=None, + current_goal_ref=None, + ) + return + current = _source_goal_ref(goal_id, goal) + yield CollaborationGoalScope( + registry_path=registry_path, + goal_id=goal_id, + goal=goal, + profile_id=SOURCE_SESSION_PROFILE_ID, + caller_goal_ref=_caller_ref( + goal_id=goal_id, + requested=caller_goal_ref, + current=current, + ), + current_goal_ref=current, + ) + + +def capture_collaboration_goal_ref( + registry_path: Path, + *, + goal_id: str, + agent_id: str, +) -> dict[str, str] | None: + with collaboration_goal_scope( + registry_path, + goal_id=goal_id, + agents=(agent_id,), + ) as scope: + return ( + dict(scope.caller_goal_ref or {}) + if scope.exact + else None + ) + + +def decide_collaboration_lifecycle( + scope: CollaborationGoalScope, + *, + operation: str, + record: dict[str, Any] | None = None, + route: dict[str, Any] | None = None, + initial_delivery_proved: bool = False, +) -> dict[str, Any]: + if not scope.exact: + return {"kind": "legacy"} + result = effect_runtime_result( + "collaboration.goal_instance.decide", + { + "profile_id": scope.profile_id, + "operation": operation, + "caller_goal_ref": scope.caller_goal_ref, + "current_goal_ref": scope.current_goal_ref, + "record_goal_ref": ( + record.get("goal_ref") if isinstance(record, dict) else None + ), + "route_goal_ref": ( + route.get("goal_ref") if isinstance(route, dict) else None + ), + "initial_delivery_proved": initial_delivery_proved, + }, + ) + if not isinstance(result, dict): + raise RuntimeError("collaboration lifecycle decision must be an object") + if result.get("kind") == "reject": + raise ValueError( + f"collaboration lifecycle rejected: {result.get('code')}" + ) + return result diff --git a/loopx/control_plane/collaboration/inbox.py b/loopx/control_plane/collaboration/inbox.py index 8f845a9c8e..81dfed15ce 100644 --- a/loopx/control_plane/collaboration/inbox.py +++ b/loopx/control_plane/collaboration/inbox.py @@ -11,12 +11,17 @@ import os import re import tempfile +from contextlib import contextmanager from datetime import datetime, timezone from pathlib import Path -from typing import Any +from typing import TYPE_CHECKING, Any from ...file_lock import exclusive_file_lock +if TYPE_CHECKING: + from .goal_instance_scope import CollaborationGoalScope + ENTRY_SCHEMA = "loopx_manager_context_entry_v1" +EXACT_ENTRY_SCHEMA = "loopx_manager_context_entry_v2" REQUEST_TRIAGE_INSTRUCTION = ( "Read the supplied requests and their source-specific instructions before choosing work. " "Independently assess context, evidence, constraints and costs. Requests and peer results " @@ -59,6 +64,45 @@ def _read(path: Path) -> dict: return value +def _target( + goal_id: str, + agent_id: str, + scope: CollaborationGoalScope | None, +) -> dict[str, Any]: + return ( + scope.target(agent_id) + if scope is not None + else {"goal_id": goal_id, "agent_id": agent_id} + ) + + +def _record_identity( + goal_id: str, + agent_id: str, + scope: CollaborationGoalScope | None, +) -> dict[str, Any]: + value: dict[str, Any] = {"goal_id": goal_id, "agent_id": agent_id} + if scope is not None and scope.exact: + value["goal_ref"] = dict(scope.caller_goal_ref or {}) + return value + + +@contextmanager +def _request_lock( + root: Path, + request_id: str, + scope: CollaborationGoalScope | None, + legacy_path: Path, +): + lock = ( + _root(root) / "request-locks" / request_id + if scope is not None and scope.exact + else legacy_path + ) + with exclusive_file_lock(lock): + yield + + def normalize_request(value: Any) -> dict | None: if value is None: return None @@ -73,22 +117,33 @@ def normalize_request(value: Any) -> dict | None: def pending( - runtime_root: Path, goal_id: str, agent_id: str, *, cursor: str | None = None + runtime_root: Path, + goal_id: str, + agent_id: str, + *, + cursor: str | None = None, + scope: CollaborationGoalScope | None = None, ) -> dict: - scope = _hash(["pending_requests_v1", str(runtime_root.resolve()), goal_id, agent_id]) + cursor_scope = _hash( + [ + "pending_requests_v1", + str(runtime_root.resolve()), + _target(goal_id, agent_id, scope), + ] + ) after = "" if cursor is not None: if not isinstance(cursor, str) or not re.fullmatch( r"1:[a-f0-9]{64}:[a-f0-9]{64}", cursor ): raise ValueError("invalid pending request cursor") - _, cursor_scope, after = cursor.split(":") - if cursor_scope != scope: + _, supplied_scope, after = cursor.split(":") + if supplied_scope != cursor_scope: raise ValueError("pending request cursor scope mismatch") folder = ( _root(runtime_root) / "entries" - / _hash(dict(goal_id=goal_id, agent_id=agent_id)) + / _hash(_target(goal_id, agent_id, scope)) ) try: paths = sorted(folder.iterdir()) @@ -107,12 +162,23 @@ def pending( continue item = _read(path) if ( - item.get("schema_version") != ENTRY_SCHEMA + item.get("schema_version") + != (EXACT_ENTRY_SCHEMA if scope is not None and scope.exact else ENTRY_SCHEMA) or item.get("goal_id") != goal_id or item.get("agent_id") != agent_id or item.get("request_id") != path.stem ): raise ValueError("context inbox scope mismatch") + if scope is not None: + from .goal_instance_scope import decide_collaboration_lifecycle + + lifecycle = decide_collaboration_lifecycle( + scope, + operation="inbox_observe", + record=item, + ) + if lifecycle.get("kind") == "omit": + continue if decided: item = { **item, @@ -124,13 +190,17 @@ def pending( break from .peers import returns - peer_returns = returns(runtime_root, goal_id, agent_id) + peer_returns = returns(runtime_root, goal_id, agent_id, scope=scope) return { "ok": True, **({"peer_returns": peer_returns} if peer_returns["items"] else {}), "items": items[:20], "has_more": len(items) > 20, - "next_cursor": f"1:{scope}:{items[19]['request_id']}" if len(items) > 20 else None, + "next_cursor": ( + f"1:{cursor_scope}:{items[19]['request_id']}" + if len(items) > 20 + else None + ), "instruction": ( REQUEST_TRIAGE_INSTRUCTION + " Follow next_cursor to read later pending requests. Restart without a cursor " @@ -146,24 +216,56 @@ def acknowledge( request_id: str, decision: str, reason: str, + *, + registry: Path | None = None, + caller_goal_ref: dict[str, str] | None = None, + scope: CollaborationGoalScope | None = None, ) -> dict: + if scope is None and registry is not None: + from .goal_instance_scope import collaboration_goal_scope + + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(agent_id,), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + return acknowledge( + runtime_root, + goal_id, + agent_id, + request_id, + decision, + reason, + scope=goal_scope, + ) if not re.fullmatch(r"[a-f0-9]{64}", request_id): raise ValueError("invalid context request id") - target = dict(goal_id=goal_id, agent_id=agent_id) - entry = _read( - _root(runtime_root) / "entries" / _hash(target) / (request_id + ".json") + target = _record_identity(goal_id, agent_id, scope) + entry = _entry( + runtime_root, + goal_id, + agent_id, + request_id, + scope=scope, ) - if any(entry.get(k) != v for k, v in target.items()): - raise ValueError("context inbox scope mismatch") if ( decision not in {"adopt", "defer", "reject", "no_change"} or not reason.strip() or len(reason) > 2000 ): raise ValueError("a bounded replan decision and reason are required") + if scope is not None: + from .goal_instance_scope import decide_collaboration_lifecycle + + decide_collaboration_lifecycle( + scope, + operation="receiver_decide", + record=entry, + ) value = {"request_id": request_id, **target, "decision": decision, "reason": reason} path = _root(runtime_root) / "decisions" / (request_id + ".json") - with exclusive_file_lock(path.with_suffix(".lock")): + with _request_lock(runtime_root, request_id, scope, path.with_suffix(".lock")): if path.exists(): if {k: v for k, v in _read(path).items() if k != "decided_at"} != value: raise ValueError("context decision already recorded") @@ -176,26 +278,58 @@ def _now(): return datetime.now(timezone.utc).isoformat() -def _entry(root, goal_id, agent_id, request_id): +def _entry(root, goal_id, agent_id, request_id, *, scope=None): if not isinstance(request_id, str) or not re.fullmatch(r"[a-f0-9]{64}", request_id): raise ValueError("invalid context request id") - target = dict(goal_id=goal_id, agent_id=agent_id) + target = _target(goal_id, agent_id, scope) + identity = _record_identity(goal_id, agent_id, scope) row = _read(_root(root) / "entries" / _hash(target) / (request_id + ".json")) - if any(row.get(k) != v for k, v in {**target, "request_id": request_id}.items()): + if any( + row.get(k) != v + for k, v in {**identity, "request_id": request_id}.items() + ): raise ValueError("context receipt scope mismatch") return row -def record_read(root: Path, items: list[dict]) -> None: +def record_read( + root: Path, + items: list[dict], + *, + scope: CollaborationGoalScope | None = None, +) -> None: """CLI supplied these messages to the receiver; not proof of comprehension.""" for row in items: - _entry(root, row["goal_id"], row["agent_id"], row["request_id"]) + _entry( + root, + row["goal_id"], + row["agent_id"], + row["request_id"], + scope=scope, + ) + if scope is not None: + from .goal_instance_scope import decide_collaboration_lifecycle + + decide_collaboration_lifecycle( + scope, + operation="read_record", + record=row, + ) path = _root(root) / "reads" / (row["request_id"] + ".json") - with exclusive_file_lock(path.with_suffix(".lock")): + with _request_lock( + root, + row["request_id"], + scope, + path.with_suffix(".lock"), + ): if not path.exists(): _write( path, - {k: row[k] for k in ("request_id", "goal_id", "agent_id")} + { + k: row[k] + for k in ("request_id", "goal_id", "agent_id", "goal_ref") + if k in row + } | { "read_at": _now(), "kind": "receiver_cli_read", @@ -210,7 +344,9 @@ def _receipt(root, lane, row): try: value = _read(path) if any( - value.get(k) != row.get(k) for k in ("goal_id", "agent_id", "request_id") + value.get(k) != row.get(k) + for k in ("goal_id", "agent_id", "request_id", "goal_ref") + if k in row ): raise ValueError("receipt identity conflict") if lane == "decisions" and value.get("decision") not in { @@ -251,7 +387,15 @@ def needs_conclusion(root, request_id): ).exists() -def record_result(root, row, phase, text): +def record_result( + root, + row, + phase, + text, + *, + scope: CollaborationGoalScope | None = None, + route: dict[str, Any] | None = None, +): """Persist a receiver conclusion; the transport adapter validates its audience.""" request_id = row["request_id"] decision, error = _receipt(root, "decisions", row) @@ -267,9 +411,33 @@ def record_result(root, row, phase, text): "a decision/conclusion phase and bounded reply text are required" ) text = text.strip() + if scope is not None: + from .goal_instance_scope import decide_collaboration_lifecycle + + decide_collaboration_lifecycle( + scope, + operation="result_publish", + record=row, + route=route, + ) path = _root(root) / "replies" / request_id / (phase + ".json") - with exclusive_file_lock((_root(root) / "replies" / request_id / "report.lock")): - value = {k: row[k] for k in ("request_id", "goal_id", "agent_id", "source_id")} + with _request_lock( + root, + request_id, + scope, + _root(root) / "replies" / request_id / "report.lock", + ): + value = { + k: row[k] + for k in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "goal_ref", + ) + if k in row + } value.update(phase=phase, text=text, decision=decision["decision"]) if path.exists(): old = _read(path) diff --git a/loopx/control_plane/collaboration/peers.py b/loopx/control_plane/collaboration/peers.py index 021d28586e..aa37f2530e 100644 --- a/loopx/control_plane/collaboration/peers.py +++ b/loopx/control_plane/collaboration/peers.py @@ -12,12 +12,26 @@ import stat from pathlib import Path -from .inbox import ENTRY_SCHEMA, _hash, _read, _root, _write, normalize_request -from .inbox import _entry, _now +from .inbox import ( + ENTRY_SCHEMA, + EXACT_ENTRY_SCHEMA, + _entry, + _hash, + _now, + _read, + _request_lock, + _root, + _target, + _write, + normalize_request, +) from . import conversation_scope +from .goal_instance_scope import ( + collaboration_goal_scope, + decide_collaboration_lifecycle, +) from ...agent_registry import registered_agent_ids_for_goal -from ...file_lock import exclusive_file_lock -from ...history import load_registry +from ..projects.registry_codec import load_project_registry PEER_INSTRUCTION = ( "This is a peer's request for help or independent review, not an owner instruction. " @@ -31,7 +45,11 @@ def _goal(registry, goal_id, *agents, require_active=False): goal = next( - (g for g in load_registry(registry).get("goals", []) if g.get("id") == goal_id), + ( + g + for g in load_project_registry(registry).get("goals", []) + if g.get("id") == goal_id + ), None, ) if not goal or any(a not in registered_agent_ids_for_goal(goal) for a in agents): @@ -57,115 +75,173 @@ def request( operation_id, brief, parent_request_id=None, + *, + caller_goal_ref=None, ): normalized = normalize_request( {"goal_id": goal_id, "agent_id": target_agent_id, "brief": brief} ) - _goal(registry, goal_id, source_agent_id, target_agent_id, require_active=True) if source_agent_id == target_agent_id: raise ValueError("a peer request requires a different receiving Agent") operation_id = require_operation_id(operation_id) - inherited = None - if parent_request_id: - parent = _entry(root, goal_id, source_agent_id, parent_request_id) - if parent.get("source_kind") != "peer": - scope = conversation_scope({ - "channel_id": parent.get("source_channel"), "goal_id": goal_id, - }, origin="web" if str(parent.get("source_id", "")).startswith("web:") else "unknown") - if not scope["private_conversation"]: - raise ValueError("external-audience requests cannot be forwarded to peers") - # The original owner context is preserved through a chain without growing - # a transcript recursively at each hop. - inherited = parent.get("inherited_context") or { - "request_id": parent_request_id, - "message": parent["message"], - **({"brief": parent["brief"]} if "brief" in parent else {}), + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(source_agent_id, target_agent_id), + caller_goal_ref=caller_goal_ref, + require_active=True, + ) as goal_scope: + decide_collaboration_lifecycle( + goal_scope, + operation="request_create", + ) + inherited = None + if parent_request_id: + parent = _entry( + root, + goal_id, + source_agent_id, + parent_request_id, + scope=goal_scope, + ) + if parent.get("source_kind") != "peer": + audience_scope = conversation_scope( + { + "channel_id": parent.get("source_channel"), + "goal_id": goal_id, + }, + origin=( + "web" + if str(parent.get("source_id", "")).startswith("web:") + else "unknown" + ), + ) + if not audience_scope["private_conversation"]: + raise ValueError( + "external-audience requests cannot be forwarded to peers" + ) + inherited = parent.get("inherited_context") or { + "request_id": parent_request_id, + "message": parent["message"], + **({"brief": parent["brief"]} if "brief" in parent else {}), + } + source_identity = ( + [goal_scope.caller_goal_ref, source_agent_id, operation_id] + if goal_scope.exact + else [goal_id, source_agent_id, operation_id] + ) + source_id = "peer:" + _hash(source_identity) + request_id = _hash( + [source_id, _target(goal_id, target_agent_id, goal_scope)] + ) + row = { + "schema_version": ( + EXACT_ENTRY_SCHEMA if goal_scope.exact else ENTRY_SCHEMA + ), + **normalized, + **goal_scope.record_identity(), + "request_id": request_id, + "source_id": source_id, + "source_kind": "peer", + "source_agent_id": source_agent_id, + "parent_request_id": parent_request_id, + "inherited_context": inherited, + "message": normalized["brief"]["purpose"], + "instruction": PEER_INSTRUCTION, } - source_id = "peer:" + _hash([goal_id, source_agent_id, operation_id]) - request_id = _hash([source_id, {"goal_id": goal_id, "agent_id": target_agent_id}]) - row = { - "schema_version": ENTRY_SCHEMA, - **normalized, - "request_id": request_id, - "source_id": source_id, - "source_kind": "peer", - "source_agent_id": source_agent_id, - "parent_request_id": parent_request_id, - "inherited_context": inherited, - "message": normalized["brief"]["purpose"], - "instruction": PEER_INSTRUCTION, - } - # Lock the operation, not its recipient: retargeting a retry is a conflict. - operation_path = ( - _root(root) - / "peer-operations" - / _hash({"goal_id": goal_id, "agent_id": source_agent_id}) - / (source_id[5:] + ".json") - ) - path = ( - _root(root) - / "entries" - / _hash({"goal_id": goal_id, "agent_id": target_agent_id}) - / (request_id + ".json") - ) - with exclusive_file_lock(operation_path.with_suffix(".lock")): - if operation_path.exists() and _read(operation_path) != row: - raise ValueError("peer request operation identity conflict") - if not operation_path.exists(): - _write(operation_path, row) - replayed = path.exists() - if ( - replayed - and {k: v for k, v in _read(path).items() if k != "delivered_at"} != row + operation_path = ( + _root(root) + / "peer-operations" + / _hash(_target(goal_id, source_agent_id, goal_scope)) + / (source_id[5:] + ".json") + ) + path = ( + _root(root) + / "entries" + / _hash(_target(goal_id, target_agent_id, goal_scope)) + / (request_id + ".json") + ) + with _request_lock( + root, + request_id, + goal_scope, + operation_path.with_suffix(".lock"), ): - raise ValueError("peer request identity conflict") - route = { - k: row[k] - for k in ( - "request_id", - "goal_id", - "agent_id", - "source_id", - "source_agent_id", - ) + if operation_path.exists() and _read(operation_path) != row: + raise ValueError("peer request operation identity conflict") + if not operation_path.exists(): + _write(operation_path, row) + replayed = path.exists() + if ( + replayed + and { + key: value + for key, value in _read(path).items() + if key != "delivered_at" + } + != row + ): + raise ValueError("peer request identity conflict") + route = { + key: row[key] + for key in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "source_agent_id", + "goal_ref", + ) + if key in row + } + route.update(kind="peer", channel_id="peer") + route_path = _root(root) / "roundtrips" / (request_id + ".json") + if route_path.exists() and _read(route_path) != route: + raise ValueError("peer return route identity conflict") + _write(route_path, route) + if not replayed: + _write(path, row | {"delivered_at": _now()}) + return { + "ok": True, + "request_id": request_id, + "goal_id": goal_id, + "agent_id": target_agent_id, + "status": "delivered", + "replayed": replayed, + "todo_created": False, + "priority_changed": False, + "execution_interrupted": False, } - route.update(kind="peer", channel_id="peer") - route_path = _root(root) / "roundtrips" / (request_id + ".json") - if route_path.exists() and _read(route_path) != route: - raise ValueError("peer return route identity conflict") - _write(route_path, route) - if not replayed: - _write(path, row | {"delivered_at": _now()}) - return { - "ok": True, - "request_id": request_id, - "goal_id": goal_id, - "agent_id": target_agent_id, - "status": "delivered", - "replayed": replayed, - "todo_created": False, - "priority_changed": False, - "execution_interrupted": False, - } -def returns(root, goal_id, agent_id, *, mark_read=False): +def returns(root, goal_id, agent_id, *, mark_read=False, scope=None): """Re-offer results until the requester explicitly acknowledges consumption.""" items = [] folder = ( _root(root) / "peer-operations" - / _hash({"goal_id": goal_id, "agent_id": agent_id}) + / _hash(_target(goal_id, agent_id, scope)) ) for operation_path in sorted(folder.glob("*.json")): operation = _read(operation_path) if ( operation.get("goal_id") != goal_id or operation.get("source_agent_id") != agent_id + or ( + scope is not None + and scope.exact + and operation.get("goal_ref") != scope.caller_goal_ref + ) ): raise ValueError("peer return scope mismatch") try: - row = _entry(root, goal_id, operation["agent_id"], operation["request_id"]) + row = _entry( + root, + goal_id, + operation["agent_id"], + operation["request_id"], + scope=scope, + ) except FileNotFoundError: continue # A reserved send without an entry is repaired by its exact retry. if row.get("source_agent_id") != agent_id or row.get( @@ -179,11 +255,28 @@ def returns(root, goal_id, agent_id, *, mark_read=False): if ( any( reply.get(k) != row.get(k) - for k in ("request_id", "goal_id", "agent_id", "source_id") + for k in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "goal_ref", + ) + if k in row ) or reply.get("phase") != "conclusion" ): raise ValueError("peer reply identity conflict") + route = _read(_root(root) / "roundtrips" / (row["request_id"] + ".json")) + if scope is not None: + lifecycle = decide_collaboration_lifecycle( + scope, + operation="peer_return_observe", + record=row, + route=route, + ) + if lifecycle.get("kind") == "omit": + continue consumed = path.parent / "conclusion.consumed.json" if consumed.exists(): value = _read(consumed) @@ -193,6 +286,11 @@ def returns(root, goal_id, agent_id, *, mark_read=False): "request_id": row["request_id"], "goal_id": goal_id, "agent_id": agent_id, + **( + {"goal_ref": row["goal_ref"]} + if "goal_ref" in row + else {} + ), }.items() ): raise ValueError("peer consumption receipt scope mismatch") @@ -212,7 +310,12 @@ def returns(root, goal_id, agent_id, *, mark_read=False): break if mark_read: state = path.with_name("conclusion.delivery.json") - with exclusive_file_lock(path.with_suffix(".lock")): + with _request_lock( + root, + row["request_id"], + scope, + path.with_suffix(".lock"), + ): if not state.exists(): _write( state, @@ -220,12 +323,40 @@ def returns(root, goal_id, agent_id, *, mark_read=False): "status": "delivered", "delivered_at": _now(), "kind": "requester_cli_read", + **( + {"goal_ref": row["goal_ref"]} + if "goal_ref" in row + else {} + ), }, ) return {"items": items[:20], "has_more": len(items) > 20} -def consume_return(root, goal_id, agent_id, request_id): +def consume_return( + root, + goal_id, + agent_id, + request_id, + *, + registry=None, + caller_goal_ref=None, + scope=None, +): + if scope is None and registry is not None: + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(agent_id,), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + return consume_return( + root, + goal_id, + agent_id, + request_id, + scope=goal_scope, + ) route = _read(_root(root) / "roundtrips" / (_request_id(request_id) + ".json")) if ( route.get("kind") != "peer" @@ -233,10 +364,24 @@ def consume_return(root, goal_id, agent_id, request_id): or route.get("source_agent_id") != agent_id ): raise ValueError("peer return scope mismatch") - row = _entry(root, goal_id, route["agent_id"], request_id) + row = _entry( + root, + goal_id, + route["agent_id"], + request_id, + scope=scope, + ) if row.get("source_kind") != "peer" or any( route.get(k) != row.get(k) - for k in ("request_id", "goal_id", "agent_id", "source_id", "source_agent_id") + for k in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "source_agent_id", + "goal_ref", + ) + if k in row ): raise ValueError("peer return scope mismatch") folder = _root(root) / "replies" / request_id @@ -245,19 +390,41 @@ def consume_return(root, goal_id, agent_id, request_id): reply = _read(folder / "conclusion.json") if reply.get("phase") != "conclusion" or any( reply.get(k) != row.get(k) - for k in ("request_id", "goal_id", "agent_id", "source_id") + for k in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "goal_ref", + ) + if k in row ): raise ValueError("peer reply identity conflict") - if _read(folder / "conclusion.delivery.json").get("status") != "delivered": + delivery = _read(folder / "conclusion.delivery.json") + if delivery.get("status") != "delivered" or ( + "goal_ref" in row and delivery.get("goal_ref") != row["goal_ref"] + ): raise ValueError("read the peer conclusion before acknowledging it") + if scope is not None: + decide_collaboration_lifecycle( + scope, + operation="peer_return_consume", + record=row, + route=route, + ) path = folder / "conclusion.consumed.json" - with exclusive_file_lock(path.with_suffix(".lock")): + with _request_lock(root, request_id, scope, path.with_suffix(".lock")): if path.exists() and any( _read(path).get(k) != v for k, v in { "request_id": request_id, "goal_id": goal_id, "agent_id": agent_id, + **( + {"goal_ref": row["goal_ref"]} + if "goal_ref" in row + else {} + ), }.items() ): raise ValueError("peer consumption receipt scope mismatch") @@ -269,6 +436,11 @@ def consume_return(root, goal_id, agent_id, request_id): "goal_id": goal_id, "agent_id": agent_id, "consumed_at": _now(), + **( + {"goal_ref": row["goal_ref"]} + if "goal_ref" in row + else {} + ), }, ) return { @@ -354,20 +526,50 @@ def input_readiness( return result -def read_inbox(root, registry, goal_id, agent_id, *, workspace=None, cursor=None): +def read_inbox( + root, + registry, + goal_id, + agent_id, + *, + workspace=None, + cursor=None, + caller_goal_ref=None, +): from .inbox import pending from .inbox import record_read - _goal(registry, goal_id, agent_id) - result = pending(root, goal_id, agent_id, cursor=cursor) + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(agent_id,), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + result = pending( + root, + goal_id, + agent_id, + cursor=cursor, + scope=goal_scope, + ) + peer_returns = 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) + + # Input hashing can touch arbitrary workspace files and does not participate + # 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 ) - peer_returns = returns(root, goal_id, agent_id, mark_read=True) - if peer_returns["items"]: - result["peer_returns"] = peer_returns result["followthrough"] = ( "Independently assess requests and actual input versions before accepting work. " "Use request_peer for help or independent review. Assess peer conclusions against " @@ -375,24 +577,63 @@ def read_inbox(root, registry, goal_id, agent_id, *, workspace=None, cursor=None "Finish the original request with return_result, including evidence and remaining gaps. " "Adoption, file hashes and returned opinions are not independent acceptance or Todo completion." ) - record_read(root, result["items"]) return result -def return_result(root, goal_id, agent_id, request_id, text): +def return_result( + root, + goal_id, + agent_id, + request_id, + text, + *, + registry=None, + caller_goal_ref=None, + scope=None, +): """Route by the saved recipient, never by an Agent's coordinator role.""" - row = _entry(root, goal_id, agent_id, request_id) + if scope is None and registry is not None: + with collaboration_goal_scope( + registry, + goal_id=goal_id, + agents=(), + caller_goal_ref=caller_goal_ref, + ) as goal_scope: + return return_result( + root, + goal_id, + agent_id, + request_id, + text, + scope=goal_scope, + ) + row = _entry(root, goal_id, agent_id, request_id, scope=scope) if row.get("source_kind") != "peer": raise ValueError("peer result requires a peer return route") route = _read(_root(root) / "roundtrips" / (request_id + ".json")) if route.get("kind") != "peer" or any( route.get(k) != row.get(k) - for k in ("request_id", "goal_id", "agent_id", "source_id", "source_agent_id") + for k in ( + "request_id", + "goal_id", + "agent_id", + "source_id", + "source_agent_id", + "goal_ref", + ) + if k in row ): raise ValueError("peer return route identity mismatch") from .inbox import record_result return { - **record_result(root, row, "conclusion", text), + **record_result( + root, + row, + "conclusion", + text, + scope=scope, + route=route, + ), "status": "queued_for_requester", } diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 730c3b9c8a..d1123553b1 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -221,6 +221,7 @@ import { classifyManagerReturnVerification, normalizeManagerReturnDeliveryAttempt, } from "./collaboration/return_delivery.ts"; +import { decideCollaborationLifecycle } from "./collaboration/goal_instance_lifecycle.ts"; import { normalizeCollaborationRequest } from "./collaboration/semantic_request.ts"; import { @@ -719,6 +720,10 @@ export function createEffectRuntimeHandlers( "collaboration.request.normalize", (params) => normalizeCollaborationRequest(params.request), ], + [ + "collaboration.goal_instance.decide", + (params) => decideCollaborationLifecycle(params), + ], ["external_evidence.discover", projectExternalEvidenceDiscovery], ["external_evidence.plan", planExternalEvidenceRequest], ["external_evidence.receipt", recordExternalEvidenceReceiptObservation], diff --git a/tests/control_plane_ts/collaboration_goal_instance_lifecycle.test.ts b/tests/control_plane_ts/collaboration_goal_instance_lifecycle.test.ts new file mode 100644 index 0000000000..122bdff2b8 --- /dev/null +++ b/tests/control_plane_ts/collaboration_goal_instance_lifecycle.test.ts @@ -0,0 +1,180 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { decideCollaborationLifecycle } from "../../loopx/control_plane/collaboration/goal_instance_lifecycle.ts"; + +const INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; +const INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; +const GOAL_A = { goal_id: "delivery", goal_instance_id: INSTANCE_A }; +const GOAL_B = { goal_id: "delivery", goal_instance_id: INSTANCE_B }; + +function facts( + operation: string, + patch: Record = {}, +): Record { + return { + profile_id: "source_session_v1", + operation, + caller_goal_ref: GOAL_A, + current_goal_ref: GOAL_A, + record_goal_ref: GOAL_A, + route_goal_ref: GOAL_A, + initial_delivery_proved: true, + ...patch, + }; +} + +test("legacy registries preserve the existing lifecycle without parsing refs", () => { + assert.deepEqual( + decideCollaborationLifecycle({ + profile_id: null, + operation: "result_publish", + caller_goal_ref: "legacy", + }), + { kind: "legacy" }, + ); + assert.deepEqual( + decideCollaborationLifecycle({ + profile_id: "future_profile", + operation: "result_publish", + }), + { kind: "reject", code: "unsupported_profile" }, + ); +}); + +test("request creation requires the captured instance to remain current", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("request_create", { + record_goal_ref: null, + route_goal_ref: null, + }), + ), + { kind: "allow", mode: "current_instance", goal_ref: GOAL_A }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("request_create", { + current_goal_ref: GOAL_B, + record_goal_ref: null, + route_goal_ref: null, + }), + ), + { kind: "reject", code: "stale_goal_instance" }, + ); +}); + +test("current inbox observations omit unstamped and other-instance records", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("inbox_observe", { record_goal_ref: null }), + ), + { kind: "omit", code: "legacy_unbound" }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("inbox_observe", { record_goal_ref: GOAL_B }), + ), + { kind: "omit", code: "different_goal_instance" }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("inbox_observe", { current_goal_ref: GOAL_B }), + ), + { kind: "omit", code: "different_goal_instance" }, + ); +}); + +test("current mutations reject after recreation", () => { + for (const operation of [ + "read_record", + "receiver_decide", + "artifact_link", + "peer_return_consume", + ]) { + assert.deepEqual( + decideCollaborationLifecycle( + facts(operation, { current_goal_ref: GOAL_B }), + ), + { kind: "reject", code: "historical_mutation_forbidden" }, + ); + } +}); + +test("late results stay bound to their saved exact route", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("result_publish", { current_goal_ref: GOAL_B }), + ), + { kind: "allow", mode: "historical_result", goal_ref: GOAL_A }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("result_publish", { + current_goal_ref: GOAL_B, + route_goal_ref: GOAL_B, + }), + ), + { kind: "reject", code: "route_instance_mismatch" }, + ); +}); + +test("original return admission requires trusted initial delivery evidence", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("original_return_admit", { + current_goal_ref: GOAL_B, + initial_delivery_proved: false, + }), + ), + { kind: "reject", code: "initial_delivery_unproved" }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("original_return_admit", { current_goal_ref: GOAL_B }), + ), + { kind: "allow", mode: "historical_result", goal_ref: GOAL_A }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("original_return_settle", { current_goal_ref: GOAL_B }), + ), + { kind: "allow", mode: "historical_result", goal_ref: GOAL_A }, + ); +}); + +test("historical inspection remains read-only", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("history_inspect", { + current_goal_ref: GOAL_B, + route_goal_ref: null, + }), + ), + { kind: "allow", mode: "historical_read", goal_ref: GOAL_A }, + ); + assert.deepEqual( + decideCollaborationLifecycle( + facts("history_inspect", { + caller_goal_ref: GOAL_B, + current_goal_ref: GOAL_B, + }), + ), + { kind: "reject", code: "record_instance_mismatch" }, + ); +}); + +test("malformed strict facts fail closed", () => { + assert.deepEqual( + decideCollaborationLifecycle( + facts("receiver_decide", { caller_goal_ref: null }), + ), + { kind: "reject", code: "missing_exact_goal_ref" }, + ); + assert.throws(() => + decideCollaborationLifecycle({ + profile_id: "source_session_v1", + operation: "unknown", + }) + ); +}); diff --git a/tests/test_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py new file mode 100644 index 0000000000..92a24f9968 --- /dev/null +++ b/tests/test_collaboration_goal_instance.py @@ -0,0 +1,685 @@ +import json +import subprocess +import sys +import threading +from concurrent.futures import ThreadPoolExecutor +from pathlib import Path + +import pytest + +from loopx.capabilities.manager_context import ( + ENTRY_SCHEMA, + INSTRUCTION, + _hash, + _root, + _write, + acknowledge, + deliver, + POLICY_SCHEMA, + register_ingress, +) +from loopx.capabilities.manager_context.roundtrip import drain, report +from loopx.capabilities.manager_context.tracking import query +from loopx.chat_store import ChatSessionStore +from loopx.collaboration_mcp import register_collaboration_tools +from loopx.control_plane.collaboration.peers import read_inbox, request +from loopx.control_plane.goals.source_session_registry_state import guard_path +from loopx.control_plane.projects.registry_codec import ( + source_session_registry_transaction, +) +from loopx.file_lock import exclusive_cross_runtime_file_lock + + +INSTANCE_A = "ginst_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" +INSTANCE_B = "ginst_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb" + + +def _source_payload(root: Path, instance_id: str) -> dict: + return { + "profile_id": "source_session_v1", + "goals": [ + { + "id": "delivery", + "repo": str(root), + "status": "active", + "goal_instance_id": instance_id, + "coordination": { + "registered_agents": ["builder", "reviewer"], + }, + } + ], + "session_bindings": [], + "session_receipts": [], + "lifetime_receipts": [], + } + + +def _create_source_registry(root: Path) -> Path: + registry = root / ".loopx" / "registry.json" + with source_session_registry_transaction( + registry, + operation="test_create_source_registry", + create=lambda: _source_payload(root, INSTANCE_A), + ) as transaction: + transaction.commit(transaction.payload_copy()) + return registry + + +def _recreate(registry: Path) -> None: + with exclusive_cross_runtime_file_lock( + guard_path(registry, "delivery"), + operation="test_recreate_goal", + ): + with source_session_registry_transaction( + registry, + operation="test_recreate_goal_registry", + ) as transaction: + payload = transaction.payload_copy() + payload["goals"][0]["goal_instance_id"] = INSTANCE_B + transaction.commit(payload) + + +def _brief(purpose: str) -> dict: + return { + "schema_version": "collaboration_brief_v0", + "purpose": purpose, + "context": "Use the current Goal instance only.", + "constraints": ["Do not infer continuity from the Goal alias."], + "inputs": [], + "acceptance": ["The result remains bound to its originating instance."], + "return_requirement": "Return one evidence-backed conclusion.", + } + + +def _manager_request(root: Path, registry: Path): + store = ChatSessionStore(root) + session = store.create_session( + goal_id="loopx-manager", + agent_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="fixture", + channel_id="manager", + ) + turn, _ = store.create_turn( + session["session_id"], + client_turn_id="owner-request", + message="Check the current delivery plan.", + origin="web", + ) + receipt = deliver( + root, + registry, + session=session, + turn=turn, + request={ + "goal_id": "delivery", + "agent_id": "builder", + "brief": _brief("Check the delivery plan"), + }, + ) + store.update_turn( + session["session_id"], + turn["turn_id"], + status="completing", + response={ + "message": "Delegated", + "context_handoff_receipt": receipt, + }, + ) + store.finalize_managed_turn_completion( + session["session_id"], + turn["turn_id"], + ) + return store, session, receipt + + +def _external_manager_request(root: Path, registry: Path): + store = ChatSessionStore(root) + session = store.create_session( + goal_id="loopx-manager", + agent_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="fixture", + channel_id="manager.external.fixture", + ) + turn, _ = store.create_turn( + session["session_id"], + client_turn_id="owner-request", + message="Check the current delivery plan.", + origin="lark", + ) + target = {"goal_id": "delivery", "agent_id": "builder"} + _write( + _root(root) / "policy.json", + { + "schema_version": POLICY_SCHEMA, + "sources": { + session["channel_id"]: { + "sender_ids": ["owner"], + "targets": [target], + } + }, + }, + ) + register_ingress( + root, + session_id=session["session_id"], + client_turn_id=turn["client_turn_id"], + channel=session["channel_id"], + sender_id="owner", + message=turn["message"], + source_id="lark:source", + ) + receipt = deliver( + root, + registry, + session=session, + turn=turn, + request=target, + ) + store.update_turn( + session["session_id"], + turn["turn_id"], + status="completing", + response={ + "message": "Delegated", + "context_handoff_receipt": receipt, + }, + ) + store.finalize_managed_turn_completion( + session["session_id"], + turn["turn_id"], + ) + acknowledge( + root, + "delivery", + "builder", + receipt["request_id"], + "adopt", + "Work in A.", + registry=registry, + caller_goal_ref=receipt["goal_ref"], + ) + report( + root, + "delivery", + "builder", + receipt["request_id"], + "conclusion", + "A completed its bounded review.", + registry=registry, + caller_goal_ref=receipt["goal_ref"], + ) + return store, session, receipt + + +def test_recreated_goal_cannot_observe_or_mutate_prior_instance_requests( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + _, _, receipt = _manager_request(tmp_path, registry) + goal_ref_a = receipt["goal_ref"] + + first = read_inbox( + tmp_path, + registry, + "delivery", + "builder", + caller_goal_ref=goal_ref_a, + ) + assert [item["request_id"] for item in first["items"]] == [ + receipt["request_id"] + ] + acknowledge( + tmp_path, + "delivery", + "builder", + receipt["request_id"], + "adopt", + "Work in A.", + registry=registry, + caller_goal_ref=goal_ref_a, + ) + + _recreate(registry) + + assert read_inbox( + tmp_path, + registry, + "delivery", + "builder", + )["items"] == [] + with pytest.raises((FileNotFoundError, ValueError)): + acknowledge( + tmp_path, + "delivery", + "builder", + receipt["request_id"], + "reject", + "B must not mutate A.", + registry=registry, + ) + + +def test_late_prior_instance_result_returns_only_to_its_saved_conversation( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + store, session, receipt = _manager_request(tmp_path, registry) + goal_ref_a = receipt["goal_ref"] + acknowledge( + tmp_path, + "delivery", + "builder", + receipt["request_id"], + "adopt", + "Work in A.", + registry=registry, + caller_goal_ref=goal_ref_a, + ) + _recreate(registry) + + report( + tmp_path, + "delivery", + "builder", + receipt["request_id"], + "conclusion", + "A completed its bounded review.", + registry=registry, + caller_goal_ref=goal_ref_a, + ) + assert drain(tmp_path, registry, store, None) == 1 + returned = [ + item + for item in store.messages(session["session_id"]) + if item.get("origin") == "manager_followup" + ] + assert len(returned) == 1 + assert "A completed its bounded review." in returned[0]["text"] + state = json.loads( + ( + _root(tmp_path) + / "replies" + / receipt["request_id"] + / "conclusion.delivery.json" + ).read_text(encoding="utf-8") + ) + assert state["status"] == "delivered" + assert state["goal_ref"] == goal_ref_a + + rows = query( + tmp_path, + registry, + goal_ids=["delivery"], + owner_scope=True, + request_id=receipt["request_id"], + )["rows"] + assert rows[0]["goal_ref"] == goal_ref_a + with pytest.raises((FileNotFoundError, ValueError)): + report( + tmp_path, + "delivery", + "builder", + receipt["request_id"], + "conclusion", + "B cannot replace A's result.", + registry=registry, + ) + + +def test_same_peer_operation_id_is_distinct_after_goal_recreation( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + first = request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + "review-one", + _brief("Review A"), + ) + _recreate(registry) + second = request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + "review-one", + _brief("Review B"), + ) + + assert first["request_id"] != second["request_id"] + current = read_inbox(tmp_path, registry, "delivery", "reviewer") + assert [item["request_id"] for item in current["items"]] == [ + second["request_id"] + ] + + +def test_exact_return_releases_lifetime_lock_around_provider_io( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + store, _, receipt = _external_manager_request(tmp_path, registry) + + entered = threading.Event() + release = threading.Event() + + class Transport: + calls = 0 + + def __call__(self, *_args): + self.calls += 1 + entered.set() + assert release.wait(timeout=5) + return { + "reply_verified": True, + "idempotency_key": "sha256:provider-proof", + } + + transport = Transport() + with ThreadPoolExecutor(max_workers=2) as executor: + first = executor.submit( + drain, + tmp_path, + registry, + store, + transport, + ) + assert entered.wait(timeout=5) + _recreate(registry) + second = executor.submit( + drain, + tmp_path, + registry, + store, + transport, + ) + assert second.result(timeout=5) == 0 + release.set() + assert first.result(timeout=5) == 1 + + assert transport.calls == 1 + state = json.loads( + ( + _root(tmp_path) + / "replies" + / receipt["request_id"] + / "conclusion.delivery.json" + ).read_text(encoding="utf-8") + ) + assert state["status"] == "delivered" + assert state["goal_ref"] == receipt["goal_ref"] + assert "admission" not in state + + +def test_exact_external_return_verifies_after_recreation_without_resend( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + store, _, receipt = _external_manager_request(tmp_path, registry) + + 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_exact_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": True, + "verification_performed": True, + "reply_verified": True, + } + + transport = Transport() + assert drain(tmp_path, registry, store, transport) == 1 + first = json.loads( + ( + _root(tmp_path) + / "replies" + / receipt["request_id"] + / "conclusion.delivery.json" + ).read_text(encoding="utf-8") + ) + assert first["status"] == "verification_required" + assert first["goal_ref"] == receipt["goal_ref"] + + _recreate(registry) + assert drain(tmp_path, registry, store, transport) == 1 + assert transport.send_calls == 1 + assert transport.verify_calls == 1 + state = json.loads( + ( + _root(tmp_path) + / "replies" + / receipt["request_id"] + / "conclusion.delivery.json" + ).read_text(encoding="utf-8") + ) + assert state["status"] == "delivered" + assert state["goal_ref"] == receipt["goal_ref"] + assert state["verification"] == "reconciled_after_restart" + + +def test_long_lived_mcp_keeps_its_captured_instance_after_recreation( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + _, _, receipt = _manager_request(tmp_path, registry) + + class Server: + def __init__(self) -> None: + self.tools = {} + + def tool(self): + def register(function): + self.tools[function.__name__] = function + return function + + return register + + server = Server() + register_collaboration_tools( + server, + tmp_path, + registry, + "delivery", + "builder", + tmp_path, + ) + server.tools["assess_request"]( + receipt["request_id"], + "adopt", + "Work in A.", + ) + _recreate(registry) + + assert server.tools["read_context"]()["items"] == [] + with pytest.raises(ValueError, match="stale_goal_instance"): + server.tools["request_peer"]( + "reviewer", + "stale-review", + _brief("Do not create work in B"), + ) + result = server.tools["return_result"]( + receipt["request_id"], + "A completed its bounded review.", + ) + assert result["status"] == "queued_for_original_conversation" + + +def test_manager_inbox_cli_reads_and_decides_current_exact_request( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + receipt = request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + "cli-review", + _brief("Review through the CLI"), + ) + base = [ + sys.executable, + "-m", + "loopx.cli", + "--runtime-root", + str(tmp_path), + "--registry", + str(registry), + "manager-inbox", + ] + + read = subprocess.run( + [ + *base, + "read", + "--goal-id", + "delivery", + "--agent-id", + "reviewer", + ], + check=False, + capture_output=True, + text=True, + timeout=30, + ) + assert read.returncode == 0, (read.stdout, read.stderr) + assert json.loads(read.stdout)["items"][0]["goal_ref"] == { + "goal_id": "delivery", + "goal_instance_id": INSTANCE_A, + } + + decision = subprocess.run( + [ + *base, + "acknowledge", + "--goal-id", + "delivery", + "--agent-id", + "reviewer", + "--request-id", + receipt["request_id"], + "--decision", + "adopt", + "--reason", + "Review the exact request.", + ], + check=False, + capture_output=True, + text=True, + timeout=30, + ) + assert decision.returncode == 0, (decision.stdout, decision.stderr) + assert json.loads(decision.stdout)["goal_ref"]["goal_instance_id"] == INSTANCE_A + + +def test_legacy_delivery_keeps_v1_paths_ids_and_json_bytes( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + import loopx.capabilities.manager_context.roundtrip as roundtrip + import loopx.capabilities.manager_context.tracking as tracking + + fixed_time = "2026-09-26T00:00:00+00:00" + monkeypatch.setattr(tracking, "_now", lambda: fixed_time) + monkeypatch.setattr(roundtrip, "_now", lambda: fixed_time) + registry = tmp_path / "registry.json" + registry.write_text( + json.dumps( + { + "goals": [ + { + "id": "delivery", + "repo": str(tmp_path), + "coordination": {"registered_agents": ["builder"]}, + } + ] + } + ), + encoding="utf-8", + ) + session = {"session_id": "manager-session", "channel_id": "manager"} + turn = { + "client_turn_id": "request-one", + "origin": "web", + "message": "Check the plan.", + } + target = {"goal_id": "delivery", "agent_id": "builder"} + source_id = "web:" + _hash( + [session["session_id"], turn["client_turn_id"]] + ) + request_id = _hash([source_id, target]) + + receipt = deliver( + tmp_path, + registry, + session=session, + turn=turn, + request=target, + ) + + assert receipt["request_id"] == request_id + assert "goal_ref" not in receipt + entry = { + "schema_version": ENTRY_SCHEMA, + "request_id": request_id, + **target, + "source_id": source_id, + "message": turn["message"], + "instruction": INSTRUCTION, + "delivered_at": fixed_time, + "source_channel": "manager", + } + route = { + "request_id": request_id, + **target, + "source_id": source_id, + "session_id": session["session_id"], + "client_turn_id": turn["client_turn_id"], + "channel_id": session["channel_id"], + "registered_at": fixed_time, + } + entry_path = ( + _root(tmp_path) + / "entries" + / _hash(target) + / f"{request_id}.json" + ) + route_path = _root(tmp_path) / "roundtrips" / f"{request_id}.json" + assert entry_path.read_bytes() == json.dumps( + entry, + ensure_ascii=False, + ).encode() + assert route_path.read_bytes() == json.dumps( + route, + ensure_ascii=False, + ).encode() From d118c7752353e5d45d379390ad918c40d955d52c Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sat, 26 Sep 2026 22:01:32 +0800 Subject: [PATCH 2/4] fix(collaboration): defer GoalRef capture without a registry Signed-off-by: duanjialing.777 --- loopx/collaboration_mcp.py | 35 +++++++++++++++++++++++++--------- tests/test_local_delegation.py | 23 ++++++++++++++++++++++ 2 files changed, 49 insertions(+), 9 deletions(-) diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index aee2bd5cca..63e59295be 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -20,6 +20,7 @@ import time from functools import lru_cache from pathlib import Path +from threading import Lock from typing import TYPE_CHECKING, Literal if TYPE_CHECKING: @@ -309,11 +310,27 @@ class Delegations: def __init__(self, root: Path, registry: Path, goal_id: str, agent_id: str, config: Path): self.root, self.registry = root.resolve(), registry.resolve() self.goal_id, self.agent_id, self.config = goal_id, agent_id, config.resolve() - self.goal_ref = capture_collaboration_goal_ref( - self.registry, - goal_id=self.goal_id, - agent_id=self.agent_id, - ) + self._goal_ref_lock = Lock() + try: + self.goal_ref = capture_collaboration_goal_ref( + self.registry, + goal_id=self.goal_id, + agent_id=self.agent_id, + ) + except FileNotFoundError: + self.goal_ref = None + + def _caller_goal_ref(self) -> dict[str, str] | None: + if self.goal_ref is not None: + return self.goal_ref + with self._goal_ref_lock: + if self.goal_ref is None: + self.goal_ref = capture_collaboration_goal_ref( + self.registry, + goal_id=self.goal_id, + agent_id=self.agent_id, + ) + return self.goal_ref def binding(self, binding_id: str, *, require_active: bool = False) -> dict: _goal(self.registry, self.goal_id, self.agent_id, require_active=require_active) @@ -435,7 +452,7 @@ def start(self, binding_id: str, operation_id: str, brief: dict, delegation_results.require_dependencies(self, binding, brief) delivered = request(self.root, self.registry, self.goal_id, self.agent_id, binding["agent_id"], operation_id, brief, parent_request_id, - caller_goal_ref=self.goal_ref) + caller_goal_ref=self._caller_goal_ref()) identity = {"binding": binding, "request_id": delivered["request_id"], "operation_id": operation_id} if exists: if _read(path).get("identity") != identity: @@ -767,7 +784,7 @@ def _receiver_adopted(self, row: dict, binding: dict) -> bool: self.registry, goal_id=self.goal_id, agents=(), - caller_goal_ref=self.goal_ref, + caller_goal_ref=self._caller_goal_ref(), ) as goal_scope: entry = _entry( self.root, @@ -794,7 +811,7 @@ def _delegation_bootstrap(self, row: dict, binding: dict) -> dict: self.registry, goal_id=self.goal_id, agents=(), - caller_goal_ref=self.goal_ref, + caller_goal_ref=self._caller_goal_ref(), ) as goal_scope: entry = _entry( self.root, @@ -1054,7 +1071,7 @@ def _execute(self, path: Path, row: dict, binding: dict) -> None: } ), registry=self.registry, - caller_goal_ref=self.goal_ref, + caller_goal_ref=self._caller_goal_ref(), ) self._observe(path, row, "accepted", canonical_done=True, acceptance_ready=True, artifacts_current=True) except (ValueError, KeyError, subprocess.TimeoutExpired, EffectRuntimeRemoteError) as exc: diff --git a/tests/test_local_delegation.py b/tests/test_local_delegation.py index 3326324140..d29b13fd7e 100644 --- a/tests/test_local_delegation.py +++ b/tests/test_local_delegation.py @@ -83,6 +83,29 @@ def test_worker_rejects_unbounded_operation_arguments(tmp_path, monkeypatch, ope assert calls == [] +def test_delegation_captures_goal_ref_after_registry_becomes_available(tmp_path, monkeypatch): + from loopx import collaboration_mcp as delegation + + runner = Delegations( + tmp_path, + tmp_path / "registry.json", + "goal", + "lead", + tmp_path / "config.json", + ) + captured = {"goal_id": "goal", "goal_instance_id": "instance-a"} + calls = [] + monkeypatch.setattr( + delegation, + "capture_collaboration_goal_ref", + lambda *args, **kwargs: calls.append((args, kwargs)) or captured, + ) + + assert runner._caller_goal_ref() == captured + assert runner._caller_goal_ref() == captured + assert len(calls) == 1 + + def test_worker_waits_for_a_transient_status_probe(service, monkeypatch): """A reader temporarily holding the lock must not discard admitted work.""" from loopx import collaboration_mcp as delegation From 4e962db9346965564200c0506a56133db9e2814d Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 28 Sep 2026 16:55:40 +0800 Subject: [PATCH 3/4] fix(collaboration): fence exact return delivery Signed-off-by: duanjialing.777 --- .../capabilities/manager_context/roundtrip.py | 39 ++++- loopx/control_plane/collaboration/inbox.py | 19 ++- tests/test_collaboration_goal_instance.py | 138 ++++++++++++++++-- tests/test_inbox_pagination.py | 24 +++ 4 files changed, 203 insertions(+), 17 deletions(-) diff --git a/loopx/capabilities/manager_context/roundtrip.py b/loopx/capabilities/manager_context/roundtrip.py index 5fb6218656..89f95a1a8f 100644 --- a/loopx/capabilities/manager_context/roundtrip.py +++ b/loopx/capabilities/manager_context/roundtrip.py @@ -9,12 +9,17 @@ import logging import re import threading +from contextlib import ExitStack from datetime import datetime, timezone, timedelta from uuid import uuid4 from . import _root, _read, _write, _hash, authority from .tracking import _entry, _now -from ...file_lock import exclusive_file_lock +from ...file_lock import ( + LockAcquisitionPolicy, + LockAcquireTimeoutError, + exclusive_file_lock, +) from ...control_plane.collaboration import conversation_scope from ...control_plane.collaboration.goal_instance_scope import ( collaboration_goal_scope, @@ -457,6 +462,20 @@ def _exact_return_context(root, registry, store, path, state_path, now): and now.isoformat() < admission["expires_at"] ): return None + if ( + state.get("status") == "admitted" + and state.get("attempt") is None + and not conversation_scope(session)["private_conversation"] + ): + _write( + state_path, + { + "status": "explicit_unverified", + "error": "provider_delivery_unverified", + "goal_ref": row["goal_ref"], + }, + ) + return None if ( state.get("retry_at") and state.get("status") != "admitted" @@ -588,6 +607,20 @@ def _drain_exact(root, registry, store, external_sender, *, now, cancelled): if path.stem not in PHASES: continue state_path = path.with_name(path.stem + ".delivery.json") + effect_locks = ExitStack() + try: + effect_locks.enter_context( + exclusive_file_lock( + _root(root) + / "return-effect-locks" + / path.parent.name + / path.stem, + policy=LockAcquisitionPolicy.SINGLE_FLIGHT, + operation="manager_return_delivery", + ) + ) + except LockAcquireTimeoutError: + continue try: context = _exact_return_context( root, @@ -598,11 +631,13 @@ def _drain_exact(root, registry, store, external_sender, *, now, cancelled): now, ) except (OSError, ValueError, KeyError, TypeError, RuntimeError): + effect_locks.close() logging.getLogger(__name__).warning( "Exact manager return admission unavailable" ) continue if context is None: + effect_locks.close() continue row = context["row"] route = context["route"] @@ -833,6 +868,8 @@ def record_attempt(value): logging.getLogger(__name__).warning( "Exact manager return settlement unavailable" ) + finally: + effect_locks.close() processed += 1 if processed >= 20: break diff --git a/loopx/control_plane/collaboration/inbox.py b/loopx/control_plane/collaboration/inbox.py index 81dfed15ce..797c9658c8 100644 --- a/loopx/control_plane/collaboration/inbox.py +++ b/loopx/control_plane/collaboration/inbox.py @@ -125,11 +125,20 @@ def pending( scope: CollaborationGoalScope | None = None, ) -> dict: cursor_scope = _hash( - [ - "pending_requests_v1", - str(runtime_root.resolve()), - _target(goal_id, agent_id, scope), - ] + ( + [ + "pending_requests_v1", + str(runtime_root.resolve()), + _target(goal_id, agent_id, scope), + ] + if scope is not None and scope.exact + else [ + "pending_requests_v1", + str(runtime_root.resolve()), + goal_id, + agent_id, + ] + ) ) after = "" if cursor is not None: diff --git a/tests/test_collaboration_goal_instance.py b/tests/test_collaboration_goal_instance.py index 92a24f9968..1b50b58512 100644 --- a/tests/test_collaboration_goal_instance.py +++ b/tests/test_collaboration_goal_instance.py @@ -3,6 +3,7 @@ import sys import threading from concurrent.futures import ThreadPoolExecutor +from datetime import datetime, timedelta, timezone from pathlib import Path import pytest @@ -18,7 +19,11 @@ POLICY_SCHEMA, register_ingress, ) -from loopx.capabilities.manager_context.roundtrip import drain, report +from loopx.capabilities.manager_context.roundtrip import ( + EXACT_DELIVERY_ADMISSION_SECONDS, + drain, + report, +) from loopx.capabilities.manager_context.tracking import query from loopx.chat_store import ChatSessionStore from loopx.collaboration_mcp import register_collaboration_tools @@ -359,11 +364,39 @@ def test_same_peer_operation_id_is_distinct_after_goal_recreation( ] -def test_exact_return_releases_lifetime_lock_around_provider_io( +def test_exact_inbox_cursor_cannot_cross_goal_recreation(tmp_path: Path) -> None: + registry = _create_source_registry(tmp_path) + for index in range(21): + request( + tmp_path, + registry, + "delivery", + "builder", + "reviewer", + f"review-a-{index}", + _brief(f"Review A item {index}"), + ) + first = read_inbox(tmp_path, registry, "delivery", "reviewer") + assert first["has_more"] + + _recreate(registry) + + with pytest.raises(ValueError, match="cursor scope mismatch"): + read_inbox( + tmp_path, + registry, + "delivery", + "reviewer", + cursor=first["next_cursor"], + ) + + +def test_exact_return_single_flight_outlives_admission_without_holding_lifetime( tmp_path: Path, ) -> None: registry = _create_source_registry(tmp_path) store, _, receipt = _external_manager_request(tmp_path, registry) + admitted_at = datetime(2026, 9, 28, tzinfo=timezone.utc) entered = threading.Event() release = threading.Event() @@ -373,8 +406,9 @@ class Transport: def __call__(self, *_args): self.calls += 1 - entered.set() - assert release.wait(timeout=5) + if self.calls == 1: + entered.set() + assert release.wait(timeout=5) return { "reply_verified": True, "idempotency_key": "sha256:provider-proof", @@ -388,6 +422,7 @@ def __call__(self, *_args): registry, store, transport, + now=admitted_at, ) assert entered.wait(timeout=5) _recreate(registry) @@ -397,6 +432,8 @@ def __call__(self, *_args): registry, store, transport, + now=admitted_at + + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS + 1), ) assert second.result(timeout=5) == 0 release.set() @@ -416,11 +453,72 @@ def __call__(self, *_args): assert "admission" not in state +def test_exact_external_return_does_not_resend_after_expired_unknown_admission( + tmp_path: Path, +) -> None: + registry = _create_source_registry(tmp_path) + store, _, receipt = _external_manager_request(tmp_path, registry) + admitted_at = datetime(2026, 9, 28, tzinfo=timezone.utc) + state_path = ( + _root(tmp_path) + / "replies" + / receipt["request_id"] + / "conclusion.delivery.json" + ) + _write( + state_path, + { + "status": "admitted", + "goal_ref": receipt["goal_ref"], + "admission": { + "token": "interrupted-writer", + "prior_status": "queued", + "admitted_at": admitted_at.isoformat(), + "expires_at": ( + admitted_at + + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS) + ).isoformat(), + }, + }, + ) + + class Transport: + calls = 0 + + def __call__(self, *_args): + self.calls += 1 + return { + "reply_verified": True, + "idempotency_key": "sha256:unexpected-resend", + } + + transport = Transport() + assert ( + drain( + tmp_path, + registry, + store, + transport, + now=admitted_at + + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS + 1), + ) + == 0 + ) + assert transport.calls == 0 + state = json.loads(state_path.read_text(encoding="utf-8")) + assert state == { + "status": "explicit_unverified", + "error": "provider_delivery_unverified", + "goal_ref": receipt["goal_ref"], + } + + def test_exact_external_return_verifies_after_recreation_without_resend( tmp_path: Path, ) -> None: registry = _create_source_registry(tmp_path) store, _, receipt = _external_manager_request(tmp_path, registry) + admitted_at = datetime(2026, 9, 28, tzinfo=timezone.utc) class Transport: send_calls = 0 @@ -444,10 +542,7 @@ def send_with_attempt( "provider_receipt": "sha256:" + "b" * 64, } ) - return { - "external_write_performed": True, - "reply_verified": False, - } + raise SystemExit("simulated sender crash after attempt persistence") def verify(self, *_args): self.verify_calls += 1 @@ -458,7 +553,17 @@ def verify(self, *_args): } transport = Transport() - assert drain(tmp_path, registry, store, transport) == 1 + with pytest.raises( + SystemExit, + match="simulated sender crash after attempt persistence", + ): + drain( + tmp_path, + registry, + store, + transport, + now=admitted_at, + ) first = json.loads( ( _root(tmp_path) @@ -467,11 +572,22 @@ def verify(self, *_args): / "conclusion.delivery.json" ).read_text(encoding="utf-8") ) - assert first["status"] == "verification_required" + assert first["status"] == "admitted" + assert first["attempt"]["message_ref"] == "om_exact_reply" assert first["goal_ref"] == receipt["goal_ref"] _recreate(registry) - assert drain(tmp_path, registry, store, transport) == 1 + assert ( + drain( + tmp_path, + registry, + store, + transport, + now=admitted_at + + timedelta(seconds=EXACT_DELIVERY_ADMISSION_SECONDS + 1), + ) + == 1 + ) assert transport.send_calls == 1 assert transport.verify_calls == 1 state = json.loads( diff --git a/tests/test_inbox_pagination.py b/tests/test_inbox_pagination.py index efcb664e8b..b57ef86294 100644 --- a/tests/test_inbox_pagination.py +++ b/tests/test_inbox_pagination.py @@ -1,6 +1,7 @@ """Receiver recovery through real CLI processes and identity-bound MCP stdio.""" import asyncio +import hashlib import json import subprocess import sys @@ -114,6 +115,29 @@ def test_cli_page_boundaries(inbox, count): assert read_ids(root) == set(expected[:20]) +def test_cli_resumes_cursor_issued_before_exact_goal_scoping(inbox): + root, registry, seed, _ = inbox + expected = seed(45) + legacy_scope = hashlib.sha256( + json.dumps( + [ + "pending_requests_v1", + str(root.resolve()), + "delivery", + "receiver", + ], + ensure_ascii=False, + sort_keys=True, + ).encode() + ).hexdigest() + cursor = f"1:{legacy_scope}:{expected[19]}" + + resumed = cli(root, registry, "read", "--cursor", cursor) + + assert ids(resumed) == expected[20:40] + assert resumed["has_more"] + + def test_cli_rejects_wrong_scope_and_malformed_cursors_before_receipts(inbox): root, registry, seed, _ = inbox expected = seed(21) From 8b7af90a1381408da516fc164fbe898cd5f18c1c Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Mon, 28 Sep 2026 17:34:14 +0800 Subject: [PATCH 4/4] test(architecture): register collaboration registry reads Signed-off-by: duanjialing.777 --- .../project_registry_io_manifest_v1.json | 56 +++++++++++++------ .../test_source_session_registry_denial.py | 12 ++++ 2 files changed, 52 insertions(+), 16 deletions(-) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 09c271ee70..bfe5ebe6f6 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -174,27 +174,27 @@ "classification": "codec_api" }, { - "site": "loopx/capabilities/manager_context/__init__.py::.authority::codec_read:load_registry#1", - "line": 74, + "site": "loopx/capabilities/manager_context/__init__.py::.authority::codec_read:load_project_registry#1", + "line": 86, "column": 20, "kind": "codec_read", - "api": "load_registry", + "api": "load_project_registry", "classification": "codec_api" }, { - "site": "loopx/capabilities/manager_context/__init__.py::.configure_delivery_target::codec_read:load_registry#1", - "line": 321, + "site": "loopx/capabilities/manager_context/__init__.py::.configure_delivery_target::codec_read:load_project_registry#1", + "line": 406, "column": 20, "kind": "codec_read", - "api": "load_registry", + "api": "load_project_registry", "classification": "codec_api" }, { - "site": "loopx/capabilities/manager_context/__init__.py::.configure_evidence_scope::codec_read:load_registry#1", - "line": 279, + "site": "loopx/capabilities/manager_context/__init__.py::.configure_evidence_scope::codec_read:load_project_registry#1", + "line": 364, "column": 16, "kind": "codec_read", - "api": "load_registry", + "api": "load_project_registry", "classification": "codec_api" }, { @@ -213,6 +213,14 @@ "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", + "line": 883, + "column": 22, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/capabilities/periodic_report/cadence_runtime.py::.periodic_report_cadence_hooks::codec_read:load_registry#1", "line": 27, @@ -437,6 +445,14 @@ "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/cli.py::.main::codec_read:load_project_registry#1", + "line": 837, + "column": 17, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/cli_commands/authority_archive.py::.authority_upgrade_roots::codec_read:load_registry#1", "line": 123, @@ -654,11 +670,11 @@ "classification": "codec_api" }, { - "site": "loopx/cli_commands/manager_inbox.py::.handle_manager_inbox::codec_read:load_registry#1", + "site": "loopx/cli_commands/manager_inbox.py::.handle_manager_inbox::codec_read:load_project_registry#1", "line": 105, "column": 20, "kind": "codec_read", - "api": "load_registry", + "api": "load_project_registry", "classification": "codec_api" }, { @@ -887,12 +903,20 @@ }, { "site": "loopx/contract.py::.check_contract::codec_read:load_registry#1", - "line": 1014, + "line": 1027, "column": 20, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/control_plane/collaboration/goal_instance_scope.py::.collaboration_goal_scope::codec_read:load_project_registry#1", + "line": 108, + "column": 20, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/control_plane/collaboration/peer_host_route.py::.resolve_peer_host_route::codec_read:load_registry#1", "line": 65, @@ -902,11 +926,11 @@ "classification": "codec_api" }, { - "site": "loopx/control_plane/collaboration/peers.py::._goal::codec_read:load_registry#1", - "line": 35, - "column": 21, + "site": "loopx/control_plane/collaboration/peers.py::._goal::codec_read:load_project_registry#1", + "line": 51, + "column": 22, "kind": "codec_read", - "api": "load_registry", + "api": "load_project_registry", "classification": "codec_api" }, { diff --git a/tests/architecture/test_source_session_registry_denial.py b/tests/architecture/test_source_session_registry_denial.py index d4887729f4..cf6404d9be 100644 --- a/tests/architecture/test_source_session_registry_denial.py +++ b/tests/architecture/test_source_session_registry_denial.py @@ -8,8 +8,14 @@ DIRECT_LOADER_ALLOWLIST = { "loopx/authority.py", "loopx/bootstrap.py", + "loopx/capabilities/manager_context/__init__.py", + "loopx/capabilities/manager_context/roundtrip.py", "loopx/claude_goal_mode/scripts/connect.py", + "loopx/cli.py", + "loopx/cli_commands/manager_inbox.py", "loopx/configure_goal.py", + "loopx/control_plane/collaboration/goal_instance_scope.py", + "loopx/control_plane/collaboration/peers.py", "loopx/control_plane/goals/first_party_host_admission.py", "loopx/control_plane/projects/registry.py", "loopx/kunluncode_goal_mode/cli.py", @@ -33,6 +39,12 @@ def test_direct_project_registry_loaders_have_source_session_denial() -> None: assert callers == DIRECT_LOADER_ALLOWLIST source_session_owners = { + "loopx/capabilities/manager_context/__init__.py", + "loopx/capabilities/manager_context/roundtrip.py", + "loopx/cli.py", + "loopx/cli_commands/manager_inbox.py", + "loopx/control_plane/collaboration/goal_instance_scope.py", + "loopx/control_plane/collaboration/peers.py", "loopx/control_plane/goals/first_party_host_admission.py", "loopx/control_plane/projects/registry.py", }