diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md index 9232d7eaf5..d0c0476008 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.md @@ -237,6 +237,15 @@ The steward sends one source-linked consultation or delegated-work request to th **Discovery implementation checkpoint.** The existing manager/context read tool now has an `agents` view over the complete permitted registry, with responsibility search, pagination and explicit stopped-history opt-in. It reads independently of the progress snapshot's Agent cap and sender-bound delivery list. Local owner scope is broad by default; Goal Chat and external audiences retain their scope. CLI and authorized SSH exports share the reader. Registration, declared responsibility, context delivery permission and unchecked execution readiness remain distinct. This qualifies a bounded read/diagnostic slice of A24, not worker selection, launch, adoption or original-route completion; prompt-only adapters and older remote installations remain explicit coverage gaps. Continue A24 through existing delivery configuration, actual worker execution and result return before claiming the golden query passed. +The shared local peer-route resolver now excludes alternatives only on explicit +host archive evidence before selecting a unique readable binding. This removes +manual task-link lookup for an Agent with archived session history; unknown, +missing, unsupported and multiple readable alternatives retain a gap. The rule +lives in TypeScript collaboration, while Python adapts registry and host reads. +Real disposable host-store and CLI tests qualify request/replay pinning, not +native worker submission, freshness, capacity, receiver adoption or A24. Keep +the App-first original-conversation pilot open until those facts are proven. + ### 5.6 One exchange, independent durable facts The user-facing exchange is **received → assessed/working → result**, with meaningful updates when needed. Internally, keep transport and work facts separate: diff --git a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md index aa7f91eb76..efb4d12913 100644 --- a/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md +++ b/docs/architecture/rfcs/capable-manager-semantic-handoff-v0.zh-CN.md @@ -212,6 +212,12 @@ LoopX 不是只有任务队列。交接应让接收方结合权威状态和持 **职责发现实现检查点。** 现有管家/项目对话读取工具新增 `agents` 视图,可搜索并分页读取权限范围内的完整注册目录,显式选择后可查已停止 Goal 的历史身份。它不受进度快照中 Agent 数量截断或发送方委派名单限制。主人本地管家默认广泛发现,Goal Chat 和外部群聊保留原有范围;CLI 与已授权 SSH 导出复用同一读取实现。注册、声明职责、上下文投递权限与尚未检查的执行就绪状态分别呈现。这只验证 A24 的发现和诊断切片,不代表已完成选人、启动、采用或原路返回;仅提示词适配器和旧版远端仍保留明确缺口。后续沿现有投递配置、真实 worker 执行和结果回传验收,不能据此宣称 golden query 已通过。 +共享本机 peer-route 解析器只有读到宿主明确归档证据才排除备选,再选择唯一可读绑定。 +这让保留已归档历史 session 的 Agent 无需主人手工找任务链接;未知、缺失、不支持或多个 +可读备选仍保留缺口。规则归属 TypeScript collaboration,Python 只适配注册与宿主读取。 +真实隔离宿主库和 CLI 验证请求与重试固定同一路径,不证明原生 worker 投递、新鲜度、容量、 +接收方采用或 A24。App 优先的原对话闭环继续保留,直到这些事实通过验收。 + ### 5.6 一次交互,分开的持久事实 用户体验是**收到 → 已评估/工作中 → 结果**,必要时有实质中间反馈。内部不能混淆工作与传输: diff --git a/docs/reference/protocols/peer-agent-directory-and-observation-v0.md b/docs/reference/protocols/peer-agent-directory-and-observation-v0.md index f0a50ec717..7ad3a211e1 100644 --- a/docs/reference/protocols/peer-agent-directory-and-observation-v0.md +++ b/docs/reference/protocols/peer-agent-directory-and-observation-v0.md @@ -160,8 +160,13 @@ Rules: For a named existing peer, `resolve-peer-route --goal-id ... --agent-id ...` reads the binding owner and the local host observer before any request is -recorded. `ambiguous` preserves all accepted candidate identities and asks for -an exact task link; it never selects the last or newest binding. An explicit +recorded. Without a task link, the shared typed selector may resolve exactly one +readable local task when every alternative has an explicit host `archived` +observation. It observes all matching candidates within the existing 32-thread +budget, independently of the three-row publication cap. Missing, unsupported, +failed or withheld observations remain unknown alternatives; multiple remaining +candidates or an over-budget inventory stay `ambiguous`. It never selects the +last, newest or only visible binding. An explicit `--thread-link` must resolve to that same Goal and Agent across the project registry. `unavailable` includes archived, unknown, unsupported and missing host observations. `not_authorized` means the named peer is outside the Goal or @@ -171,6 +176,16 @@ verify its own profile and submission permission. Every preview says `host_delivery: not_attempted`, and no route preview grants a claim, lease, session resume or message-send permission. +This changes the previous default refusal for *all* multiple bindings: archived +history no longer requires the owner to supply a link. It applies to the local +CLI and `manager-inbox request --require-host-route`, including trusted steward +and peer callers. Readable records are still neither fresh runtime liveness nor +model/capacity qualification. Restricted Chat and remote hosts do not acquire a +new observer or execution path. Registration, the scoped candidate set and the +selected identity are rechecked after host reads; a changed set remains +unresolved, while a mere reorder does not change the route. Request replay +keeps its original exact route. + ## Presence Vocabulary Presence answers "is this Agent runnable right now", not "is its work done". diff --git a/loopx/capabilities/manager_context/README.md b/loopx/capabilities/manager_context/README.md index cad60d244b..29661c7cf3 100644 --- a/loopx/capabilities/manager_context/README.md +++ b/loopx/capabilities/manager_context/README.md @@ -362,8 +362,10 @@ When the intended recipient is an existing Codex host task, resolve that peer before substituting a temporary child. First inspect `agent-directory` for the named Agent and its candidate count. Use `loopx resolve-peer-route --goal-id allocation --agent-id reviewer` for an -observed local route. If several historical bindings exist, supply the exact -user-selected task link with `--thread-link codex://threads/`; never choose +observed local route. Archived history permits a unique readable route without +a task link; missing or unknown alternatives preserve ambiguity. If the result +is still ambiguous, supply the exact user-selected task link with +`--thread-link codex://threads/`; never choose the last binding by order or recency. Then record the request with `manager-inbox request ... --require-host-route --peer-thread-link codex://threads/`. The result contains the same stable request id and a diff --git a/loopx/capabilities/manager_context/skills/loopx-manager/SKILL.md b/loopx/capabilities/manager_context/skills/loopx-manager/SKILL.md index 00986221f5..071cd9e480 100644 --- a/loopx/capabilities/manager_context/skills/loopx-manager/SKILL.md +++ b/loopx/capabilities/manager_context/skills/loopx-manager/SKILL.md @@ -118,8 +118,9 @@ the work it shows. When a user names an existing peer for review or other work, keep that Agent identity. The directory's `peer_route` is a candidate list, not a selected or verified host task. On a trusted host with CLI access, use `resolve-peer-route` -for the exact Goal and Agent; when several historical bindings exist, use the -user's exact task link or report the ambiguity. If a host message tool is +for the exact Goal and Agent. It may select one readable task when every other +binding is explicitly archived. If it still reports ambiguity, use the user's +exact task link or report the gap; never treat unknown as archived. If a host message tool is authorized, verify the selected task through that host, record the stable `manager-inbox request --require-host-route`, then send its content-free `host_delivery.message` to the selected task. Treat its `not_attempted` state diff --git a/loopx/control_plane/collaboration/peer_host_route.py b/loopx/control_plane/collaboration/peer_host_route.py index 7e05e651ea..a499dc40d2 100644 --- a/loopx/control_plane/collaboration/peer_host_route.py +++ b/loopx/control_plane/collaboration/peer_host_route.py @@ -14,12 +14,13 @@ from ...agent_registry import registered_agent_ids_for_goal from ...codex_app_thread_activity import codex_thread_observers from ...control_plane.agents.host_thread_activity import ( + HostThreadActivity, HostThreadObserver, HostThreadState, HostThreadUnknownReason, + MAX_OBSERVED_THREADS_PER_GOAL, ) from ...control_plane.runtime.public_safety import validate_public_safe_value -from ...history import load_registry from ...registry import find_registry_goal from ...thread_agent_binding import ( CODEX_THREAD_HOST_SURFACES, @@ -27,12 +28,69 @@ resolve_registry_thread_agent_binding, summarize_agent_binding_routes, ) +from ..effect_runtime import effect_runtime_result +from ..projects.registry_codec import load_registry PEER_HOST_ROUTE_SCHEMA_VERSION = "loopx_peer_host_route_v0" MAX_PUBLISHED_CANDIDATES = 3 +def _matching_bindings( + candidates: list[dict[str, str]], *, thread_id: str | None, host_surface: str | None, +) -> list[dict[str, str]]: + return [ + item for item in candidates + if (host_surface is None or item["host_surface"] == host_surface) + and (thread_id is None or ( + item["thread_id"] == thread_id and item["host_surface"] in CODEX_THREAD_HOST_SURFACES + )) + ] + + +def _observe_bindings( + matching: list[dict[str, str]], visible: list[dict[str, str]], + available_observers: Mapping[str, HostThreadObserver], +) -> tuple[list[dict[str, str]], list[HostThreadActivity | None]]: + # Read each host once, including every accepted alternative. A publication + # cap or an unreadable host must not turn an unknown binding into history. + requested: dict[str, set[str]] = {} + for candidate in matching: + if candidate in visible: + requested.setdefault(candidate["host_surface"], set()).add(candidate["thread_id"]) + observed: dict[str, Mapping[str, HostThreadActivity]] = {} + failures: dict[str, str] = {} + for surface, ids in requested.items(): + observer = available_observers.get(surface) + if observer is None: + failures[surface] = "host_observer_unavailable" + continue + try: + observed[surface] = observer(sorted(ids)) + except Exception: # noqa: BLE001 - host failure remains an unknown alternative. + failures[surface] = "host_observation_failed" + facts: list[dict[str, str]] = [] + activity: list[HostThreadActivity | None] = [] + for candidate in matching: + surface, thread = candidate["host_surface"], candidate["thread_id"] + item = observed.get(surface, {}).get(thread) + activity.append(item) + if candidate not in visible: + reason = "route_candidate_withheld" + elif surface in failures: + reason = failures[surface] + elif item is None: + reason = HostThreadUnknownReason.THREAD_NOT_FOUND.value + elif item.state is HostThreadState.UNKNOWN: + assert item.reason is not None # HostThreadActivity's constructor invariant. + reason = item.reason.value + else: + facts.append({"state": item.state.value}) + continue + facts.append({"state": "unavailable", "reason": reason}) + return facts, activity + + def resolve_peer_host_route( registry_path: Path, *, @@ -84,38 +142,50 @@ def resolve_peer_host_route( if withheld: result["withheld_candidate_count"] = withheld + thread_id = None if thread_link is not None: try: thread_id = codex_thread_deep_link_locator(thread_link)["thread_id"] except ValueError: result["reason"] = "invalid_thread_link" return result - matching = [ - item - for item in candidates - if item["thread_id"] == thread_id - and item["host_surface"] in CODEX_THREAD_HOST_SURFACES - and (host_surface is None or item["host_surface"] == host_surface) - ] - if not matching: - result.update(status="not_authorized", reason="thread_not_bound_to_peer") - return result - else: - matching = [ - item - for item in candidates - if host_surface is None or item["host_surface"] == host_surface - ] + matching = _matching_bindings(candidates, thread_id=thread_id, host_surface=host_surface) + if thread_link is not None and not matching: + result.update(status="not_authorized", reason="thread_not_bound_to_peer") + return result if not matching: return result - if len(matching) != 1: + if len(matching) > MAX_OBSERVED_THREADS_PER_GOAL: result.update(status="ambiguous", reason="multiple_binding_candidates") return result - selected = matching[0] - if selected not in visible: + if len(matching) == 1 and matching[0] not in visible: result["reason"] = "route_candidate_withheld" return result + available_observers = codex_thread_observers() if observers is None else observers + facts, activity = _observe_bindings(matching, visible, available_observers) + if len(activity) == 1 and activity[0] is not None: + result["host_observation"] = activity[0].to_payload() + selection = effect_runtime_result("collaboration.peer_host_route.select", {"observations": facts}) + if selection["status"] != "resolved": + result.update(status=selection["status"], reason=selection["reason"]) + return result + index = selection["selected_index"] + selected = matching[index] + # Host I/O must not hide a newly bound alternative or a revoked registration. + latest_goal = find_registry_goal(load_registry(registry_path), goal_id) + if latest_goal is None or agent_id not in registered_agent_ids_for_goal(latest_goal): + result.update(status="not_authorized", reason="peer_not_registered") + return result + latest = _matching_bindings( + summarize_agent_binding_routes([latest_goal], agent_id=agent_id)["candidates"], + thread_id=thread_id, host_surface=host_surface, + ) + if {(item["host_surface"], item["thread_id"]) for item in latest} != { + (item["host_surface"], item["thread_id"]) for item in matching + }: + result.update(status="ambiguous", reason="binding_identity_conflict") + return result exact = resolve_registry_thread_agent_binding( registry_path=registry_path, host_surface=selected["host_surface"], @@ -128,26 +198,8 @@ def resolve_peer_host_route( result.update(status="ambiguous", reason="binding_identity_conflict") return result - available_observers = codex_thread_observers() if observers is None else observers - observer = available_observers.get(selected["host_surface"]) - if observer is None: - result["reason"] = "host_observer_unavailable" - return result - try: - observation = observer([selected["thread_id"]]).get(selected["thread_id"]) - except Exception: # noqa: BLE001 - host read failure cannot authorize delivery. - result["reason"] = "host_observation_failed" - return result - if observation is None: - result["reason"] = HostThreadUnknownReason.THREAD_NOT_FOUND.value - return result - result["host_observation"] = observation.to_payload() - if observation.state is HostThreadState.ARCHIVED: - result["reason"] = "host_thread_archived" - elif observation.state is HostThreadState.UNKNOWN: - result["reason"] = ( - observation.reason.value if observation.reason else "host_unknown" - ) - else: - result.update(status="resolved", reason=None, selected_route=selected) + selected_activity = activity[index] + assert selected_activity is not None # The typed selector requires a readable observation. + result["host_observation"] = selected_activity.to_payload() + result.update(status="resolved", reason=None, selected_route=selected) return result diff --git a/loopx/control_plane/collaboration/peer_route_selection.ts b/loopx/control_plane/collaboration/peer_route_selection.ts new file mode 100644 index 0000000000..2f1b7309b3 --- /dev/null +++ b/loopx/control_plane/collaboration/peer_route_selection.ts @@ -0,0 +1,38 @@ +import type { JsonObject } from "../effect_program.ts"; +import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; +import { requireJsonObject, requireStringLiteral } from "../runtime_decode.ts"; + +type Observation = + | Readonly<{ state: "idle" | "turn_open" | "archived" }> + | Readonly<{ state: "unavailable"; reason: string }>; + +function observation(value: unknown): Observation { + const row = requireJsonObject(value, "peer host observation"); + const state = requireStringLiteral(row.state, + ["idle", "turn_open", "archived", "unavailable"], "peer host observation state"); + return state === "unavailable" ? { state, reason: requireStringLiteral(row.reason, + ["host_observer_unavailable", "host_observation_failed", "thread_not_found", + "store_unavailable", "record_unrecognized", "no_turn_marker", "unsupported_host", + "route_candidate_withheld"], "unavailable peer host reason") } : { state }; +} + +/** Choose a locator, not an executor or permission grant. Unknown is never retired. */ +export function selectObservedPeerHostRoute(params: JsonObject): JsonObject { + if (!Array.isArray(params.observations) || !params.observations.length + || params.observations.length > 32) { + throw new EffectRuntimeRequestError("peer host observations must contain 1..32 candidates"); + } + const rows = params.observations.map(observation); + const remaining = rows.flatMap((row, index) => row.state === "archived" ? [] : [index]); + if (remaining.length > 1) { + return { status: "ambiguous", reason: "multiple_binding_candidates", selected_index: null }; + } + if (!remaining.length) { + return { status: "unavailable", reason: "host_thread_archived", selected_index: null }; + } + const index = remaining[0]; + const row = rows[index]; + return row.state === "unavailable" + ? { status: "unavailable", reason: row.reason, selected_index: null } + : { status: "resolved", reason: null, selected_index: index }; +} diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 492dbf1291..ea98c7d2f7 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -235,6 +235,7 @@ import { } from "./collaboration/return_delivery.ts"; import { decideCollaborationLifecycle } from "./collaboration/goal_instance_lifecycle.ts"; import { inspectCollaborationInboxReceipts } from "./collaboration/inbox_receipts.ts"; +import { selectObservedPeerHostRoute } from "./collaboration/peer_route_selection.ts"; import { normalizeCollaborationRequest } from "./collaboration/semantic_request.ts"; import { @@ -761,6 +762,7 @@ export function createEffectRuntimeHandlers( (params) => normalizeCollaborationRequest(params.request), ], ["collaboration.inbox.inspect_receipts", inspectCollaborationInboxReceipts], + ["collaboration.peer_host_route.select", selectObservedPeerHostRoute], [ "collaboration.goal_instance.decide", (params) => decideCollaborationLifecycle(params), diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 767ddb1777..b8bc76963c 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -967,12 +967,20 @@ }, { "site": "loopx/control_plane/collaboration/peer_host_route.py::.resolve_peer_host_route::codec_read:load_registry#1", - "line": 65, + "line": 123, "column": 16, "kind": "codec_read", "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/control_plane/collaboration/peer_host_route.py::.resolve_peer_host_route::codec_read:load_registry#2", + "line": 176, + "column": 38, + "kind": "codec_read", + "api": "load_registry", + "classification": "codec_api" + }, { "site": "loopx/control_plane/collaboration/peers.py::._goal::codec_read:load_project_registry#1", "line": 52, diff --git a/tests/control_plane/test_peer_host_route.py b/tests/control_plane/test_peer_host_route.py index 0358141c95..a39fa359d4 100644 --- a/tests/control_plane/test_peer_host_route.py +++ b/tests/control_plane/test_peer_host_route.py @@ -5,8 +5,14 @@ import json import contextlib import io +import os +import sqlite3 +import subprocess +import sys from pathlib import Path +import pytest + from loopx.cli import main as cli_main from loopx.control_plane.agents.host_thread_activity import ( HostThreadActivity, @@ -38,8 +44,8 @@ def _registry(tmp_path: Path, bindings: list[dict[str, str]]) -> Path: return path -def _binding(thread_id: str, agent_id: str = "reviewer") -> dict[str, str]: - return {"agent_id": agent_id, "host_surface": "codex-app", "thread_id": thread_id} +def _binding(thread_id: str, agent_id: str = "reviewer", *, host_surface: str = "codex-app") -> dict[str, str]: + return {"agent_id": agent_id, "host_surface": host_surface, "thread_id": thread_id} def _observer(state: HostThreadState): @@ -48,7 +54,152 @@ def _observer(state: HostThreadState): } -def test_named_peer_with_several_bindings_needs_an_exact_selection( +@pytest.mark.parametrize("current", [HostThreadState.IDLE, HostThreadState.TURN_OPEN]) +def test_archived_history_does_not_require_the_owner_to_find_a_task_link(tmp_path, current): + bindings = [_binding(f"old-{i}") for i in range(4)] + [_binding("current")] + registry = _registry(tmp_path, bindings) + observed = [] + + def observe(ids): + observed.append(set(ids)) + return {item: HostThreadActivity(state=current if item == "current" else HostThreadState.ARCHIVED) + for item in ids} + + route = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", + observers={"codex-app": observe}) + assert route["status"] == "resolved" + assert route["selected_route"]["thread_id"] == "current" + assert route["candidate_count"] == 5 and len(route["candidates"]) == 3 + assert route["host_observation"]["state"] == current.value + assert route["host_delivery"] == "not_attempted" + assert observed == [{b["thread_id"] for b in bindings}] + + +@pytest.mark.parametrize("other", ["missing", "unsupported", "failed", "idle"]) +def test_unknown_or_another_readable_binding_never_becomes_archived(tmp_path, other): + registry = _registry(tmp_path, [_binding("current"), _binding("other", host_surface=other)]) + + def fail(ids): + raise OSError("private host failure") + + observers = _observer(HostThreadState.IDLE) + if other != "unsupported": + observers[other] = fail if other == "failed" else ( + lambda ids: {} if other == "missing" else {i: HostThreadActivity(state=HostThreadState.IDLE) for i in ids} + ) + route = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", observers=observers) + assert route["status"] == "ambiguous" + assert route["selected_route"] is None and route["host_delivery"] == "not_attempted" + assert "private host failure" not in json.dumps(route) + + +def test_all_archived_is_unavailable_but_explicit_archived_link_stays_exact(tmp_path): + registry = _registry(tmp_path, [_binding("old"), _binding("current")]) + archived = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", + observers=_observer(HostThreadState.ARCHIVED)) + assert archived["status"] == "unavailable" and archived["reason"] == "host_thread_archived" + assert archived["selected_route"] is None + explicit = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", + thread_link="codex://threads/old", observers={"codex-app": lambda ids: { + i: HostThreadActivity(state=HostThreadState.ARCHIVED if i == "old" else HostThreadState.IDLE) for i in ids}}) + assert explicit["status"] == "unavailable" and explicit["selected_route"] is None + + +def test_observation_cannot_rebind_the_selected_peer(tmp_path): + registry = _registry(tmp_path, [_binding("old"), _binding("current")]) + + def observe(ids): + registry.write_text(json.dumps({"goals": [{"id": "goal", "coordination": { + "registered_agents": ["reviewer", "builder"], + "thread_agent_bindings": [_binding("current", "builder")], + }}]})) + return {i: HostThreadActivity(state=HostThreadState.ARCHIVED if i == "old" else HostThreadState.IDLE) for i in ids} + + route = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", observers={"codex-app": observe}) + assert route["status"] == "ambiguous" and route["reason"] == "binding_identity_conflict" + assert route["selected_route"] is None + + +@pytest.mark.parametrize("change", ["new_binding", "unregistered", "reordered"]) +def test_observation_rechecks_the_candidate_set_and_registration(tmp_path, change): + registry = _registry(tmp_path, [_binding("old"), _binding("current")]) + + def observe(ids): + content = json.loads(registry.read_text()) + coordination = content["goals"][0]["coordination"] + if change == "new_binding": + coordination["thread_agent_bindings"].append(_binding("unobserved")) + elif change == "unregistered": + coordination["registered_agents"].remove("reviewer") + else: + coordination["thread_agent_bindings"].reverse() + registry.write_text(json.dumps(content)) + return {i: HostThreadActivity(state=HostThreadState.ARCHIVED if i == "old" else HostThreadState.IDLE) + for i in ids} + + route = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", + observers={"codex-app": observe}) + expected = {"new_binding": "ambiguous", "unregistered": "not_authorized", "reordered": "resolved"} + assert route["status"] == expected[change] + assert route["selected_route"] == ( + {"host_surface": "codex-app", "thread_id": "current"} if change == "reordered" else None + ) + assert route["host_delivery"] == "not_attempted" + + +def test_withheld_and_over_budget_candidates_cannot_hide_a_second_binding(tmp_path): + private_id = "ghp_" + "1234567890abcdefghijklmnopqrstuvwxyz1234" + registry = _registry(tmp_path, [_binding("current"), _binding(private_id)]) + result = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", + observers=_observer(HostThreadState.IDLE)) + assert result["status"] == "ambiguous" and result["selected_route"] is None + assert result["withheld_candidate_count"] == 1 and private_id not in json.dumps(result) + registry = _registry(tmp_path, [_binding(f"task-{i}") for i in range(33)]) + + def forbidden(ids): + pytest.fail("over-budget inventory must remain ambiguous without an unbounded host read") + + result = resolve_peer_host_route(registry, goal_id="goal", agent_id="reviewer", observers={"codex-app": forbidden}) + assert result["status"] == "ambiguous" and result["candidate_count"] == 33 + assert result["selected_route"] is None + + +def test_real_cli_uses_read_only_host_store_and_pins_one_request_after_archived_history(tmp_path): + registry = _registry(tmp_path, [_binding("old"), _binding("current")]) + home = tmp_path / "codex-home" + home.mkdir() + rollout = home / "current.jsonl" + rollout.write_text(json.dumps({"timestamp": "2026-01-01T00:00:00Z", "type": "event_msg", + "payload": {"type": "task_complete"}}) + "\n") + db = home / "state_5.sqlite" + with sqlite3.connect(db) as conn: + conn.execute("CREATE TABLE threads(id TEXT, rollout_path TEXT, archived INTEGER)") + conn.executemany("INSERT INTO threads VALUES(?,?,?)", [("old", "", 1), ("current", str(rollout), 0)]) + before = (registry.read_bytes(), db.read_bytes(), rollout.read_bytes()) + brief = tmp_path / "brief.json" + brief.write_text(json.dumps({"schema_version": "collaboration_brief_v0", "purpose": "Check the draft", + "context": "Independent review", "constraints": ["Do not publish"], "inputs": [], + "acceptance": ["Return findings"], "return_requirement": "Report to requester"})) + env = {**os.environ, "LOOPX_CODEX_HOMES": str(home)} + + def cli(*args): + result = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), + "--runtime-root", str(tmp_path / "runtime"), "--format", "json", *args], + env=env, capture_output=True, text=True, check=True) + return json.loads(result.stdout) + + preview = cli("resolve-peer-route", "--goal-id", "goal", "--agent-id", "reviewer") + assert preview["status"] == "resolved" and preview["selected_route"]["thread_id"] == "current" + args = ("manager-inbox", "request", "--goal-id", "goal", "--agent-id", "builder", + "--peer-agent-id", "reviewer", "--operation-id", "review-draft", "--brief-file", str(brief), "--require-host-route") + requested, replay = cli(*args), cli(*args) + assert requested["request_id"] == replay["request_id"] and replay["replayed"] is True + assert requested["host_delivery"]["thread_id"] == "current" + assert requested["host_delivery"]["status"] == "not_attempted" + assert before == (registry.read_bytes(), db.read_bytes(), rollout.read_bytes()) + + +def test_named_peer_with_two_readable_bindings_needs_an_exact_selection( tmp_path: Path, ) -> None: registry = _registry(tmp_path, [_binding("older"), _binding("current")]) diff --git a/tests/control_plane_ts/peer_host_route.test.ts b/tests/control_plane_ts/peer_host_route.test.ts new file mode 100644 index 0000000000..9a79dc7664 --- /dev/null +++ b/tests/control_plane_ts/peer_host_route.test.ts @@ -0,0 +1,43 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { selectObservedPeerHostRoute } from "../../loopx/control_plane/collaboration/peer_route_selection.ts"; + +const archived = { state: "archived" }; +const idle = { state: "idle" }; +const open = { state: "turn_open" }; +const unknown = { state: "unavailable", reason: "thread_not_found" }; +const resolved = (index: number) => ({ status: "resolved", reason: null, selected_index: index }); +const ambiguous = { status: "ambiguous", reason: "multiple_binding_candidates", selected_index: null }; +const absent = { status: "unavailable", reason: "host_thread_archived", selected_index: null }; + +test("one readable task survives only positively archived alternatives, in either order", () => { + for (const [observations, expected] of [ + [[idle], resolved(0)], [[open], resolved(0)], [[archived], absent], + [[archived, idle], resolved(1)], [[open, archived], resolved(0)], + [[archived, archived, open], resolved(2)], [[archived, archived], absent], + [[idle, open], ambiguous], [[unknown, idle], ambiguous], + [[idle, archived, unknown], ambiguous], [[unknown, unknown], ambiguous], + [[archived, unknown], { status: "unavailable", reason: "thread_not_found", selected_index: null }], + ] as const) { + assert.deepEqual(selectObservedPeerHostRoute({ observations: [...observations] }), expected); + } +}); + +test("unsupported, failed and withheld alternatives retain uncertainty", () => { + for (const reason of ["host_observer_unavailable", "host_observation_failed", "route_candidate_withheld", + "store_unavailable", "record_unrecognized", "no_turn_marker", "unsupported_host"]) { + const unavailable = { state: "unavailable", reason }; + assert.deepEqual(selectObservedPeerHostRoute({ observations: [idle, unavailable] }), ambiguous); + assert.deepEqual(selectObservedPeerHostRoute({ observations: [archived, unavailable] }), { + status: "unavailable", reason, selected_index: null, + }); + } +}); + +test("malformed or over-budget observations cannot choose a route", () => { + for (const observations of [[], Array(33).fill(archived), [{ state: "ready" }], + [{ state: "unavailable" }], [{ state: "unavailable", reason: "assume_archived" }]]) { + assert.throws(() => selectObservedPeerHostRoute({ observations })); + } + assert.deepEqual(selectObservedPeerHostRoute({ observations: [...Array(31).fill(archived), idle] }), resolved(31)); +});