diff --git a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md index f8c791a9db..7dbe0e6f27 100644 --- a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md +++ b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md @@ -453,3 +453,79 @@ Before claiming the investment minimum loop, still prove installation, genuine human approval, bound native consumption, domain preflight and original-system evidence, accepted result and original-card/audience readback. The core PR requires owner review and is not self-installed before merge. + +### Preparation and canonical return-audience readback + +Proposal preparation is not domain execution. An admitted task with supplied +immutable terms may call `prepare` before confirmation or consumption; +`execution_allowed: false` is expected. `context/pending/inspect` also require +no consumption. Reconcile existing preparations first; neither preparation, +waiting nor final-answer prose proves task completion. Domain effects still +require the first successful consumption on the original native tool connection. + +For a managed executor, one distinct registered Goal/Agent return audience is +selected automatically. Multiple historical audiences require an explicit host +option, `--codex-operation-source-route-json +'{"host_surface":"codex-app","thread_id":"REGISTERED_THREAD"}'`. +The selector is registered-audience routing, not authentication, a session +replacement or an execution permit; the model cannot retarget it. Managed duplicate +bindings collapse. No registered audience retains the historical null route; +non-managed adapters retain their previous no-source-route projection and do not +invoke this managed resolver. Store +preparation freezes the selected audience in the original confirmation digest. + +Source-audience selection and the preparation prompt form a bounded backend +slice, not a completed product delivery. The companion personal-workspace +change below remains a separate first-screen-reviewed delivery in the same plan. +It reads the canonical action list while visible, with +scope-keyed cancellation, a bounded request timeout and no background interval +reads. The existing Manager brief links known-Goal pending operations to their +original drawer; overflow remains reachable through the existing conversation. +The drawer follows the same proposal ID through delivery, confirmation, +consumption, outcome and cancellation. A prepared request without a transport +receipt does not claim that a group card exists. Cancellation is administrative, +not an execution result. Failed readback marks cached state as potentially stale +and provides retry; browsing never confirms or executes an operation. + +Companion UI qualification uses isolated canonical backend snapshots through a loopback +HTTP fixture and the packaged EN/ZH desktop/mobile workspace. These synthetic +receipts do not establish genuine approval, live delivery or confirmation-to-host +wakeup latency. Actual click-to-original-session dispatch remains a separate +required acceptance condition using the existing scheduling/session owner. +The backend slice's focused native/normalization tests do not certify that the +companion UI ships in its commit or that source delivery is connected. + +### Automatic continuation and original-audience return: remaining work + +An authenticated confirmation callback currently records the claim and projects +the managed handoff; it does not launch the original host. Native `report` +commits an outcome and returns on the tool connection, while delivery recovery +updates the original Lark card. Neither is delivery to a selected source +conversation. An `outcome_observed` operation is absent from the pending-work +Inbox, so exhausting that Inbox cannot certify original-audience return. + +The same implementation Todo and capability owner retain both gaps. The next +bounded deliveries must prove the following in order; these are planned work, +not qualified features of the preparation/readback slice: + +| Delivery | Reused owner and required boundary | Exit evidence | +| --- | --- | --- | +| Confirmation-triggered continuation | The original Turn/session owner and an explicit operator-owned launch binding; the callback passes only an operation locator, never private terms or a new execution grant | Automatic bounded start of the same native session/profile after normal Goal, Todo, quota, lease and expiry checks; duplicate callbacks produce at most one active continuation | +| Restart recovery | The existing Turn journal and kernel single-flight lock; a terminal waiting Turn is not blindly resumed as a new invocation | Crash before/after launch and lost-response tests recover the original attempt; consumed/unknown work reconciles evidence without another submission; stop, session replacement and profile drift deny new launch | +| Source-audience return | The canonical operation/outcome and existing return-delivery semantics, with a qualified adapter for the selected audience | A durable attempt followed by actual readback of the same operation ID, outcome stage and digest at the original audience; unknown delivery is not blindly re-sent; Lark-card and source-conversation receipts remain distinct | +| Original-consumer qualification | The installed, pinned runtime and a real non-financial confirmation on the original test Todo | Actual confirmation, host-start, consumption, outcome and source-delivery timestamps; one consumption; healthy confirmation-to-host-start at most 60 seconds for the initial qualification | + +The 60-second threshold is an initial acceptance target, not a proven latency or +a scheduler guarantee. Host-start evidence must come from the native provider's +accepted Turn and matching connection metadata, not just process creation or a +locally stamped `started_at`. The adjacent [delegation lead-wake work](https://github.com/loopx-project/loopx/pull/5304) +targets an originating internal Goal Chat conversation; its planner/recovery +boundary is a reuse candidate, not proof of managed-operation confirmation wake +or external source-audience delivery. Frontend refresh intervals, a later heartbeat, manual +resume and engineering relay messages do not satisfy automatic continuation or +source return. A registered route alone is not a qualified delivery adapter; +an unrelated BotMux binding or attached Desktop identity is not a substitute. +Any new automatic launch configuration must reuse the existing configuration +owner and expose its effective state in the affected product entry points. +Keep the original consumer Todo open until these receipts are observed; do not +convert a cancelled test card or synthetic receipt into genuine approval. diff --git a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.zh-CN.md b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.zh-CN.md index 5dabf38c26..4ff7719db4 100644 --- a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.zh-CN.md +++ b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.zh-CN.md @@ -353,3 +353,58 @@ context 调用,结果均由既有 typed result validator 接受)、规范 不启动平行 resumed 执行者。宣称投研最小闭环前,仍须证明安装、真实用户批准、 绑定原生消费、垂域提交前检查与原系统证据、结果验收及原卡/受众读回。 Core PR 仍须 owner review,不在合并前自行安装。 + +### 准备提案与规范返回受众读回 + +准备提案不是领域执行。已准入任务提供不可变条款时,可以在确认和消费前调用 +`prepare`,其 `execution_allowed: false` 是正常结果;`context/pending/inspect` +也不要求先消费。先核对已有提案,不重复创建;准备、等待和最终答复文字都不能 +证明任务完成。领域副作用仍要求原生工具连接上的首次成功消费回执。 + +受管执行器只有一个不同的已登记 Goal/Agent 返回受众时自动选择;存在多个历史 +受众时,必须由宿主显式传入 `--codex-operation-source-route-json +'{"host_surface":"codex-app","thread_id":"REGISTERED_THREAD"}'`。 +这只是已登记的回传受众,不是认证、session 替换或执行许可;模型不能改投。 +受管重复绑定会去重;无登记受众保持历史 null 路由,非受管 adapter 保留既有 +无 source-route 投影,不调用这条受管解析。准备落盘后,返回受众进入原确认摘要,不可修改。 + +源受众选择与准备提示构成有界后台切片,不是完整产品交付。下面的个人工作台 +配套改动仍在同一计划内单独完成首屏评审后交付。工作台在可见时读取规范 +action list,按范围隔离查询并取消旧请求,限制 +单次请求时长,不在后台做间隔读取。复用管家简报将已知 Goal 的待确认操作连到 +原抽屉,溢出项通过既有对话可达。抽屉按同一提案 ID 跟随投递、确认、消费、 +结果和取消。没有投递回执的请求不能声称群卡已存在;取消只是行政终态,不是 +执行结果。读回失败明确提示缓存可能过期并允许重试;浏览不确认也不执行操作。 + +配套 UI 验收通过隔离的规范后端快照、loopback HTTP 夹具和打包页面,覆盖中英文及 +桌面/移动端。这些合成回执不证明真实用户批准、真实投递或确认后的即时唤醒。 +真实点击到原受管 session 的派发延迟仍是独立的必需验收项,必须复用现有调度 +和 session owner。 +后台切片的聚焦原生/规范化测试不证明其提交包含配套 UI,也不证明源投递已经接通。 + +### 自动续接与原受众返回:剩余交付 + +当前经过认证的确认回调会登记 claim 并投影受管交接,但不会启动原宿主。 +原生 `report` 将结果写回规范存储并在工具连接上返回;投递恢复更新原 Lark 卡片。 +两者都不等于已向选定的源会话投递。`outcome_observed` 操作不在待办 Inbox 中, +因此读完 Inbox 也不能证明已返回原受众。 + +这两个缺口继续由同一个实现 Todo 和 capability owner 负责。后续有界交付按以下 +顺序验收;这是规划,不是准备/读回切片已经具备或验证的功能: + +| 交付 | 复用的权威与必需边界 | 退出证据 | +| --- | --- | --- | +| 确认事件触发续接 | 原 Turn/session owner 与显式的 operator-owned 启动绑定;回调只传操作定位信息,不复制私有条款或新增执行许可 | 经过正常 Goal、Todo、quota、租约及过期核对后自动启动同一个原生 session/profile;重复回调最多产生一个活动续接 | +| 重启恢复 | 既有 Turn journal 与内核 single-flight 锁;不把已终止的等待 Turn 盲目当成新调用续跑 | 覆盖启动前后崩溃与响应丢失,恢复原尝试;已消费/未知操作只核对证据、不重新提交;停止、session 替换及 profile 漂移拒绝新启动 | +| 原受众返回 | 规范操作/结果与既有 return-delivery 语义,选定受众须有合格 adapter | 先登记投递尝试,再真实读回原受众处相同 operation ID、结果阶段及 digest;未知投递不盲目重发;Lark 原卡与源会话回执分开 | +| 原消费者验收 | 已安装且固定的 runtime,以及原测试 Todo 上真实的非金融确认 | 实际确认、宿主启动、消费、结果及源投递时间;一次消费;首轮健康路径从确认到宿主启动不超过 60 秒 | + +60 秒是首轮验收目标,不是已证明的延迟或调度保证。宿主启动必须由原生 provider +接受的 Turn 与匹配的连接元数据证明,仅创建进程或本地写入 `started_at` 不够。 +相邻的 [delegation lead 唤醒改动](https://github.com/loopx-project/loopx/pull/5304) +面向原发起的内部 Goal Chat 会话;其判定/恢复边界可供复用,不证明受管操作的 +确认唤醒或外部源受众已送达。前端刷新间隔、后续 heartbeat、 +手动 resume 及工程转述消息都不算自动续接或源返回。路由登记不等于投递 adapter +已合格;无关的 BotMux 绑定或附着 Desktop 身份不能替代。新增自动启动配置必须 +复用既有配置 owner,并在受影响的产品入口展示生效状态。看到上述回执前保留原 +消费者 Todo 为开放状态,不把已取消测试卡或合成回执当成真实批准。 diff --git a/loopx/chat_action_normalization.py b/loopx/chat_action_normalization.py index bc66278ec7..0a9dcdff85 100644 --- a/loopx/chat_action_normalization.py +++ b/loopx/chat_action_normalization.py @@ -50,6 +50,7 @@ def _normalize( "expires_at", "authorized_principals", "executor", + "source_route", }, ) if values.get("schema_version") != "loopx_operation_request_v0": @@ -233,26 +234,39 @@ def _normalize( field: _opaque(raw_executor.get(field), field=f"executor.{field}") for field in ("extension_id", "protocol", "permission", "revision") } - source_routes = [ - route - for route in (goal.get("coordination") or {}).get( - "thread_agent_bindings", [] + managed_source = executor.get("kind") == "managed_turn" + if not managed_source and "source_route" in values: + raise ValueError("source route selection requires a managed executor") + source_route = None + if managed_source: + from .control_plane.effect_runtime import effect_runtime_result + from .thread_agent_binding import ( + collect_accepted_bindings, + resolve_thread_agent_binding, ) - if isinstance(route, Mapping) and route.get("agent_id") == agent_id - and all(isinstance(route.get(key), str) and route[key] - for key in ("host_surface", "thread_id")) - ] - source_route = ( - { - "goal_id": goal_id, - **{ - key: source_routes[0][key] - for key in ("agent_id", "host_surface", "thread_id") + + selected_route = values.get("source_route") + if isinstance(selected_route, Mapping): + selected_binding = resolve_thread_agent_binding( + goal, + host_surface=selected_route.get("host_surface"), + thread_id=selected_route.get("thread_id"), + ) + selected_route = { + **selected_route, + "host_surface": selected_binding["host_surface"], + "thread_id": selected_binding["thread_id"], + } + + source_route = effect_runtime_result( + "operation.source_route.resolve", + { + "goal_id": goal_id, + "agent_id": agent_id, + "bindings": collect_accepted_bindings([goal]), + "selected_route": selected_route, }, - } - if len(source_routes) == 1 - else None - ) + )["source_route"] expires_at = parse_timestamp( _text(values.get("expires_at"), field="expires_at", limit=80) ) diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 8e9b908103..82d2c8160e 100644 --- a/loopx/cli_commands/turn.py +++ b/loopx/cli_commands/turn.py @@ -131,6 +131,8 @@ def handle_turn_command( strict_goal_admission = goal_admission if goal_admission.enabled else None if getattr(args, "codex_operation_tools", False) and args.host != "codex-cli": raise ValueError("--codex-operation-tools requires the codex-cli host") + if getattr(args, "codex_operation_source_route_json", None) is not None and not getattr(args, "codex_operation_tools", False): + raise ValueError("--codex-operation-source-route-json requires --codex-operation-tools") # Planning and dry-run execution inspect existing admitted intents. # Only an executing wake may sync inboxes or reserve a calendar window. turn_start_hook_dispatch = {} @@ -1002,7 +1004,8 @@ def run_built_in_host( ) return run_codex_operation_host( - request, registry_path=registry_path, **options + request, registry_path=registry_path, + source_route=getattr(args, "codex_operation_source_route_json", None), **options ) return run_codex_cli_host(request, **options) diff --git a/loopx/cli_commands/turn_registration.py b/loopx/cli_commands/turn_registration.py index 67ee763808..9f1be95320 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -224,6 +224,11 @@ def register_turn_commands( action="store_true", help="Opt in to the owned app-server operation transport for this admitted codex-cli Turn. Reuses the original Todo/session; does not authenticate an attached Desktop or grant domain effects.", ) + run_once.add_argument( + "--codex-operation-source-route-json", + type=json.loads, + help="Registered return audience {host_surface,thread_id} for operation proposals. Required when this Agent has several source routes; not executor authentication or execution permission.", + ) run_once.add_argument( "--codex-reasoning-effort", choices=list(REASONING_EFFORTS), diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index cf04abcb75..492dbf1291 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -1,5 +1,5 @@ import {manageNewGoalStorage} from "./coordination/local_authority_defaults.ts"; -import {deriveAgentOperationActor, managedOperationBindingCurrent, normalizeAgentOperationExecutor, planAgentOperationHandoff, projectAgentOperationInbox, projectManagedOperationTransport} from "./work_items/operation_agent_handoff.ts"; +import {deriveAgentOperationActor, managedOperationBindingCurrent, normalizeAgentOperationExecutor, planAgentOperationHandoff, projectAgentOperationInbox, projectManagedOperationTransport, resolveOperationSourceRoute} from "./work_items/operation_agent_handoff.ts"; import {projectDecisionNotice} from "./presentation/decision_notice.ts"; import {normalizeResearchObservation, validateResearchAttribution, projectResearchFrontier} from "./capabilities/explore_research.ts"; import {projectTodoSummary} from "./todos/summary_projection.ts"; @@ -632,6 +632,7 @@ export function createEffectRuntimeHandlers( ["presentation.action_review_plan.compile", (params) => compileActionReviewPlan(params.proposal)], ["operation.agent_executor.normalize", normalizeAgentOperationExecutor], + ["operation.source_route.resolve", resolveOperationSourceRoute], ["operation.managed_binding.current", managedOperationBindingCurrent], ["operation.managed_transport.project", projectManagedOperationTransport], ["operation.agent_handoff.actor", deriveAgentOperationActor], diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 783c10f436..b65abbe045 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -24,6 +24,7 @@ pending_operation_handoffs, ) from ..goals.first_party_host_admission import FirstPartyHostGoalAdmission +from ..effect_runtime import EffectRuntimeRejected from .codex_cli import ( _lineage, _prompt, @@ -77,6 +78,7 @@ def operation_tool_handler( profile_digest: str, model: str, reasoning_effort: str, + source_route: Mapping[str, str] | None = None, goal_admission: FirstPartyHostGoalAdmission | None = None, ): """Private closure installed only on a process owned by the Turn driver. @@ -154,6 +156,10 @@ def handle(tool: str, arguments: Any, native: dict[str, Any]) -> dict[str, Any]: if action == "prepare": request = dict(arguments["request"]) terms = dict(request.get("normalized_parameters") or {}) + if "source_route" in terms and terms["source_route"] != source_route: + raise ValueError("the model cannot retarget the host-selected return audience") + if source_route is not None: + terms["source_route"] = dict(source_route) for key, expected in { "goal_id": lineage["goal_id"], "agent_id": lineage["agent_id"], @@ -197,6 +203,12 @@ def handle(tool: str, arguments: Any, native: dict[str, Any]) -> dict[str, Any]: outcome=arguments.get("outcome"), ), } + except EffectRuntimeRejected as exc: + if exc.diagnostic_code == "operation_source_route_ambiguous": + return {"ok": False, "error": "operation_source_route_ambiguous", + "execution_allowed": False, + "next_action": "Select a registered return audience with --codex-operation-source-route-json in the host invocation; source routing is not executor authority."} + return {"ok": False, "error": "operation_admission_rejected", "execution_allowed": False} except (ValueError, KeyError, TypeError, RuntimeError): # Private payloads, paths and adapter error text never enter the model tool error. return { @@ -218,6 +230,7 @@ def run_codex_operation_host( sandbox: str = "read-only", model: str | None = None, reasoning_effort: str | None = None, + source_route: Mapping[str, str] | None = None, mcp_server: Mapping[str, Any] | None = None, timeout_seconds: float = 115, goal_admission: FirstPartyHostGoalAdmission | None = None, @@ -340,12 +353,18 @@ def store_binding(): profile_digest=profile_digest, model=model, reasoning_effort=reasoning_effort, + source_route=source_route, goal_admission=goal_admission, ) return session.send( _prompt(request) - + "\nUse loopx_operation for pending/prepare/inspect/consume/report. " - "Source conversations are not executor identity. Execute only after the first receipt says execution_allowed=true. " + + "\nUse loopx_operation for context/pending/prepare/inspect/consume/report. " + "Source conversations are not executor identity. context/pending/inspect do not require consumption. " + "When the admitted task authorizes proposal preparation and supplies its terms, prepare may run before human confirmation or consume; " + "prepare writes a canonical proposal, not a domain/external effect. execution_allowed=false is expected for prepare, not a reason to refuse it. " + "Use pending/inspect to reconcile an existing proposal before preparing another; do not invent missing terms or repeat an already granted preparation approval. " + "Only domain/external effects require the first successful consume receipt with execution_allowed=true. " + "Preparation or waiting for human confirmation is not task completion. Report outcomes only with original execution evidence. " "Already consumed/unknown effects require evidence reconciliation, never retry. Never treat final-answer prose as an outcome receipt.", output_schema=codex_cli_result_schema(request), ) diff --git a/loopx/control_plane/work_items/operation_agent_handoff.ts b/loopx/control_plane/work_items/operation_agent_handoff.ts index 0845ca29ec..3b93c87da2 100644 --- a/loopx/control_plane/work_items/operation_agent_handoff.ts +++ b/loopx/control_plane/work_items/operation_agent_handoff.ts @@ -27,6 +27,41 @@ function timestamp(value: unknown): number { return parsed; } +/** Python supplies owner-normalized thread bindings and selector tokens. A + * return audience is not an executor ID: do not narrow the registry vocabulary. + * Select explicitly under ambiguity; never freeze an ambiguous null route. */ +export function resolveOperationSourceRoute(input: JsonObject): JsonObject { + const goal = id(input.goal_id, "goal_id"); + const agent = id(input.agent_id, "agent_id"); + const bindings = Array.isArray(input.bindings) ? input.bindings : []; + const routes = new Map(); + for (const raw of bindings) { + if (raw === null || typeof raw !== "object" || Array.isArray(raw)) continue; + const route = raw as JsonObject; + if (route.agent_id !== agent) continue; + const audience = {goal_id: goal, agent_id: agent, + host_surface: requireNonEmptyString(route.host_surface, "registered source host_surface"), + thread_id: requireNonEmptyString(route.thread_id, "registered source thread_id")}; + routes.set(JSON.stringify(audience), audience); + } + if (input.selected_route != null) { + const selected = requireJsonObject(input.selected_route, "source route"); + const keys = Object.keys(selected); + if (keys.length !== 2 || !keys.includes("host_surface") || !keys.includes("thread_id")) { + throw new EffectRuntimeRequestError("source route selects only host_surface and thread_id"); + } + const audience = {goal_id: goal, agent_id: agent, + host_surface: requireNonEmptyString(selected.host_surface, "source host_surface"), + thread_id: requireNonEmptyString(selected.thread_id, "source thread_id")}; + requireThat(routes.has(JSON.stringify(audience)), "operation source route is not registered for this Goal and Agent"); + return {source_route: audience}; + } + if (routes.size > 1) throw new EffectRuntimeRequestError( + "operation source route is ambiguous; select a registered return audience in the host invocation", + "operation_source_route_ambiguous"); + return {source_route: routes.size === 1 ? routes.values().next().value : null}; +} + /** Readback of an operator-selected transport. This is not a session binding, * a runtime qualification or an execution permit. Python supplies argv facts. * Accept the CLI's provider-neutral effort vocabulary; actual model support diff --git a/tests/control_plane_ts/operation_agent_handoff.test.ts b/tests/control_plane_ts/operation_agent_handoff.test.ts index bd65cce56c..a61cabefd2 100644 --- a/tests/control_plane_ts/operation_agent_handoff.test.ts +++ b/tests/control_plane_ts/operation_agent_handoff.test.ts @@ -2,7 +2,35 @@ import assert from "node:assert/strict"; import test from "node:test"; import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; import {AGENT_OPERATION_REVISION, MANAGED_OPERATION_REVISION, managedOperationBindingCurrent, deriveAgentOperationActor, normalizeAgentOperationExecutor, planAgentOperationHandoff, - projectAgentOperationInbox, projectManagedOperationTransport} from "../../loopx/control_plane/work_items/operation_agent_handoff.ts"; + projectAgentOperationInbox, projectManagedOperationTransport, resolveOperationSourceRoute} from "../../loopx/control_plane/work_items/operation_agent_handoff.ts"; + +test("registered return audience is explicit under ambiguity and never executor identity", () => { + const route = {agent_id: "agent", host_surface: "codex-app", thread_id: "original"}; + const value = {goal_id: "goal", agent_id: "agent", bindings: [route]}; + assert.deepEqual(resolveOperationSourceRoute(value).source_route, {goal_id: "goal", ...route}); + assert.deepEqual(resolveOperationSourceRoute({...value, bindings: [route, route]}), resolveOperationSourceRoute(value)); + assert.equal(resolveOperationSourceRoute({...value, bindings: []}).source_route, null); + const multiple = {...value, bindings: [route, {...route, thread_id: "historical"}]}; + assert.throws(() => resolveOperationSourceRoute(multiple), {code: "operation_source_route_ambiguous"}); + const selector = {host_surface: route.host_surface, thread_id: route.thread_id}; + assert.deepEqual(resolveOperationSourceRoute({...multiple, selected_route: selector}), resolveOperationSourceRoute(value)); + for (const selected_route of ["original", {...selector, thread_id: "unregistered"}, + {...selector, host_surface: "other-host"}, {...selector, verified: true}]) { + assert.throws(() => resolveOperationSourceRoute({...multiple, selected_route})); + } + assert.throws(() => resolveOperationSourceRoute({...multiple, agent_id: "other-agent", selected_route: selector})); +}); + +test("source routing preserves owner-normalized audience tokens, not executor ID grammar", () => { + const route = {agent_id: "agent", host_surface: "custom-host", thread_id: "source@current"}; + const value = {goal_id: "goal", agent_id: "agent", bindings: [route]}; + assert.deepEqual(resolveOperationSourceRoute(value).source_route, {goal_id: "goal", ...route}); + const multiple = {...value, bindings: [route, {...route, thread_id: "safe-second"}]}; + assert.throws(() => resolveOperationSourceRoute(multiple), {code: "operation_source_route_ambiguous"}); + const selected_route = {host_surface: route.host_surface, thread_id: route.thread_id}; + assert.deepEqual(resolveOperationSourceRoute({...multiple, selected_route}), resolveOperationSourceRoute(value)); + assert.throws(() => resolveOperationSourceRoute({...multiple, selected_route: {...selected_route, thread_id: "unregistered+source"}})); +}); function input(): JsonObject { const executor = {kind: "agent_session", host_surface: "codex-app", thread_id: "original-thread", diff --git a/tests/test_chat_operation_actions.py b/tests/test_chat_operation_actions.py index 3e9e0f06c3..21a980db9c 100644 --- a/tests/test_chat_operation_actions.py +++ b/tests/test_chat_operation_actions.py @@ -182,6 +182,71 @@ def test_owned_managed_tool_uses_canonical_approval_once_without_desktop_binding assert recovered["ok"] is True and recovered["needs_reconciliation"] is False +@pytest.mark.parametrize("selector", [None, {"host_surface": "codex-app", "thread_id": "source"}]) +def test_non_managed_prepare_rejects_source_selector_even_when_null(tmp_path: Path, selector) -> None: + service, store = _service(tmp_path) + request = _request() + request["normalized_parameters"]["source_route"] = selector + with pytest.raises(ValueError, match="source route selection requires a managed executor"): + service.preview(request) + assert store.list() == [] + + +def test_non_managed_prepare_does_not_mount_managed_source_resolution(tmp_path: Path, monkeypatch) -> None: + from loopx.control_plane import effect_runtime + + original = effect_runtime.effect_runtime_result + calls = [] + + def record(method, *args, **kwargs): + calls.append(method) + return original(method, *args, **kwargs) + + monkeypatch.setattr(effect_runtime, "effect_runtime_result", record) + service, store = _service(tmp_path) + proposal = service.preview(_request()) + assert "source_route" not in store.load(proposal["proposal_id"])["normalized_parameters"] + assert "operation.source_route.resolve" not in calls + + +def test_managed_prepare_selects_registered_return_audience_without_rebinding_executor(tmp_path: Path) -> None: + from loopx.control_plane.turn_driver.codex_operation_host import operation_tool_handler + + service, store = _service(tmp_path) + registry = json.loads(service.registry_path.read_text()) + route = {"agent_id": "finance-fixture-agent", "host_surface": "codex-app", "thread_id": "source-current"} + registry["goals"][0]["coordination"]["thread_agent_bindings"] = [ + route, {**route, "thread_id": "source-historical"}, + ] + service.registry_path.write_text(json.dumps(registry)) + original = _managed_handler(service, store) + request = _request() + request["normalized_parameters"].pop("executor") + request["normalized_parameters"]["projection"]["simulated"] = False + native = {"thread_id": "owned-managed-thread", "host_turn_id": "native-turn-1"} + rejected = original("loopx_operation", {"action": "prepare", "request": request}, native) + assert rejected["error"] == "operation_source_route_ambiguous" + assert store.list() == [] + selector = {key: route[key] for key in ("host_surface", "thread_id")} + selected = operation_tool_handler( + runtime_root=store.root.parent.parent, registry_path=service.registry_path, + lineage={"goal_id": GOAL_ID, "agent_id": "finance-fixture-agent", "todo_id": "todo-managed"}, + session_id=native["thread_id"], profile_digest="c" * 64, model="test-model", + reasoning_effort="xhigh", source_route=selector, + ) + model_override = {**request, "normalized_parameters": {**request["normalized_parameters"], + "source_route": {**selector, "thread_id": "source-historical"}}} + assert selected("loopx_operation", {"action": "prepare", "request": model_override}, native)["ok"] is False + prepared = selected("loopx_operation", {"action": "prepare", "request": request}, native) + assert prepared["ok"] and not prepared["execution_allowed"] + proposal = prepared["proposal"] + assert proposal["normalized_parameters"]["source_route"] == {"goal_id": GOAL_ID, **route} + assert proposal["normalized_parameters"]["executor"]["session_id"] == native["thread_id"] + assert not selected("loopx_operation", {"action": "consume", "proposal_id": proposal["proposal_id"], + "consumption_id": "not-approved"}, native)["ok"] + assert len(store.list()) == 1 + + def test_managed_pending_reuses_registered_agent_and_goal_instance_scope( tmp_path: Path, ) -> None: diff --git a/tests/test_codex_operation_host.py b/tests/test_codex_operation_host.py index 7256e3710f..1d6faa56d0 100644 --- a/tests/test_codex_operation_host.py +++ b/tests/test_codex_operation_host.py @@ -58,6 +58,14 @@ def emit(value): if os.environ.get("FAKE_OPERATION_HANG") == "1": while True: time.sleep(.1) text = row["params"]["input"][0]["text"] + # The actual native prompt must separate proposal preparation from + # domain execution; otherwise approval can never get a proposal. + assert "context/pending/inspect do not require consumption" in text + assert "prepare may run before human confirmation or consume" in text + assert "prepare writes a canonical proposal" in text + assert "execution_allowed=false is expected for prepare" in text + assert "Only domain/external effects require the first successful consume" in text + assert "Execute only after the first receipt" not in text import re key = re.search(r'"turn_key":"([^"]+)"', text).group(1) emit({"id": row["id"], "result": {"turn": {"id": turn}}}) @@ -77,6 +85,174 @@ def emit(value): """ +SOURCE_PREPARE_SERVER = """#!/usr/bin/env python3 +import json, os, sys +thread, turn = "owned-source-test-thread", "native-source-test-turn" +calls = json.loads(os.environ["FAKE_SOURCE_TOOL_CALLS"]) +results = [] +def emit(value): + print(json.dumps(value), flush=True) +def call_next(): + if len(results) < len(calls): + emit({"id": 100 + len(results), "method": "item/tool/call", "params": { + "threadId": thread, "turnId": turn, "tool": "loopx_operation", + "arguments": calls[len(results)]}}) + else: + emit({"method": "item/agentMessage/delta", "params": { + "threadId": thread, "turnId": turn, "delta": json.dumps({"results": results})}}) + emit({"method": "turn/completed", "params": { + "threadId": thread, "turn": {"id": turn, "status": "completed"}}}) +for line in sys.stdin: + row = json.loads(line) + if row.get("method") == "initialize": + emit({"id": row["id"], "result": {}}) + elif row.get("method") == "thread/start": + emit({"id": row["id"], "result": {"thread": {"id": thread}, + "model": "test-model", "reasoningEffort": "xhigh"}}) + elif row.get("method") == "turn/start": + emit({"id": row["id"], "result": {"turn": {"id": turn}}}) + call_next() + elif isinstance(row.get("id"), int) and row["id"] >= 100: + result = json.loads(row["result"]["contentItems"][0]["text"]) + results.append(result) + if result.get("proposal") and len(results) == 1: + calls.append({"action": "consume", "proposal_id": result["proposal"]["proposal_id"], + "consumption_id": "unapproved-native-attempt"}) + call_next() +""" + + +@pytest.mark.parametrize( + "thread_id", ["source@current", "accepted+thread", "source:alternate"] +) +@pytest.mark.parametrize( + "multiple,selected", [(False, False), (True, False), (True, True)] +) +def test_registered_source_tokens_survive_owned_native_prepare_and_store_readback( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + thread_id: str, + multiple: bool, + selected: bool, +) -> None: + from examples.operation_action_fixtures import ( + GOAL_ID, + managed_handler, + request, + service, + ) + from loopx.chat_agent import CodexChatAgentSession + from loopx.control_plane.turn_driver.codex_operation_host import ( + OPERATION_TOOL, + operation_tool_handler, + ) + from loopx.thread_agent_binding import ( + bind_thread_agent_in_registry, + resolve_registry_thread_agent_binding, + ) + + api, store = service(tmp_path) + route = {"host_surface": "custom-host", "thread_id": thread_id} + for source in [ + route, + *( + [{"host_surface": "custom-host", "thread_id": "safe-second"}] + if multiple + else [] + ), + ]: + bound = bind_thread_agent_in_registry( + registry_path=api.registry_path, + goal_id=GOAL_ID, + agent_id="finance-fixture-agent", + execute=True, + **source, + ) + assert bound["ok"] and bound["written"] + assert ( + resolve_registry_thread_agent_binding( + registry_path=api.registry_path, **source + )["status"] + == "bound" + ) + value = request() + value["normalized_parameters"].pop("executor") + value["normalized_parameters"]["projection"]["simulated"] = False + arguments = {"action": "prepare", "request": value} + monkeypatch.setenv("FAKE_SOURCE_TOOL_CALLS", json.dumps([arguments, arguments])) + executable = tmp_path / "fake-source-prepare-server" + executable.write_text(SOURCE_PREPARE_SERVER) + executable.chmod(0o700) + session = CodexChatAgentSession.start( + codex_bin=str(executable), + work_dir=tmp_path, + goal_id=GOAL_ID, + objective="Synthetic source preparation", + execution_mode=True, + sandbox="read-only", + model="test-model", + reasoning_effort="xhigh", + dynamic_tools=[OPERATION_TOOL], + hard_timeout_sec=10, + response_timeout_sec=5, + ) + try: + session.bound_tool_handler = managed_handler( + api, store, session_id=session.thread_id + ) + if selected: + session.bound_tool_handler = operation_tool_handler( + runtime_root=store.root.parent.parent, + registry_path=api.registry_path, + lineage={ + "goal_id": GOAL_ID, + "agent_id": "finance-fixture-agent", + "todo_id": "todo-managed", + }, + session_id=session.thread_id, + profile_digest="c" * 64, + model="test-model", + reasoning_effort="xhigh", + source_route=route, + ) + results = session.send( + "Prepare the supplied synthetic proposal; no domain effects.", + output_schema={ + "type": "object", + "properties": {"results": {"type": "array"}}, + "required": ["results"], + }, + )["results"] + finally: + session.bound_tool_handler = None + session.close() + if multiple and not selected: + assert all( + result.get("error") == "operation_source_route_ambiguous" + for result in results + ) + assert store.list() == [] + else: + first, replay, refused = results + assert first["ok"] and first["execution_allowed"] is False + persisted = store.load(first["proposal"]["proposal_id"]) + assert persisted["normalized_parameters"]["source_route"] == { + "goal_id": GOAL_ID, + "agent_id": "finance-fixture-agent", + **route, + } + executor = persisted["normalized_parameters"]["executor"] + assert ( + executor["session_id"] == "owned-source-test-thread" + and executor["profile_digest"] == "c" * 64 + ) + assert ( + replay["proposal"]["proposal_id"] == persisted["proposal_id"] + and len(store.list()) == 1 + ) + assert refused["ok"] is False + + def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same_profile( tmp_path: Path, ) -> None: @@ -90,6 +266,7 @@ def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same "codex_bin": str(executable), "model": "test-model", "reasoning_effort": "xhigh", + "source_route": {"host_surface": "codex-app", "thread_id": "source-one"}, "timeout_seconds": 5, } first = run_codex_operation_host(_request(), **options) @@ -103,7 +280,8 @@ def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same ) assert binding["operation_transport"] == "app-server-operation-tools-v0" second = run_codex_operation_host( - _request(session_action="resume", turn_key="sha256:" + "b" * 64), **options + _request(session_action="resume", turn_key="sha256:" + "b" * 64), + **{**options, "source_route": {"host_surface": "codex-app", "thread_id": "source-two"}}, ) assert second["turn_key"] == "sha256:" + "b" * 64 assert ( @@ -117,7 +295,7 @@ def test_owned_app_server_process_authenticates_native_metadata_and_resumes_same with pytest.raises(ValueError, match="original managed transport"): run_codex_cli_host( _request(session_action="resume"), - **{key: value for key, value in options.items() if key != "registry_path"}, + **{key: value for key, value in options.items() if key not in {"registry_path", "source_route"}}, ) @@ -135,11 +313,20 @@ def forbidden_plain_cli(*args, **kwargs): pytest.fail("Operation opt-in must not downgrade to plain Codex exec") monkeypatch.setattr("loopx.cli_commands.turn.run_codex_cli_host", forbidden_plain_cli) + selected_route = {"host_surface": "codex-app", "thread_id": "source-thread"} + observed_routes = [] + + def operation_host(*args, **kwargs): + observed_routes.append(kwargs["source_route"]) + return run_codex_operation_host(*args, **kwargs) + + monkeypatch.setattr("loopx.control_plane.turn_driver.codex_operation_host.run_codex_operation_host", operation_host) arguments = [ "--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", "turn", "run-once", "--goal-id", "loopx-turn-fixture", "--agent-id", "codex-fixture", "--host", "codex-cli", "--project", str(project), "--scan-root", str(project), "--no-global-sync", "--codex-operation-tools", "--codex-bin", str(executable), + "--codex-operation-source-route-json", json.dumps(selected_route), "--codex-model", "test-model", "--codex-reasoning-effort", "xhigh", "--codex-sandbox", "read-only", "--timeout-seconds", "5", "--validation-command-json", json.dumps([sys.executable, "-c", "import json,sys; json.load(sys.stdin)"]), @@ -161,6 +348,22 @@ def forbidden_plain_cli(*args, **kwargs): assert binding["operation_transport"] == "app-server-operation-tools-v0" assert binding["session_id"] == "owned-app-server-thread" assert payload["effects"]["quota_spent"] is False + assert observed_routes == [selected_route] + + +def test_source_route_cli_requires_owned_operation_transport(tmp_path, capsys): + from loopx.cli import main as cli_main + + project, runtime, registry = _write_live_fixture(tmp_path) + rc = cli_main([ + "--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "turn", "run-once", "--goal-id", "loopx-turn-fixture", "--agent-id", "codex-fixture", + "--host", "codex-cli", "--project", str(project), "--scan-root", str(project), + "--no-global-sync", "--codex-operation-source-route-json", + json.dumps({"host_surface": "codex-app", "thread_id": "source-thread"}), + ]) + assert rc != 0 + assert "requires --codex-operation-tools" in capsys.readouterr().out @pytest.mark.parametrize(