diff --git a/docs/architecture/rfcs/agent-session-execution-modes-v0.md b/docs/architecture/rfcs/agent-session-execution-modes-v0.md index a21ce5aeb7..7290d198d7 100644 --- a/docs/architecture/rfcs/agent-session-execution-modes-v0.md +++ b/docs/architecture/rfcs/agent-session-execution-modes-v0.md @@ -182,7 +182,7 @@ Audited at `6c3da75ca`. These are current facts, not proposed behavior. | Attached broker | [`loopx/attached_session.py`](../../../loopx/attached_session.py) implements bind, claim, and complete under `loopx_attached_agent_session_broker_v0`, adapter kind `attached_host_session`, upstream mode `host_broker`, with a bounded claim wait of 1800 seconds, duplicate-safe claim and completion receipts, and per-binding file locks. | | Runtime fencing | [`loopx/chat_runtime.py`](../../../loopx/chat_runtime.py) never starts a managed adapter for an attached session and fails closed with typed errors such as `attached_session_live_steering_unavailable`, `live_steering_requires_active_turn`, and `live_steering_session_not_attached`. | | CLI surface | `loopx worker-bridge attached-session-bind`, `-list`, `-claim`, and `-complete` exist in [`loopx/cli_commands/worker_bridge.py`](../../../loopx/cli_commands/worker_bridge.py), documented in the [broker guide](../../integrations/attached-agent-session-broker.md) and the [worker-bridge install contract](../../integrations/worker-bridge-install-contract.md). | -| Existing-session delegation | [`loopx delegation`](../../reference/local-delegation.md#use-an-existing-agent-conversation-through-its-shell) exposes the same explicitly bound work as MCP to an existing shell-capable Agent. It retains the caller conversation and original operation on reconnect; it does not provision an Agent, migrate a host or install an automatic wake policy. | +| Existing-session delegation | [`loopx delegation`](../../reference/local-delegation.md#use-an-existing-agent-conversation-through-its-shell) exposes the same explicitly bound work as MCP to an existing shell-capable Agent. `operations` recovers requester-scoped work without remembered IDs, rechecks acceptance and exposes unavailable items and further pages. Newly tool-equipped Goal Chat consumes the same inventory; resumed native threads keep their original tool schema. It does not provision an Agent, migrate a host or install an automatic wake policy. | | Focused tests | [`tests/test_attached_session_cli.py`](../../../tests/test_attached_session_cli.py) and `tests/test_chat_codex_home.py::test_attached_session_uses_existing_host_not_managed_adapter` cover bind/claim/complete and the no-managed-adapter fence. | | Product-level proposal | The [Desktop execution frontends RFC](desktop-execution-frontends-v0.md) owns the Mode A/Mode B product comparison, the connector and event-source orthogonality, and the Desktop non-goals. | | Host-side loop guidance | [Codex CLI TUI loop](../../product/runtimes/codex-cli/codex-cli-tui-loop.md) documents session-attached automation and resume options for one visible host. | diff --git a/docs/architecture/rfcs/agent-session-execution-modes-v0.zh-CN.md b/docs/architecture/rfcs/agent-session-execution-modes-v0.zh-CN.md index 27f4fe768a..dedc3604dc 100644 --- a/docs/architecture/rfcs/agent-session-execution-modes-v0.zh-CN.md +++ b/docs/architecture/rfcs/agent-session-execution-modes-v0.zh-CN.md @@ -141,7 +141,7 @@ LoopX 启动,另一种已经属于其他宿主。当绑定没有说明自己 | 挂接 broker | [`loopx/attached_session.py`](../../../loopx/attached_session.py) 在 `loopx_attached_agent_session_broker_v0` 下实现 bind/claim/complete,适配器类型 `attached_host_session`,上游模式 `host_broker`,claim 等待上限 1800 秒,claim 与完成回执去重,并按绑定加文件锁。 | | 运行时围栏 | [`loopx/chat_runtime.py`](../../../loopx/chat_runtime.py) 绝不为挂接会话启动托管适配器,并以类型化错误失败关闭,例如 `attached_session_live_steering_unavailable`、`live_steering_requires_active_turn`、`live_steering_session_not_attached`。 | | CLI 面 | `loopx worker-bridge attached-session-bind`、`-list`、`-claim`、`-complete` 存在于 [`loopx/cli_commands/worker_bridge.py`](../../../loopx/cli_commands/worker_bridge.py),并在 [broker 指南](../../integrations/attached-agent-session-broker.md) 与 [worker-bridge 安装契约](../../integrations/worker-bridge-install-contract.md) 中记录。 | -| 原会话委派 | [`loopx delegation`](../../reference/local-delegation.md#use-an-existing-agent-conversation-through-its-shell) 让有 shell 能力的原 Agent 使用与 MCP 相同的显式执行绑定;重连保留原对话和操作身份,不创建 Agent、不迁移宿主,也不安装自动唤醒策略。 | +| 原会话委派 | [`loopx delegation`](../../reference/local-delegation.md#use-an-existing-agent-conversation-through-its-shell) 让有 shell 能力的原 Agent 使用与 MCP 相同的显式执行绑定;`operations` 无需记住 ID 即可找回自身委派,重新核验 accepted,明确单条不可用及剩余分页。新挂载工具的 Goal Chat 复用同一目录,已存在的原生线程恢复时保留原工具 schema;不创建 Agent、不迁移宿主,也不安装自动唤醒策略。 | | 聚焦测试 | [`tests/test_attached_session_cli.py`](../../../tests/test_attached_session_cli.py) 与 `tests/test_chat_codex_home.py::test_attached_session_uses_existing_host_not_managed_adapter` 覆盖 bind/claim/complete 与"不启动托管适配器"的围栏。 | | 产品级提案 | [桌面执行前端 RFC](desktop-execution-frontends-v0.zh-CN.md) 拥有 Mode A/Mode B 的产品对比、连接器与事件源正交性,以及桌面端非目标。 | | 宿主侧循环指引 | [Codex CLI TUI loop](../../product/runtimes/codex-cli/codex-cli-tui-loop.md) 记录了一个可见宿主的会话挂接自动化与恢复选项。 | diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 3391342d74..2fdd00c015 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -243,8 +243,10 @@ These priorities do not change live Goal quota or authorize experiments/cloud re priorities and attention; a project coordinator is an ordinary registered Agent accountable for a scoped objective, substantive investigation and synthesis. Members may coordinate narrower work through the same operations. Local Goal Chat -now reuses scoped evidence, semantic handoff and original-conversation return; -it is not the persistent coordinator or an implicit worker launch. See +reuses scoped evidence, semantic handoff and original-conversation return. +Explicit [Goal Chat LoopX mode](../../reference/goal-chat-continuation.md) now +adds native continuation, authorized member delegation and pause/recovery; +ordinary conversation does not implicitly launch workers. See [shared capabilities, local path and implementation order](../../reference/project-coordination.md). This advances R3's local entrypoint without closing R2/G1 or new Lark qualification. @@ -282,7 +284,7 @@ sessions, generic Agent creation, dynamic governed work derivation, complete inbox/queue/steer, authenticated remote authority and packaged frontend/Lark companion work remain R2/R3/R4/R6 boundaries. Existing Goals are not promoted. -An existing shell-capable coordinator can now use `delegation list/start/read/wait/resume` without replacing its session or loading new MCP tools. The synthetic example's `prepare` path creates only isolated operator bindings; the existing Agent chooses and starts the work. This completes the attached-caller entrypoint over the existing execution owner. Dynamic identity/profile provisioning, unattended lead wakeup and full inbox/queue/steer remain separate R2/R3 requirements; fixed binding readback is not fleet readiness. +An existing shell-capable coordinator uses `delegation list/operations/start/read/wait/resume` without replacing its session. Requester-scoped `operations` recovers durable work after context loss, independently rechecks accepted results and preserves unavailable branches and pagination; enabled MCP and newly tool-equipped Goal Chat use the same read model. Existing native threads retain their tool schema on resume. It starts no work and does not infer overall readiness from a display list. The example's `prepare` still only provisions isolated operator bindings. Next, feed actual execution/acceptance facts into existing R2 readiness, then extend existing registration/runtime configuration for approved identity/profile provisioning and qualify original-request return/lead continuation. Unattended wake, full cross-host inbox/queue/steer and Lark parity remain separate requirements; no G1/G3 promotion follows from this recovery entrypoint. ### R3: Semantic Requests and Automatic Return diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md index 5c5f3b93f3..ff111653c7 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.zh-CN.md @@ -234,8 +234,9 @@ R2 的一条依赖必须通过真实 LoopX Agent 间的请求/产物交接完成 **产品职责。** 管家负责所有者跨项目的上下文、取舍和注意力;项目 coordinator 是 对有范围目标、实质调查与综合负责的普通注册 Agent。成员可以用同一操作协调更小 -范围的工作。本地 Goal Chat 已复用范围内证据、语义交接和原对话返回;它不等于 -持久 coordinator,也不隐式启动 worker。见[共用能力、本地路径与实施顺序](../../reference/project-coordination.md)。 +范围的工作。本地 Goal Chat 复用范围内证据、语义交接和原对话返回;显式开启 +[LoopX 模式](../../reference/goal-chat-continuation.md)后支持原生持续推进、授权成员 +委派和暂停恢复,普通对话不隐式启动 worker。见[共用能力、本地路径与实施顺序](../../reference/project-coordination.md)。 这一进展补齐 R3 本地入口,不关闭 R2/G1,也不新增 Lark 资格声明。 - **入口与 owner:** 现有 session binding、Turn driver、quota/scheduler、manager runtime 配置;复用已有设置 editor,不新建 profile。 @@ -263,7 +264,7 @@ Todo 完成入口分别执行当前 pinned 检查,accepted 返回读 canonical 完成。长期 attached 会话、通用 Agent 创建、动态受治理工作派生、完整 inbox/queue/steer、 认证远端权威与 packaged frontend/Lark 配套仍归 R2/R3/R4/R6;不晋升已有 Goal。 -有 shell 能力的原 coordinator 现在可通过 `delegation list/start/read/wait/resume` 调用已有执行 owner,无需替换会话或重新加载 MCP 工具。合成示例的 `prepare` 只准备隔离绑定,由原 Agent 自行选择并启动工作。这闭合原会话调用入口;动态身份/profile 创建、无人值守唤醒和完整 inbox/queue/steer 仍按 R2/R3 推进,固定绑定读回不等于团队全部就绪。 +有 shell 能力的原 coordinator 可通过 `delegation list/operations/start/read/wait/resume` 调用已有执行 owner,无需替换会话。`operations` 从自身持久记录找回上下文丢失前的工作,重新核验 accepted,保留不可用分支与分页;已启用的 MCP 和新挂载工具的 Goal Chat 共用该读模型;已有原生线程恢复时保留原工具 schema。读取不启动工作,也不把展示列表当整体 readiness。合成示例 `prepare` 仍只准备隔离绑定。下一步先将实际执行/验收事实接入已有 R2 readiness,再沿现有注册及 runtime 配置扩展经授权的身份/profile 创建,并验证原请求返回与主力继续推进。无人值守唤醒、完整跨宿主 inbox/queue/steer 和 Lark 等价仍分别验收,不因新增恢复入口晋升 G1/G3。 ### R3:语义请求与自动回报 diff --git a/docs/reference/goal-chat-continuation.md b/docs/reference/goal-chat-continuation.md index bcea1fbcc1..ee049bf339 100644 --- a/docs/reference/goal-chat-continuation.md +++ b/docs/reference/goal-chat-continuation.md @@ -22,6 +22,18 @@ executable, workspace or acceptance rule. The same delegation service supports an authorized member coordinating further members. Results returned as `accepted` require current canonical task completion and unchanged artifacts. +In newly tool-equipped conversations, `action=operations` recovers the configured sender's +durable work with the same [paged inventory](local-delegation.md#recover-work-without-remembered-operation-ids) +as CLI/MCP. It includes work started outside this conversation. Follow +`next_cursor`, read the original operations and reconcile unavailable entries +before starting replacements. The recent member chips remain conversation +observations, not the complete inventory. Pausing still fences this tool. +Existing native threads retain their original tool schema when resumed; this +change does not replace an unfinished Goal to add a tool operation. Such a +thread keeps its original read/wait operations; a shell-capable caller can use +the CLI recovery entrypoint independently. Tool-schema upgrade remains a +separate session capability. + The composer shows native state, accumulated coordinator usage and last member observations. **Pause** stops the coordinator, while already delegated members continue under their independent deadlines and acceptance rules. **Continue** diff --git a/docs/reference/local-delegation.md b/docs/reference/local-delegation.md index 37a10d0283..6ed413eeb0 100644 --- a/docs/reference/local-delegation.md +++ b/docs/reference/local-delegation.md @@ -98,6 +98,53 @@ turn. The conversation remains persistent independently of whether autonomous LoopX mode is enabled. Current Dashboard/Lark setup is unchanged; those surfaces keep their existing conversation and runtime owners. +### Recover work without remembered operation ids + +After reconnecting or losing conversation context, use the same registered +requester and execution configuration: + +```bash +delegate operations --limit 10 +# When has_more is true, copy next_cursor from that response: +delegate operations --limit 10 --cursor "$NEXT_CURSOR" +delegate read --operation-id "$ORIGINAL_OPERATION_ID" +``` + +This reads the existing requester-scoped journal, including work created from +another conversation under that identity. Each item includes its original +operation/request/task identity and current execution readback. Accepted items +are independently rechecked against current canonical completion and artifacts; +the page includes artifact references/hashes, while `read` supplies full content. +One changed binding, corrupt record or invalid artifact yields `unavailable` +for that item and `page_readback_complete: false`; healthy siblings remain +visible. This is a reconciliation case, not permission to dispatch a replacement. +Failure to read the journal itself fails the command instead of returning empty. + +Pages contain at most 50 items. `has_more` is independent of page readback +completeness. Accepted-item checks rerun the existing pinned validators; use a +smaller page when those checks are expensive. Inventory is requested on demand, +not added to the dashboard polling loop. The cursor follows stable record addresses, not business priority; +this is a live listing, so restart paging to discover new records inserted before +the cursor. An empty page for one requester says nothing about other members or +whether the Goal is complete. Only explicit `start`/`resume` can launch execution. + +Enabled MCP exposes the same operation as `list_delegations`. Newly tool-equipped Goal Chat +uses `loopx_collaboration` with `action=operations`, optional `limit` and `cursor`. +It retains its existing sender/configuration pin and pause fence. Both the lead +and a coordinating member recover their own operations; creation ancestry grants +no access to another requester's journal. No new settings or background polling +are required, and disabling execution tools removes this tool with them. +Already enrolled native Chat threads keep their original tool schema on resume; +they are not replaced to install this new operation. Recovery guidance is part +of the new tool description, not injected into those older threads' shared prompt. + +中文:原对话重连后执行 `delegate operations`,不用先记住每个 operation ID。 +主力与承担协调的成员各自找回自己的工作,再用原 ID 读取完整结果;需要恢复时 +仍显式调用 `resume --execute`。分页回读会重新核验 accepted,单条失效显示 +`unavailable`,不能当成失败重派或静默隐藏。`has_more` 表示还有下一页, +`page_readback_complete` 只表示本页是否均成功读取;二者都不代表整个团队已完成。 +此入口不创建 Agent、不扩大授权,也不唤醒闲置的 Codex 对话。 + ## Use the same bindings through MCP Start the existing stdio server with the explicit opt-in: diff --git a/loopx/chat_loopx_mode.py b/loopx/chat_loopx_mode.py index e4dffa33cb..2ec10d1c5e 100644 --- a/loopx/chat_loopx_mode.py +++ b/loopx/chat_loopx_mode.py @@ -26,18 +26,22 @@ "Use action=bindings first; start requires binding_id, stable operation_id and brief with " "schema_version=collaboration_brief_v0, purpose, context, constraints (strings), inputs " "(relative ref, description, optional sha256), acceptance (strings), return_requirement. " - "Read/wait/resume use the original operation_id. Running is not failure; do not duplicate it.", + "Read/wait/resume use the original operation_id. Running is not failure; do not duplicate it. " + "After context loss, action=operations recovers this requester's durable work. Follow " + "next_cursor for more; unavailable means reconcile, not redispatch.", "inputSchema": { "type": "object", "additionalProperties": False, "properties": { "action": { "type": "string", - "enum": ["bindings", "start", "read", "wait", "resume", "messages"], + "enum": ["bindings", "operations", "start", "read", "wait", "resume", "messages"], }, "binding_id": {"type": "string"}, "operation_id": {"type": "string"}, "brief": {"type": "object"}, + "limit": {"type": "integer", "minimum": 1, "maximum": 50}, + "cursor": {"type": "string"}, }, "required": ["action"], }, @@ -505,9 +509,15 @@ def dispatch(name, arguments): "binding_id", "operation_id", "brief", + "limit", + "cursor", }: raise ValueError("invalid collaboration arguments") action = arguments.get("action") + if action != "operations" and ("limit" in arguments or "cursor" in arguments): + raise ValueError("pagination is only valid for operations") + if action == "operations" and set(arguments) - {"action", "limit", "cursor"}: + raise ValueError("operations reads a page; use read to select an operation") operation_id = arguments.get("operation_id", "") if action == "messages": rows = [ @@ -525,6 +535,8 @@ def dispatch(name, arguments): } elif action == "bindings": result = service.directory() + elif action == "operations": + result = service.operations(limit=arguments.get("limit", 20), cursor=arguments.get("cursor")) elif action == "start": result = service.start( arguments.get("binding_id", ""), diff --git a/loopx/cli_commands/delegation.py b/loopx/cli_commands/delegation.py index 341fc49146..c6db385b58 100644 --- a/loopx/cli_commands/delegation.py +++ b/loopx/cli_commands/delegation.py @@ -17,7 +17,7 @@ def register_delegation(subparsers, add_format): "delegation", help="Launch and recover authorized peer work; returns JSON." ) add_format(parser) - parser.add_argument("delegation_action", choices=("list", "start", "read", "wait", "resume")) + parser.add_argument("delegation_action", choices=("list", "operations", "start", "read", "wait", "resume")) parser.add_argument("--goal-id", required=True) parser.add_argument("--agent-id", required=True, help="Calling registered Agent, not the worker.") parser.add_argument("--execution-config", type=Path, required=True, @@ -26,6 +26,8 @@ def register_delegation(subparsers, add_format): parser.add_argument("--binding-id", help="For start: an authorized binding from list.") parser.add_argument("--brief-file", type=Path, help="For start: collaboration_brief_v0 JSON file.") parser.add_argument("--parent-request-id", help="For start: the request received by this coordinator.") + parser.add_argument("--limit", type=int, help="For operations: page size, 1–50 (default 20).") + parser.add_argument("--cursor", help="For operations: next_cursor returned by the previous page.") parser.add_argument("--execute", action="store_true", help="Required for start/resume; grants no additional authority.") @@ -39,10 +41,12 @@ def handle_delegation(args, registry_path, runtime_root): raise ValueError(f"delegation {action} requires --execute") if action not in {"start", "resume"} and args.execute: raise ValueError("--execute is only valid for start/resume") - if action != "list" and not args.operation_id: + if action not in {"list", "operations"} and not args.operation_id: raise ValueError(f"delegation {action} requires --operation-id") - if action == "list" and args.operation_id: - raise ValueError("list does not select an operation; use read") + if action in {"list", "operations"} and args.operation_id: + raise ValueError(f"{action} does not select an operation; use read") + if action != "operations" and (args.limit is not None or args.cursor is not None): + raise ValueError("limit and cursor are only supplied on operations") if action != "start" and (args.binding_id or args.brief_file or args.parent_request_id): raise ValueError("binding, brief and parent request are only supplied on start") service = Delegations(runtime_root, registry_path, args.goal_id, args.agent_id, @@ -58,6 +62,8 @@ def handle_delegation(args, registry_path, runtime_root): args.parent_request_id) elif action == "list": result = service.directory() + elif action == "operations": + result = service.operations(limit=20 if args.limit is None else args.limit, cursor=args.cursor) elif action == "read": result = service.read(args.operation_id) elif action == "wait": diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index 823d65352a..129e8a7c1c 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -156,6 +156,11 @@ def directory(self) -> dict: def path(self, operation_id: str) -> Path: return _root(self.root) / "executions" / _hash([self.goal_id, self.agent_id]) / (_hash(operation_id) + ".json") + def operations(self, *, limit: int = 20, cursor: str | None = None) -> dict: + from .control_plane.collaboration.delegation_inventory import read_delegation_inventory + + return read_delegation_inventory(self, limit=limit, cursor=cursor) + def start(self, binding_id: str, operation_id: str, brief: dict, parent_request_id: str | None = None) -> dict: binding = self.binding(binding_id, require_active=True) @@ -380,6 +385,15 @@ def list_execution_bindings() -> dict: """Read operator-authorized peer task bindings; registration alone cannot launch.""" return delegations.directory() + @server.tool() + def list_delegations(limit: int = 20, cursor: str | None = None) -> dict: + """Recover this requester's work after context loss. Follow next_cursor for more. + + Accepted items are rechecked; unavailable requires reconciliation, not duplicate + dispatch. Read the original operation for full artifacts. Listing starts no work. + """ + return delegations.operations(limit=limit, cursor=cursor) + @server.tool() def start_delegation(binding_id: str, operation_id: str, brief: dict, parent_request_id: str | None = None) -> dict: diff --git a/loopx/control_plane/collaboration/delegation.ts b/loopx/control_plane/collaboration/delegation.ts index ace85c2856..9d82f88880 100644 --- a/loopx/control_plane/collaboration/delegation.ts +++ b/loopx/control_plane/collaboration/delegation.ts @@ -37,6 +37,57 @@ const transitions: Record = { prepared: ["running", "rejected"], running: ["turn_returned", "rejected"], turn_returned: ["accepted", "rejected"], accepted: [], rejected: [], }; + +/** Page only the caller's existing journal. A cursor is not a fleet snapshot. */ +export function delegationInventoryQuery(params: JsonObject): JsonObject { + const limit = params.limit ?? 20; + const cursor = params.cursor ?? null; + requireThat(Number.isInteger(limit) && Number(limit) >= 1 && Number(limit) <= 50, + "delegation inventory limit must be between 1 and 50"); + requireThat(cursor === null || (typeof cursor === "string" && /^[a-f0-9]{64}$/.test(cursor)), + "invalid delegation inventory cursor"); + return {limit, cursor}; +} + +/** The host supplies a fresh Delegations.read result, never a saved status. */ +export function delegationInventoryItem(params: JsonObject): JsonObject { + const record = requireJsonObject(params.record, "delegation inventory record"); + requireThat(typeof record.record_id === "string" && /^[a-f0-9]{64}$/.test(record.record_id), + "invalid delegation record address"); + requireThat(record.operation_id === null || (typeof record.operation_id === "string" + && /^[A-Za-z0-9][A-Za-z0-9._-]{0,159}$/.test(record.operation_id)), "invalid delegation operation identity"); + if (params.observation === null) return { + record_id: record.record_id, operation_id: record.operation_id, + status: "unavailable", recovery_required: null, + error: "delegation_readback_unavailable", + }; + const observation = requireJsonObject(params.observation, "current delegation readback"); + requireThat(observation.operation_id === record.operation_id && record.operation_id !== null, + "delegation inventory identity mismatch"); + requireThat(Object.hasOwn(transitions, String(observation.status)), "invalid delegation observation"); + requireThat([observation.request_id, observation.agent_id, observation.todo_id].every(text), + "delegation request and task identities required"); + requireThat(typeof observation.worker_active === "boolean" + && typeof observation.recovery_required === "boolean", "current worker observation required"); + const result: JsonObject = { + record_id: record.record_id, operation_id: observation.operation_id, + request_id: observation.request_id, agent_id: observation.agent_id, todo_id: observation.todo_id, + status: observation.status, worker_active: observation.worker_active, + recovery_required: observation.recovery_required, + }; + if (observation.status === "accepted") { + requireThat(Array.isArray(observation.artifacts) && observation.artifacts.length > 0, + "accepted inventory requires current artifacts"); + result.artifacts = observation.artifacts.map(value => { + const artifact = requireJsonObject(value, "accepted artifact"); + requireThat(text(artifact.ref) && typeof artifact.sha256 === "string" + && /^[a-f0-9]{64}$/.test(artifact.sha256), "invalid accepted artifact reference"); + return {ref: artifact.ref, sha256: artifact.sha256}; + }); + } + return result; +} + export function transitionDelegationObservation(params: JsonObject): JsonObject { const from = params.from as Observation, to = params.to as Observation; requireThat(Object.hasOwn(transitions, from) && Object.hasOwn(transitions, to), "invalid delegation observation"); diff --git a/loopx/control_plane/collaboration/delegation_inventory.py b/loopx/control_plane/collaboration/delegation_inventory.py new file mode 100644 index 0000000000..193f13cf25 --- /dev/null +++ b/loopx/control_plane/collaboration/delegation_inventory.py @@ -0,0 +1,67 @@ +"""Bounded readback of an existing requester's durable delegation journal.""" +from __future__ import annotations + +import heapq +import re +from typing import TYPE_CHECKING + +from ..effect_runtime import EffectRuntimeRemoteError, effect_runtime_result +from .inbox import _read +from .peers import _goal, require_operation_id + +if TYPE_CHECKING: + from ...collaboration_mcp import Delegations + + +def read_delegation_inventory(service: Delegations, *, limit: int = 20, + cursor: str | None = None) -> dict: + query = effect_runtime_result("collaboration.delegation.inventory_query", + {"limit": limit, "cursor": cursor}) + _goal(service.registry, service.goal_id, service.agent_id) + directory = service.path("inventory").parent + + def addresses(): + try: + entries = directory.iterdir() + for path in entries: + if path.suffix != ".json": + continue + if not re.fullmatch(r"[a-f0-9]{64}", path.stem): + raise ValueError("unexpected delegation record address; reconcile inventory storage") + if query["cursor"] is None or path.stem > query["cursor"]: + yield path.stem + except FileNotFoundError: + # No journal yet is a valid empty inventory; other IO errors propagate. + if directory.exists(): + raise + + keys = heapq.nsmallest(query["limit"] + 1, addresses()) + items = [] + for key in keys[:query["limit"]]: + record = {"record_id": key, "operation_id": None} + try: + path = directory / (key + ".json") + if path.is_symlink() or not path.is_file(): + raise ValueError("delegation record unavailable") + operation_id = require_operation_id(_read(path)["identity"]["operation_id"]) + if service.path(operation_id) != path: + raise ValueError("delegation record identity mismatch") + record["operation_id"] = operation_id + observed = service.read(operation_id) + item = effect_runtime_result("collaboration.delegation.inventory_item", + {"record": record, "observation": observed}) + except (OSError, ValueError, KeyError, TypeError, EffectRuntimeRemoteError): + # A corrupt/stale branch must not hide healthy siblings or be called accepted. + item = effect_runtime_result("collaboration.delegation.inventory_item", + {"record": record, "observation": None}) + items.append(item) + more = len(keys) > query["limit"] + return { + "schema_version": "loopx_delegation_inventory_v0", + "items": items, "has_more": more, + "next_cursor": keys[query["limit"] - 1] if more else None, + "page_readback_complete": all(item["status"] != "unavailable" for item in items), + "note": "Live requester-scoped page, not a snapshot or proof that the whole Goal is complete. " + "Read original operations for full results; resume only when recovery is required. " + "Restart paging to discover work added before the cursor. Listing starts no work.", + } diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 34aa8a84f8..edc387e188 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,5 +1,5 @@ import {projectTodoSummaryLanes, projectLegacyTodoWorkCounts} from "./todos/summary_lanes.ts"; -import {selectDelegationBinding, transitionDelegationObservation} from "./collaboration/delegation.ts"; +import {delegationInventoryItem, delegationInventoryQuery, selectDelegationBinding, transitionDelegationObservation} from "./collaboration/delegation.ts"; import {planChatMode} from "./collaboration/chat_mode.ts"; import {resolveConversationScope} from "./collaboration/conversation_scope.ts"; import {previewTeamPlan, planTeamTransaction, teamTransactionIdentity} from "./work_items/team_plan.ts"; @@ -632,6 +632,8 @@ export function createEffectRuntimeHandlers( evaluatePostWritebackHookTransaction, ], ["collaboration.delegation.binding", selectDelegationBinding], + ["collaboration.delegation.inventory_query", delegationInventoryQuery], + ["collaboration.delegation.inventory_item", delegationInventoryItem], ["collaboration.chat_mode", planChatMode], ["collaboration.conversation.scope", resolveConversationScope], ["collaboration.delegation.observe", transitionDelegationObservation], diff --git a/tests/control_plane_ts/delegation.test.ts b/tests/control_plane_ts/delegation.test.ts index 9b1ebf8bd0..a20ed89400 100644 --- a/tests/control_plane_ts/delegation.test.ts +++ b/tests/control_plane_ts/delegation.test.ts @@ -1,6 +1,6 @@ import test from "node:test"; import assert from "node:assert/strict"; -import {selectDelegationBinding, transitionDelegationObservation} from "../../loopx/control_plane/collaboration/delegation.ts"; +import {delegationInventoryItem, delegationInventoryQuery, selectDelegationBinding, transitionDelegationObservation} from "../../loopx/control_plane/collaboration/delegation.ts"; const binding = {id: "review", agent_id: "reviewer", todo_id: "todo_review", workspace: "/fixture", requesters: ["coordinator", "analyst"], host_args: ["--host", "dsh"], timeout_seconds: 60, output_refs: ["output.json"]}; @@ -30,3 +30,23 @@ test("message receipt and model return do not imply accepted work", () => { assert.deepEqual(transitionDelegationObservation({from: "turn_returned", to: "accepted", canonical_done: true, acceptance_ready: true, artifacts_current: true}), {status: "accepted"}); }); + +test("inventory paging is bounded and never interprets a missing result as accepted", () => { + assert.deepEqual(delegationInventoryQuery({}), {limit: 20, cursor: null}); + for (const limit of [0, 51, true, "2"]) assert.throws(() => delegationInventoryQuery({limit})); + assert.throws(() => delegationInventoryQuery({cursor: "../other"})); + const record = {record_id: "a".repeat(64), operation_id: "review-1"}; + const observation = {operation_id: "review-1", request_id: "request", agent_id: "reviewer", + todo_id: "todo_review", status: "accepted", worker_active: false, recovery_required: false, + artifacts: [{ref: "output.json", sha256: "b".repeat(64), text: "private body"}]}; + const accepted = delegationInventoryItem({record, observation}); + assert.equal(accepted.status, "accepted"); + assert.equal(JSON.stringify(accepted).includes("private body"), false); + assert.throws(() => delegationInventoryItem({record, observation: {...observation, artifacts: []}})); + assert.throws(() => delegationInventoryItem({record, observation: {...observation, status: "done"}})); + assert.throws(() => delegationInventoryItem({record, observation: {...observation, operation_id: "other"}})); + const unavailable = delegationInventoryItem({record, observation: null}); + assert.equal(unavailable.status, "unavailable"); + assert.equal(unavailable.recovery_required, null); + assert.equal(unavailable.artifacts, undefined); +}); diff --git a/tests/test_chat_loopx_mode.py b/tests/test_chat_loopx_mode.py index ec33b3273f..e40f23742e 100644 --- a/tests/test_chat_loopx_mode.py +++ b/tests/test_chat_loopx_mode.py @@ -245,6 +245,11 @@ def test_pause_fences_dispatch_and_does_not_cancel_members(mode, monkeypatch): ][0]["id"] == "review" ) + inventory = adapter.session.read_tool_handler(TOOL["name"], {"action": "operations"}) + assert inventory["ok"] and inventory["items"] == [] + assert inventory["page_readback_complete"] and not inventory["has_more"] + invalid = adapter.session.read_tool_handler(TOOL["name"], {"action": "operations", "operation_id": "one"}) + assert invalid["error"] == "collaboration_request_rejected" # Pause persists before attempting potentially slow provider interruption. def interrupt(**_): @@ -258,6 +263,7 @@ def interrupt(**_): monkeypatch.setattr(service.controller, "interrupt_turn", interrupt) apply(mode, "pause") + assert adapter.session.read_tool_handler(TOOL["name"], {"action": "operations"})["error"] == "conversation_execution_inactive" with pytest.raises(Exception, match="active conversation execution"): apply(mode, "message", delivery_mode="queue", message="After pause") assert adapter.session.read_tool_handler("loopx_context_read", {}) == { diff --git a/tests/test_delegation_cli.py b/tests/test_delegation_cli.py index 956bca600a..ee751375ae 100644 --- a/tests/test_delegation_cli.py +++ b/tests/test_delegation_cli.py @@ -60,6 +60,17 @@ def test_attached_cli_disconnect_retry_and_verified_return(service): assert demo.canonical_tasks(root)["todo_analyst-initial"]["done"] assert result["artifacts"][0]["sha256"] + # A fresh CLI needs only the requester binding, not remembered operation ids. + status, inventory = cli(runner, "operations") + assert status == 0 and inventory["page_readback_complete"] + assert inventory["items"][0]["operation_id"] == "cli-work" + assert inventory["items"][0]["status"] == "accepted" + assert inventory["items"][0]["artifacts"][0]["sha256"] == result["artifacts"][0]["sha256"] + assert "text" not in inventory["items"][0]["artifacts"][0] + assert (root / "analyst" / "initial" / "host-invocations").read_text() == "1" + status, other = cli(runner, "operations", actor="reviewer") + assert status == 0 and other["items"] == [] + status, denied = cli(runner, "read", "--operation-id", "cli-work", actor="reviewer") assert status == 1 and not denied["ok"] # Changing an accepted artifact cannot be hidden behind the saved result. @@ -67,6 +78,10 @@ def test_attached_cli_disconnect_retry_and_verified_return(service): output.write_text("{}") status, stale = cli(runner, "read", "--operation-id", "cli-work") assert status == 1 and not stale["ok"] + status, inventory = cli(runner, "operations") + assert status == 0 and not inventory["page_readback_complete"] + assert inventory["items"][0]["status"] == "unavailable" + assert "artifacts" not in inventory["items"][0] def test_cli_invalid_inputs_do_not_launch_work(service): @@ -84,6 +99,11 @@ def test_cli_invalid_inputs_do_not_launch_work(service): status, result = cli(runner, "resume", "--operation-id", "missing") assert status == 1 and "--execute" in result["error"] assert not (root / "host-started").exists() + for arguments in [("--limit", "0"), ("--cursor", "invalid"), ("--execute",), ("--operation-id", "wrong")]: + status, result = cli(runner, "operations", *arguments) + assert status == 1 and not result["ok"] + status, result = cli(runner, "list", "--limit", "5") + assert status == 1 and not result["ok"] def test_shared_execution_host_does_not_require_optional_mcp(): diff --git a/tests/test_delegation_inventory.py b/tests/test_delegation_inventory.py new file mode 100644 index 0000000000..8486eb4a55 --- /dev/null +++ b/tests/test_delegation_inventory.py @@ -0,0 +1,77 @@ +"""Journal recovery must preserve scope and expose incomplete readback honestly.""" +import json + +import pytest + +from loopx.control_plane.collaboration.inbox import _read, _write +from test_local_delegation import brief, service as delegation_service + +service = delegation_service + + +def test_pages_find_all_original_operations_without_dispatch(service, monkeypatch): + root, runner = service + starts = [] + monkeypatch.setattr(runner, "_spawn", starts.append) + for i in range(7): + runner.start("analysis", f"work-{i}", brief()) + starts.clear() + seen = [] + cursor = None + while True: + page = runner.operations(limit=2, cursor=cursor) + assert page["page_readback_complete"] and len(page["items"]) <= 2 + seen.extend(row["operation_id"] for row in page["items"]) + if not page["has_more"]: + assert page["next_cursor"] is None + break + assert page["next_cursor"] != cursor + cursor = page["next_cursor"] + assert sorted(seen) == [f"work-{i}" for i in range(7)] + assert starts == [] and not (root / "host-started").exists() + # Revoking one binding is not authority to report prior work absent or accepted. + config = json.loads(runner.config.read_text()) + config["bindings"] = [] + runner.config.write_text(json.dumps(config)) + page = runner.operations() + assert not page["page_readback_complete"] and len(page["items"]) == 7 + assert all(row["status"] == "unavailable" for row in page["items"]) + + +def test_corruption_and_stopped_worker_do_not_hide_healthy_sibling(service, monkeypatch): + _, runner = service + monkeypatch.setattr(runner, "_spawn", lambda _: None) + for name in ["healthy", "stopped", "corrupt", "mismatched"]: + runner.start("analysis", name, brief()) + row = _read(runner.path("stopped")) + row["created_at"] = 0 + _write(runner.path("stopped"), row) + runner.path("corrupt").write_text("{") + row = _read(runner.path("mismatched")) + row["identity"]["operation_id"] = "healthy" + _write(runner.path("mismatched"), row) + page = runner.operations() + by_id = {row["operation_id"]: row for row in page["items"] if row["operation_id"]} + assert by_id["healthy"]["status"] == "prepared" + assert by_id["stopped"]["recovery_required"] + assert len(page["items"]) == 4 and not page["page_readback_complete"] + assert sum(row["status"] == "unavailable" for row in page["items"]) == 2 + + +def test_unknown_requester_and_unreadable_source_are_not_empty_inventory(service, monkeypatch): + _, runner = service + runner.agent_id = "unregistered" + with pytest.raises(ValueError, match="registered"): + runner.operations() + runner.agent_id = "lead" + directory = runner.path("inventory").parent + original = type(directory).iterdir + + def denied(path): + if path == directory: + raise PermissionError("fixture source denied") + return original(path) + + monkeypatch.setattr(type(directory), "iterdir", denied) + with pytest.raises(PermissionError): + runner.operations() diff --git a/tests/test_local_delegation.py b/tests/test_local_delegation.py index cf11688266..8b81939b4d 100644 --- a/tests/test_local_delegation.py +++ b/tests/test_local_delegation.py @@ -137,9 +137,14 @@ async def disconnect_requester(): async with stdio_client(params) as (reader, writer): async with ClientSession(reader, writer) as session: await session.initialize() + inventory = await session.call_tool("list_delegations", {}) + assert not inventory.isError and json.loads(inventory.content[0].text)["items"] == [] result = await session.call_tool("start_delegation", { "binding_id": "analysis", "operation_id": "analysis-1", "brief": brief()}) assert not result.isError + inventory = await session.call_tool("list_delegations", {}) + assert not inventory.isError + assert json.loads(inventory.content[0].text)["items"][0]["operation_id"] == "analysis-1" return json.loads(result.content[0].text) # Exiting the real stdio session closes the requesting MCP process. first = asyncio.run(disconnect_requester())