diff --git a/apps/presentation/dashboard/smoke/action-review-plan-smoke.ts b/apps/presentation/dashboard/smoke/action-review-plan-smoke.ts index fecf558d32..6e6164426f 100644 --- a/apps/presentation/dashboard/smoke/action-review-plan-smoke.ts +++ b/apps/presentation/dashboard/smoke/action-review-plan-smoke.ts @@ -146,6 +146,12 @@ check(managedFrame?.kind === "pending" && managedFrame.executionState === "manag check(managedFrame?.content.fields.some(field => field.value.includes("test-model@xhigh")) === true, "The shared managed profile survives the frontend schema transport"); check(compileActionReviewPlan(managedPending).canApply === false, "Managed approval exposes no local execute control"); +const startedManaged = typedActionProposalSchema.parse({...managedPending, + operation: {...managedPending.operation, host_start: {schema_version: "loopx_operation_host_start_v0", host_turn_id: "native-turn"}}}); +const startedFrame = compileActionReviewPlan(startedManaged).operationFrame; +check(startedFrame?.kind === "pending" && startedFrame.executionState === "managed_turn_started", + "Native accepted start survives schema transport without claiming consumption or execution"); +check(compileActionReviewPlan(startedManaged).canApply === false, "Accepted native start grants no apply control"); check(compileActionReviewPlan(unknownAgentResult).interaction === "repair", "Delivered unknown submission is not completion"); const reconciledAgentResult = typedActionProposalSchema.parse({...unknownAgentResult, operation: {...unknownAgentResult.operation, reconciliation: {outcome: "not_executed", simulation: false}}}); diff --git a/apps/presentation/dashboard/src/features/personal-workspace/channel-timeline.tsx b/apps/presentation/dashboard/src/features/personal-workspace/channel-timeline.tsx index 19a9a59bf1..7c86ea7251 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/channel-timeline.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/channel-timeline.tsx @@ -121,10 +121,12 @@ export function ChannelTimeline({ } if (item.kind === "proposal") { const appliedTeamPlan = item.proposal.actionKind === "team.plan" && item.proposal.status === "applied"; + const pendingOperation = item.proposal.reviewPlan?.operationFrame?.kind === "pending"; return ( {showManagerTeamResults && onOpenGoalEvidence && appliedTeamPlan diff --git a/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx b/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx index a0bfc08c54..43c4822f3b 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx @@ -1102,7 +1102,9 @@ export function ContextDrawer({ agents, attentionHistory = [], onSelectAttention {selection.kind === "proposal" ? ( <> {selection.item.actionKind === "team.plan" && selection.item.status === "applied" ? :
- {selection.item.actionKind} · {selection.item.status} + {selection.item.reviewPlan?.operationFrame?.kind === "pending" + ? t(`proposal.kind.${selection.item.actionKind}`) + : `${selection.item.actionKind} · ${selection.item.status}`}

{selection.item.title}

{selection.item.impact ?

{selection.item.impact}

: null} {selection.item.reviewPlan && !selection.item.reviewPlan.retryOriginal && selection.item.actionKind !== "team.plan" ?

{operationUnknown diff --git a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx index 7af8eb78bb..9f1b9eff5b 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx @@ -684,6 +684,7 @@ const en = { "proposal.field.operationState": "Operation state", "proposal.operationState.host_authentication_required": "Confirmed; original-host authentication unavailable", "proposal.operationState.managed_turn_pending": "Confirmed; waiting for the bound managed Turn", + "proposal.operationState.managed_turn_started": "Native continuation accepted; authorization not yet consumed", "proposal.operationState.consumed_outcome_pending": "Authorization consumed; waiting for the real result", "proposal.operationState.submission_unknown": "Result unknown; reconcile the original operation, do not resubmit", "proposal.field.resultDelivery": "Result delivery", @@ -730,6 +731,7 @@ const en = { "proposal.impact.operation": "The exact terms are read-only here. Confirm or reject the same immutable request in the bound Feishu group; confirmation consumes one canonical claim.", "proposal.impact.operationAuthorized": "Human confirmation is recorded, but the original host's authenticated tool transport is not connected. Thread flags or environment ids cannot authorize execution. No external result is recorded yet.", "proposal.impact.operationManagedPending": "Confirmation is bound to the selected managed session and execution profile, not the source conversation. The admitted Turn must consume it once through its owned tool connection. No execution or external result is established yet.", + "proposal.impact.operationManagedStarted": "The bound native host accepted a continuation containing this confirmed operation. This first-start receipt is not proof that the Turn is still running, that authorization was consumed, or that an external effect occurred.", "proposal.impact.operationConsumed": "The authorization has been consumed. Wait for original external evidence; a retry or lost response must not grant another submission.", "proposal.impact.operationUnknown": "Submission may have had an external effect. Reconcile the original operation using its evidence; do not resubmit or treat card delivery as execution completion.", "proposal.primary.apply": "Confirm and apply", @@ -1882,6 +1884,7 @@ const zhCN: Record = { "proposal.field.operationState": "操作状态", "proposal.operationState.host_authentication_required": "已确认,原宿主身份认证尚未接通", "proposal.operationState.managed_turn_pending": "已确认,等待绑定的受管回合", + "proposal.operationState.managed_turn_started": "原生续接已接受,授权尚未消费", "proposal.operationState.consumed_outcome_pending": "授权已消费,等待真实结果", "proposal.operationState.submission_unknown": "结果未知;核对原操作,不可重复提交", "proposal.field.resultDelivery": "结果回传", @@ -1928,6 +1931,7 @@ const zhCN: Record = { "proposal.impact.operation": "这里仅展示同一份不可变条款。请在已绑定的飞书群确认或拒绝;确认只会消费一个规范 claim。", "proposal.impact.operationAuthorized": "用户确认已记录,但原宿主的认证工具通道尚未接通。线程参数或环境变量不能授权执行;目前尚无外部执行结果。", "proposal.impact.operationManagedPending": "批准绑定选定的受管会话和执行配置,不绑定来源对话。通过准入的回合须在自有工具通道上消费一次;目前不代表已执行,也没有外部结果。", + "proposal.impact.operationManagedStarted": "绑定的原生宿主已接受包含本确认请求的续接。首次启动回执不证明回合仍在运行、授权已消费或已有外部效果。", "proposal.impact.operationConsumed": "授权已消费。等待原始外部证据;重试或响应丢失均不得重新授予提交许可。", "proposal.impact.operationUnknown": "提交可能已产生外部副作用。须以原始证据核对原操作,不可重提,也不能把卡片投递当作执行完成。", "proposal.primary.apply": "确认并应用", diff --git a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx index ae53146737..620d23995c 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx @@ -660,7 +660,8 @@ function workspaceProposal(proposal: TypedActionProposal, t: WorkspaceTranslate) ? operationFrame?.kind === "pending" && operationFrame.executionState ? t(operationFrame.executionState === "consumed_outcome_pending" ? "proposal.impact.operationConsumed" : operationFrame.executionState === "managed_turn_pending" - ? "proposal.impact.operationManagedPending" : "proposal.impact.operationAuthorized") + ? "proposal.impact.operationManagedPending" : operationFrame.executionState === "managed_turn_started" + ? "proposal.impact.operationManagedStarted" : "proposal.impact.operationAuthorized") : operationFrame?.kind === "result" && operationFrame.resultKind === "unknown" ? t("proposal.impact.operationUnknown") : t("proposal.impact.operation") : proposal.action_kind === "team.plan" diff --git a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md index 7dbe0e6f27..95cb10a56b 100644 --- a/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md +++ b/docs/architecture/rfcs/human-confirmed-domain-operations-v0.md @@ -418,6 +418,11 @@ The existing Inbox projects canonical locators directly, with recovery first, `loopx_operation pending` accepts its bound cursor; the CLI projection uses `manager-inbox read --operation-cursor CURSOR`. New/changed work restarts without a cursor. Reading or exhausting a page does not resolve obligations. +The owned native tool applies the same exact execution-subject filter as Turn +startup before pagination. Current approvals for another Todo/session/profile +of the same Agent cannot appear on that task's page or starve its continuations; +its cursor cannot be reused by another execution subject. The existing TS inbox +owner makes these decisions; this is not a separate host-owned inbox. The shared TypeScript operation frame shows executor, pinned model/effort, Goal/Agent/Todo scope, optional source context, and distinct states: confirmed @@ -445,10 +450,44 @@ remain fail-closed before private reads or writes, even with matching `CODEX_THREAD_ID`, route flags, self-signed proof or an older runtime's actor success. No proof-import shortcut is exposed. -Immediate confirmation-triggered host wakeup is not implemented: -`host_delivery: "not_attempted"` remains truthful. Continue through the -existing admitted Turn/delegation route; later durable wakeup must reuse its -original scheduling/session owner, not start a parallel resumed executor. +An admitted operation-enabled Turn now automatically includes canonical +confirmed-operation locators for its exact Goal/Agent/Todo/session/profile, +filtered before inbox pagination. After the native `turn/start` response is +accepted, the existing action store records immutable first-start evidence: +confirmation event/time, claim, LoopX Turn key, native Turn and acceptance time. +`host_delivery: "native_start_accepted"` means that observation only; it grants +no consumption or effect authority and does not claim the Turn is still live. +Receipt failure aborts before operation-tool dispatch; later Turns preserve the +first observation. CLI/Inbox, Dashboard and Lark distinguish accepted native +continuation from consumed authorization and an actual outcome. + +An authenticated confirmation callback may now request one continuation through +the existing delegation owner, **only with a separate, default-off operator +launch grant**. The grant selects one existing requester/binding, not a model, +workspace or executor supplied by the card. The operation owner checks exact +Goal/Agent/Todo/session/profile, confirmation lifetime and unconsumed state. +The ordinary Turn still owns quota, lease, validation and native startup; its +native adapter rechecks the complete effective profile before resuming the +original session. A new/replacement session is refused. Confirmation does not +grant launch configuration or domain execution authority. + +The canonical operation id determines one durable delegation identity. Callback +replays read its locator without another spawn or artifact validation. A lost +launch acknowledgement remains an original-journal recovery case; callbacks +do not auto-resume an uncertain worker. Removing the operator grant is read at +the next callback boundary. Already started work is not retroactively cancelled. +`delegation_requested` does not prove native start, validated completion or +source delivery; `host_delivery` remains `"not_attempted"` until native acceptance. +See [operator activation](../../reference/local-delegation.md#confirmed-operation-callback-continuation). + +File/SQLite qualification uses authenticated synthetic confirmation fixtures, +the real detached delegation worker and CLI/Turn path, and a synthetic native +transport. It proves original-session startup without consumption or a domain +effect; deliberate waiting is not accepted Todo completion. Genuine Lark/model +and financial probes have not been run for this increment. Frontend grant +editing and authenticated result return to the original source audience remain +separate, **partial-delivery** obligations; current grant activation is through +the operator-owned collector configuration, not a new UI authority. 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. 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 4ff7719db4..a0b40f72da 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 @@ -328,6 +328,10 @@ executor revision、consumption ID、核验投影、`simulation: false`、 及独立 operation cursor。`loopx_operation pending` 接受绑定游标;CLI 投影用 `manager-inbox read --operation-cursor CURSOR`。新增/改变工作应无游标重读, 读完一页或遍历结束不代表义务已解决。 +自有原生工具与 Turn 启动使用相同的精确执行主体过滤,并在分页前生效。 +同一 Agent 另一 Todo/session/profile 的当前批准不能出现在本任务页面或挤占其 +续接信息;游标也不能跨执行主体复用。判定仍由既有 TS Inbox owner 负责, +不新建宿主自有 Inbox。 共享 TS 操作 frame 展示执行者、固定模型/思考深度、Goal/Agent/Todo 范围、可选来源 上下文,并区分:已确认但外接认证不可用、已确认待绑定受管回合、已消费待证据、未知须 @@ -348,9 +352,33 @@ context 调用,结果均由既有 typed result validator 接受)、规范 私有读写前拒绝,即使 `CODEX_THREAD_ID`、路由、自签 proof 完全匹配或旧运行时 意外返回 actor 成功,也没有 proof-import 捷径。 -本切片不实现确认后的即时宿主唤醒,`host_delivery: "not_attempted"` 保持真实。 -沿既有已准入 Turn/delegation 续接;后续持久唤醒复用其调度/session owner, -不启动平行 resumed 执行者。宣称投研最小闭环前,仍须证明安装、真实用户批准、 +启用 operation 的已准入 Turn 自动携带其精确 Goal/Agent/Todo/session/profile 的 +规范确认请求定位信息,并在 Inbox 分页前过滤范围。仅在原生 `turn/start` 返回接受后, +既有 action store 才记录不可覆盖的首次启动证据:确认事件/时间、claim、LoopX Turn key、 +原生 Turn 和接受时间。`host_delivery: "native_start_accepted"` 只证明该观察, +不授予消费或外部效果权限,也不声称回合仍在运行。回执失败时先停止、不分发 operation 工具; +后续回合保留首次观察。CLI/Inbox、Dashboard 和 Lark 区分原生续接接受、授权消费和真实结果。 + +认证确认 callback 现在可通过既有 delegation owner 请求一次续跑,前提是具有 +**独立且默认关闭的 operator 启动 grant**。grant 指定已有 requester/binding, +不采用卡片提供的模型、工作树或执行者。操作 owner 核对精确 +Goal/Agent/Todo/session/profile、确认期限和未消费状态;普通 Turn 仍负责 quota、 +租约、验收及原生启动。原生适配器在恢复前重验完整生效配置,拒绝新建/替换会话。 +用户确认不授予启动配置或垂域执行权限。 + +规范 operation id 固定唯一持久 delegation 身份。callback 重放只读原定位信息, +不再次 spawn 或运行产物验收;丢失启动 ACK 仍由原 journal 恢复,不在回调自动续跑 +不确定的 worker。移除 operator grant 会在下一次 callback 边界读回生效, +不追溯取消已经开始的工作。`delegation_requested` 不证明原生启动、完成验收或 +来源送达;原生接受前 `host_delivery: "not_attempted"` 保持真实。 +启用方式见[原配置入口](../../reference/local-delegation.md#confirmed-operation-callback-continuation)。 + +File/SQLite 资格化使用合成认证确认夹具、真实 detached delegation worker 与 +CLI/Turn 路径,以及合成原生传输;证明原会话启动,但未消费批准或执行垂域操作。 +夹具故意等待,不冒充 Todo 完成。本增量未跑真实 Lark/模型或金融探针。 +前端 grant 编辑及经认证的原来源受众结果回传仍为独立的**部分交付**义务; +目前只能通过 operator-owned collector 配置启用,不另建 UI 状态权威。 +宣称投研最小闭环前,仍须证明安装、真实用户批准、 绑定原生消费、垂域提交前检查与原系统证据、结果验收及原卡/受众读回。 Core PR 仍须 owner review,不在合并前自行安装。 diff --git a/docs/reference/local-delegation.md b/docs/reference/local-delegation.md index 52b8ec5b6a..0eeba16163 100644 --- a/docs/reference/local-delegation.md +++ b/docs/reference/local-delegation.md @@ -110,15 +110,15 @@ workspace、Codex Session 与持久 delegation operation。`delegation inspect` 规划投影会读回例如 `gpt-5.6-sol@xhigh`,start/resume 也把同一配置送入原 Session。 这不扩大 requester grant,也不把 profile 冒充验收或结果返回回执。 -For a Codex binding, the delegation host also supplies one invocation-scoped +By default, a Codex binding also supplies one invocation-scoped `loopx_delegation` stdio MCP server to every fresh or resumed worker Session. Its command pins the selected worker `agent_id`, workspace, Goal, registry, runtime and operator execution configuration before Codex starts. The model cannot select or rewrite those values. Codex receives the server through per-invocation configuration, so LoopX does not modify the user's global Codex -MCP settings and a resumed worker keeps the same binding. The server is required -for this managed worker route and its already identity-scoped tools are approved -inside that route; failure to start the server rejects the Turn instead of +MCP settings and a resumed worker keeps the same binding. When configured, the +server and its already identity-scoped tools are approved inside that route; +failure to start the server rejects the Turn instead of silently continuing without tools. The native tools are the collaboration and authorized delegation operations from that bound server; they do not add shell, Todo or external-action authority. @@ -129,7 +129,7 @@ shell-only coordination. Both surfaces call the same `Delegations` service and preserve the same binding, operation and acceptance rules; the shell path is not a second control-plane implementation. -中文:Codex binding 会为每个新建或续接的 worker Session 注入一次调用范围内的 +中文:Codex binding 默认会为每个新建或续接的 worker Session 注入一次调用范围内的 `loopx_delegation` stdio MCP server。启动前,host 已固定 worker `agent_id`、 workspace、Goal、registry、runtime 与 operator execution configuration;模型不能 选择或改写这些值。该配置不会修改用户的全局 Codex MCP 设置,也不会授予 shell、 @@ -694,3 +694,119 @@ member counts, model choices, workspaces and requester grants in the Goal's execution configuration; reuse machine authentication without copying another Goal's assignments. A configured credential does not prove the selected SDK, model, remote environment or task acceptance is ready. + +## Confirmed operation callback continuation + +This default-off adapter lets the existing authenticated Lark operation callback +request a continuation through one existing delegation binding. Human operation +confirmation and the operator's launch grant are separate authorities. In the +original collector v1 configuration, explicitly set: + +```json +{ + "operation_callbacks": { + "enabled": true, + "managed_turn_wake": { + "registry_path": "/absolute/operator-registry.json", + "goal_id": "project-goal", + "requester_agent_id": "coordinator", + "execution_config": ".loopx/config/delegations.json", + "binding_id": "confirmed-operation" + } + } +} +``` + +This is a fragment, not a complete collector configuration. `execution_config` +is project-relative, without traversal or symlinks, and remains outside the +member workspace. The selected binding must grant that registered requester +the operation's exact Agent/Todo. Use `codex-cli`, `--codex-operation-tools` +and the original pinned model/effort. Do not select `fresh`, change the owning +home, binary, sandbox or effective MCP configuration, or retarget the Todo. +Prepare the operation in this same managed session/profile. Existing delegation +inspection and ordinary Turn admission still determine whether it can run. + +If the original standalone operation Session was prepared **without** a +delegation MCP server, its profile includes `mcp_server: null`. Preserve that +exact profile using the existing operator-owned binding option: + +```json +{ + "host_args": [ + "--host", "codex-cli", + "--codex-operation-tools", + "--codex-model", "original-model", + "--codex-reasoning-effort", "original-effort", + "--codex-sandbox", "read-only", + "--codex-mcp-server-json", "null" + ] +} +``` + +This is a fragment; replace the model/effort placeholders with their original +values and retain the original executable, workspace, home and Agent/Todo +grant as well. Binding `host_args` follow the default injected MCP +option, so the real CLI parser uses the explicit JSON `null`. Do not infer +absence by reading the first occurrence in argv. Omit this override for an +original delegation-MCP profile: removing that server would also be drift. +Inspect the original profile before enabling the collector; runtime +qualification requires same-session/profile native acceptance, not just this +configuration readback. No global Codex configuration is changed. +The native `loopx_operation` tool remains bound to the original transport; +`null` does not grant collaboration MCP tools, bypass ordinary receiver +adoption/validation, or prove consumption, outcome or source delivery. A +synthetic typed `wait` may yield native acceptance and a rejected delegation; +that is not an accepted work result or an end-to-end operation loop. + +中文:若原独立 operation Session 在没有 delegation MCP 时准备,profile 中的 +`mcp_server` 为 `null`,应在原 operator binding 的 `host_args` 显式设置 +`["--codex-mcp-server-json", "null"]`,保留 executable、workspace、home、原 +Agent/Todo grant 及模型/深度等其余字段。上例仅为片段,model/effort 占位符须替换 +为原值;binding 参数在默认注入 +参数之后,真实 CLI parser 采用后面的 JSON `null`,不能取 argv 第一个同名参数 +来判断生效配置。原 profile 使用 delegation MCP 时不要加此覆盖,删除 server +同样属于漂移。启用 collector 前核对原 profile;运行资格必须由同 session/profile +的原生接收回执证明,不能只凭配置读回。不改全局 Codex 配置。原生 +`loopx_operation` 工具保持原连接绑定,但 `null` +不会授予 collaboration MCP 工具、绕过接收方采纳/验收,也不证明消费、结果或原 +来源送达;合成 typed `wait` 可同时产生原生接收和 delegation rejected,不代表 +工作完成或端到端操作闭环。 + +`loopx lark-inbox collector-plan` and `collector-status` expose +`operation_callback_managed_wake_configured`, a configuration observation, not +runtime qualification. An authenticated callback returns `managed_turn_wake`: +`delegation_requested` gives the original operation locator; +`existing_delegation` acknowledges replay without another spawn; +`blocked` preserves the canonical confirmation but grants no launch. Read the +locator through existing delegation read/recovery commands. Callback receipts +never certify native start, consumption, Todo completion or original-audience +result delivery. The internal `--codex-confirmed-operation-id` is only an exact +resume fence, not a caller credential or domain permit. + +Remove `managed_turn_wake` or the exact delegation requester grant to stop new +admission. The collector reloads the launch grant at each callback; already +started work needs its ordinary stop/reconciliation path. Lost spawn ACKs do +not trigger automatic callback resume. Retain the original journal and inspect +it before recovery. No new scheduler, approval store or result owner is added. + +The existing Dashboard shows operation/host-start state through the shared +frame, but does **not** edit this grant. Frontend grant editing and genuine +original-source outcome delivery remain partial; synthetic startup tests do +not establish an end-to-end user or trading loop. + +中文:此适配默认关闭。认证 Lark 确认只能请求既有 delegation 续跑,不能直接 +执行金融操作;启动 grant 与用户对不可变条款的批准独立。上例是原 collector v1 +配置片段,不是完整配置。`execution_config` 必须项目内相对路径、无穿越或符号 +链接,并位于成员工作树之外;精确 binding 必须授予已注册 requester 原 Agent/Todo。 +使用 `codex-cli`、operation tools 和原固定模型/深度;不能选 `fresh`、更换 home、 +binary、sandbox、生效 MCP 或任务。在相同受管 session/profile 准备操作, +普通 Turn 的准入、quota、租约与验收仍全部生效。 + +plan/inspect 的 `operation_callback_managed_wake_configured` 仅证明配置, +不证明运行资格。callback 的 `delegation_requested` 返回原定位信息, +`existing_delegation` 表示重放未重新 spawn,`blocked` 保留确认但不授予启动。 +通过既有 delegation 读回/恢复命令检查原记录;上述状态都不是原生启动、消费、 +完成或原受众送达回执。内部 operation-id 参数只做精确恢复保护,不是身份凭证。 +移除启动配置或原 requester grant 在下一 callback 生效,已开始工作须走原停止/ +对账流程;丢失 ACK 不自动重启。Dashboard 沿用共享 frame 展示状态,但尚不能编辑 +此 grant;前端配置与真实原来源结果送达仍为部分交付,不能据合成测试称交易闭环。 diff --git a/loopx/chat_action_store.py b/loopx/chat_action_store.py index d42b3f1bf6..b3eec55761 100644 --- a/loopx/chat_action_store.py +++ b/loopx/chat_action_store.py @@ -1068,6 +1068,29 @@ def _agent_operation_plan( except EffectRuntimeRemoteError as exc: raise ActionConflictError(str(exc)) from exc + def record_agent_operation_host_start( + self, proposal_id: str, *, actor: Mapping[str, Any], + binding_current: bool, turn_key: str, + ) -> dict[str, Any]: + """Internal native transport observation; never a caller proof or permit.""" + token = _opaque_id(proposal_id, field="proposal_id") + with exclusive_file_lock(self.path, operation="observe_operation_host_start"): + payload = self._read() + proposal = payload["proposals"].get(token) + if not isinstance(proposal, dict): + raise KeyError("typed operation was not found") + plan = self._agent_operation_plan( + proposal, action="observe_host_start", actor=dict(actor), + binding_current=binding_current, turn_key=turn_key, + ) + receipt = plan.pop("write_host_start", None) + if receipt is not None: + proposal["operation"]["host_start"] = receipt + proposal["updated_at"] = _utc_now() + self._write(payload) + plan.update(host_start=receipt, host_delivery="native_start_accepted") + return plan + def consume_agent_operation( self, proposal_id: str, diff --git a/loopx/cli_commands/turn.py b/loopx/cli_commands/turn.py index 82d2c8160e..9a0d8192c4 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_confirmed_operation_id", None) and not getattr(args, "codex_operation_tools", False): + raise ValueError("confirmed operation continuation requires the owned operation transport") 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. @@ -1005,7 +1007,9 @@ def run_built_in_host( return run_codex_operation_host( request, registry_path=registry_path, - source_route=getattr(args, "codex_operation_source_route_json", None), **options + confirmed_operation_id=getattr(args, "codex_confirmed_operation_id", None), + 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 9f1be95320..ab857990d6 100644 --- a/loopx/cli_commands/turn_registration.py +++ b/loopx/cli_commands/turn_registration.py @@ -224,6 +224,10 @@ 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-confirmed-operation-id", + help="Internal exact-operation resume fence for an operator-granted callback continuation. Requires operation tools; does not authenticate a caller or permit a domain effect.", + ) run_once.add_argument( "--codex-operation-source-route-json", type=json.loads, diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index 426c0916bc..7a4a5f0894 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -521,7 +521,8 @@ def recheck_workspace() -> dict[str, object] | None: }) def start(self, binding_id: str, operation_id: str, brief: dict, - parent_request_id: str | None = None) -> dict: + parent_request_id: str | None = None, *, + confirmed_operation_id: str | None = None) -> dict: binding = self.binding(binding_id, require_active=True) require_operation_id(operation_id) brief = normalize_request({"goal_id": self.goal_id, "agent_id": binding["agent_id"], "brief": brief})["brief"] @@ -536,6 +537,10 @@ def start(self, binding_id: str, operation_id: str, brief: dict, binding["agent_id"], operation_id, brief, parent_request_id, caller_goal_ref=self._caller_goal_ref()) identity = {"binding": binding, "request_id": delivered["request_id"], "operation_id": operation_id} + if confirmed_operation_id is not None: + # Internal callback adapter only: a canonical locator/CAS fence, + # not an executor identity or domain execution permission. + identity["confirmed_operation_id"] = require_operation_id(confirmed_operation_id) if exists: if _read(path).get("identity") != identity: raise ValueError("delegation operation identity conflict") @@ -818,12 +823,18 @@ def _execution_arguments(self, binding: dict, operation_id: str) -> list[str]: ], } native_tools = ["--codex-mcp-server-json", json.dumps(mcp_server)] + continuation: list[str] = [] + path = self.path(operation_id) + if path.is_file(): + confirmed = _read(path)["identity"].get("confirmed_operation_id") + if confirmed is not None: + continuation = ["--codex-confirmed-operation-id", require_operation_id(confirmed)] return ["--execution-mode", "isolated-headless", "--project", binding["workspace"], "--scan-root", binding["workspace"], "--no-global-sync", "--timeout-seconds", str(binding["timeout_seconds"]), "--validation-command-json", json.dumps(validator), "--validation-failure-kind", "repair_required", *native_tools, - *binding["host_args"]] + *binding["host_args"], *continuation] def _record_turn_result( self, path: Path, row: dict, result: dict, *, publish: bool = True diff --git a/loopx/control_plane/collaboration/operation_handoff.py b/loopx/control_plane/collaboration/operation_handoff.py index 70c162dd5a..2bada69a3f 100644 --- a/loopx/control_plane/collaboration/operation_handoff.py +++ b/loopx/control_plane/collaboration/operation_handoff.py @@ -79,6 +79,7 @@ def pending_operation_handoffs( scope: CollaborationGoalScope | None = None, cursor: str | None = None, cursor_scope: str, + executor_route: Mapping[str, Any] | None = None, ) -> dict[str, Any]: """Project canonical tickets into the existing Inbox; do not copy authority.""" store = ( @@ -155,6 +156,7 @@ def pending_operation_handoffs( "items": result, "cursor": cursor, "cursor_scope": cursor_scope, + "executor_route": dict(executor_route) if executor_route is not None else None, }, ) ) @@ -171,6 +173,7 @@ def agent_operation_action( action: str, consumption_id: str | None = None, outcome: Mapping[str, Any] | None = None, + turn_key: str | None = None, ) -> dict[str, Any]: """Internal locked IO seam, not a public caller-authentication endpoint. @@ -194,12 +197,12 @@ def agent_operation_action( goal_id=actor["goal_id"], agents=(actor["agent_id"],), caller_goal_ref=parameters.get("origin_goal_ref"), - require_active=action == "consume", + require_active=action in {"consume", "observe_host_start"}, # Recovery-owner binding must stay valid through evidence commit too. # Use the same Goal -> registry -> action-store ordering as consumption. lock_registry=True, ) as scope: - if action == "consume": + if action in {"consume", "observe_host_start"}: decide_collaboration_lifecycle(scope, operation="request_create") if goal_is_stopped(scope.goal): raise ActionConflictError("confirmed operation Goal is stopped") @@ -258,6 +261,7 @@ def agent_operation_action( action, consumption_id, outcome, + turn_key, ) @@ -272,6 +276,7 @@ def _commit_agent_operation( action: str, consumption_id: str | None, outcome: Mapping[str, Any] | None, + turn_key: str | None, ) -> dict[str, Any]: current = _binding(registry_path, parameters, runtime_root) actor_executor = ( @@ -313,6 +318,10 @@ def _commit_agent_operation( "outcome_report": proposal["operation"].get("outcome_report"), "reconciliation_report": proposal["operation"].get("reconciliation_report"), } + if action == "observe_host_start": + return store.record_agent_operation_host_start( + proposal_id, actor=actor, binding_current=current, turn_key=str(turn_key or ""), + ) if action == "consume": return store.consume_agent_operation( proposal_id, diff --git a/loopx/control_plane/collaboration/operation_wake.py b/loopx/control_plane/collaboration/operation_wake.py new file mode 100644 index 0000000000..26d7fa9137 --- /dev/null +++ b/loopx/control_plane/collaboration/operation_wake.py @@ -0,0 +1,112 @@ +"""Callback IO into the existing operator-granted delegation/Turn owner. + +The callback confirms terms; configuration separately grants this one launch. +No new queue, scheduler, approval, result store or public identity issuer. +""" +from __future__ import annotations + +import hashlib +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from ...chat_action_store import ChatActionStore +from ..turn_driver.codex_cli import load_codex_cli_session +from ..turn_driver.host_binding import turn_host_arg_option +from .operation_handoff import managed_operation_binding_current +from .goal_instance_scope import collaboration_goal_scope, decide_collaboration_lifecycle + + +def require_current_operation_scope(registry_path: Path, parameters: Mapping[str, Any]) -> None: + """Reuse the Goal lifetime owner before either queueing or native resume.""" + with collaboration_goal_scope(registry_path, goal_id=parameters["goal_id"], + agents=(parameters["agent_id"],), require_active=True) as scope: + decision = decide_collaboration_lifecycle( + scope, operation="inbox_observe", record={"goal_ref": parameters.get("origin_goal_ref")} + ) + if decision.get("kind") == "omit": + raise ValueError("operation belongs to a different Goal lifetime") + + +def dispatch_confirmed_operation_wake( + proposal: Mapping[str, Any], *, runtime_root: Path, + configuration: Mapping[str, Any] | None, +) -> dict[str, Any]: + """Request one durable delegation; accepted native start is separate evidence. + + Read/compare failures remain a pending canonical authorization. Never resume + an uncertain launch from a callback replay or include private argv/errors. + """ + base = {"execution_allowed": False, "external_write_performed": False, + "native_start_verified": False} + if configuration is None: + return {**base, "state": "not_configured"} + try: + from ...collaboration_mcp import Delegations + from .delegation_context import _configuration_path + + parameters = proposal["normalized_parameters"] + if parameters["goal_id"] != configuration["goal_id"]: + return {**base, "state": "out_of_scope"} + service = Delegations( + runtime_root, Path(configuration["registry_path"]), + configuration["goal_id"], configuration["requester_agent_id"], + _configuration_path(Path(configuration["project"]), configuration["execution_config"]), + ) + binding = service.binding(configuration["binding_id"], require_active=True) + require_current_operation_scope(service.registry, parameters) + if service.config.is_relative_to(Path(binding["workspace"]).resolve()): + raise ValueError("operation wake configuration must be outside the worker workspace") + operation_id = "operation-wake-" + hashlib.sha256( + str(proposal["proposal_id"]).encode() + ).hexdigest() + if service.path(operation_id).is_file(): + # Start's lost ACK is resolved by readback, never another spawn. + # Do not run artifact validators on the callback's latency path or + # describe a stored accepted bit as current acceptance. CLI read + # owns result qualification/recovery after this bounded locator. + from .inbox import _read + + row = _read(service.path(operation_id)) + service._bound(row, require_active=True) + if row["identity"].get("confirmed_operation_id") != proposal["proposal_id"]: + raise ValueError("operation wake identity drifted") + return {**base, "state": "existing_delegation", "operation_id": operation_id} + projection = ChatActionStore._agent_operation_plan(proposal, action="project") + if projection["status"] != "authorized_pending" or projection["host_start"] is not None: + return {**base, "state": "no_new_launch"} + session = load_codex_cli_session(runtime_root, lineage={ + "goal_id": service.goal_id, "agent_id": binding["agent_id"], + "todo_id": binding["todo_id"], + }) or {} + argv = binding["host_args"] + ChatActionStore._agent_operation_plan( + proposal, action="wake", + binding_current=managed_operation_binding_current(runtime_root, parameters), + launch_context={"host": turn_host_arg_option(argv, "--host"), + "operation_tools": "--codex-operation-tools" in argv, + "iteration_context": turn_host_arg_option(argv, "--iteration-context")}, + executor_route={"goal_id": service.goal_id, "agent_id": binding["agent_id"], + "todo_id": binding["todo_id"], "host_surface": "loopx-managed-codex", + "thread_id": session.get("session_id"), + "profile_digest": session.get("operation_profile_digest"), + "model": turn_host_arg_option(argv, "--codex-model"), + "reasoning_effort": turn_host_arg_option(argv, "--codex-reasoning-effort")}, + ) + brief = { + "schema_version": "collaboration_brief_v0", + "purpose": "Resume the exact human-confirmed canonical operation", + "context": f"Canonical operation locator: {proposal['proposal_id']}. Read through loopx_operation in the original managed session.", + "constraints": ["Inspect current canonical terms and consume once before any effect.", + "Consumed or unknown results require reconciliation, never resubmission.", + "The source conversation is not executor authentication."], + "inputs": [], "acceptance": ["Preserve the pinned Todo validation and original evidence."], + "return_requirement": "Report the exact operation outcome through loopx_operation; startup/adoption is not a domain result.", + } + result = service.start(configuration["binding_id"], operation_id, brief, + confirmed_operation_id=proposal["proposal_id"]) + return {**base, "state": "delegation_requested", "operation_id": operation_id, + "delegation_status": result["status"], + "recovery_required": result["recovery_required"]} + except Exception: # authorization remains recorded even if launch/readback fails + return {**base, "state": "blocked", "reason_code": "operation_wake_admission_or_dispatch_failed"} diff --git a/loopx/control_plane/presentation/action_review_plan.ts b/loopx/control_plane/presentation/action_review_plan.ts index 2aec54d11b..c13468f922 100644 --- a/loopx/control_plane/presentation/action_review_plan.ts +++ b/loopx/control_plane/presentation/action_review_plan.ts @@ -42,7 +42,7 @@ export type OperationReviewFrame = OperationReviewFrameBase & ( kind: "pending"; attentionKind: "progress"; interactionMode: "inform"; - executionState?: "host_authentication_required" | "managed_turn_pending" | "consumed_outcome_pending"; + executionState?: "host_authentication_required" | "managed_turn_pending" | "managed_turn_started" | "consumed_outcome_pending"; } | { kind: "result"; @@ -332,7 +332,10 @@ export function compileOperationReviewFrame(proposalValue: unknown): OperationRe interactionMode: "inform", ...(["agent_session", "managed_turn"].includes(String(objectValue(parameters.executor)?.kind)) ? {executionState: objectValue(operation.agent_handoff) ? "consumed_outcome_pending" as const - : objectValue(parameters.executor)?.kind === "managed_turn" ? "managed_turn_pending" as const : "host_authentication_required" as const} + : objectValue(parameters.executor)?.kind === "managed_turn" + ? objectValue(operation.host_start)?.schema_version === "loopx_operation_host_start_v0" + ? "managed_turn_started" as const : "managed_turn_pending" as const + : "host_authentication_required" as const} : {}), }; } diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index b65abbe045..1643a2536c 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -95,6 +95,12 @@ def operation_tool_handler( "model": model, "reasoning_effort": reasoning_effort, } + executor_route = { + **lineage, + "host_surface": "loopx-managed-codex", + "thread_id": session_id, + "profile_digest": profile_digest, + } def handle(tool: str, arguments: Any, native: dict[str, Any]) -> dict[str, Any]: if tool != "loopx_operation" or not isinstance(arguments, dict): @@ -151,6 +157,7 @@ def handle(tool: str, arguments: Any, native: dict[str, Any]) -> dict[str, Any]: scope=scope, cursor=arguments.get("cursor"), cursor_scope=cursor_scope, + executor_route=executor_route, ), } if action == "prepare": @@ -234,6 +241,7 @@ def run_codex_operation_host( mcp_server: Mapping[str, Any] | None = None, timeout_seconds: float = 115, goal_admission: FirstPartyHostGoalAdmission | None = None, + confirmed_operation_id: str | None = None, ) -> dict[str, Any]: if request.get("schema_version") != LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION: raise ValueError("unsupported LoopX Turn host request schema") @@ -293,6 +301,29 @@ def run_codex_operation_host( raise ValueError( "managed operation profile changed; explicitly select a fresh iteration and obtain fresh approval" ) + if confirmed_operation_id is not None: + # Final admission recheck before native resume. A queued delegation must + # not create/rebind a Session if approval, Todo, profile or lifetime + # changed after the callback. This grants no tool/effect authority. + from ..collaboration.operation_handoff import managed_operation_binding_current + from ..collaboration.operation_wake import require_current_operation_scope + + store = ChatActionStore(runtime_root / "chat" / "actions") + proposal = store.load(confirmed_operation_id) + if proposal is None or binding is None or action != "resume": + raise ValueError("confirmed operation requires its original resumable session") + require_current_operation_scope(registry_path, proposal["normalized_parameters"]) + store._agent_operation_plan( + proposal, action="wake", + binding_current=managed_operation_binding_current( + runtime_root, proposal["normalized_parameters"] + ), + launch_context={"host": "codex-cli", "operation_tools": True, + "iteration_context": (session_plan.get("context_policy") or {}).get("mode")}, + executor_route={**lineage, "host_surface": "loopx-managed-codex", + "thread_id": binding["session_id"], "profile_digest": profile_digest, + "model": model, "reasoning_effort": reasoning_effort}, + ) host_config = ( { "mcp_servers": { @@ -356,6 +387,47 @@ def store_binding(): source_route=source_route, goal_admission=goal_admission, ) + route = {**lineage, "host_surface": "loopx-managed-codex", + "thread_id": session.thread_id, "profile_digest": profile_digest} + continuations = [] + if (runtime_root / "chat" / "actions" / "actions.json").is_file(): + with collaboration_goal_scope( + registry_path, goal_id=lineage["goal_id"], agents=(lineage["agent_id"],), + caller_goal_ref=request.get("goal_ref"), require_active=True, + ) as scope: + pending = pending_operation_handoffs( + runtime_root, lineage["goal_id"], lineage["agent_id"], + registry_path=registry_path, scope=scope, executor_route=route, + cursor_scope=hashlib.sha256(json.dumps(route, sort_keys=True).encode()).hexdigest(), + ) + # Bounded canonical locators only. Private terms still require the + # existing authenticated inspect tool; these locators grant nothing. + continuations = [ + {key: item[key] for key in ("operation_id", "payload_digest", "confirmation_digest", "claim_id")} + for item in pending["items"] if item["status"] == "authorized_pending" + ] + + def on_event(kind: str, event: dict[str, Any]) -> None: + if kind != "turn.started": + return + # send emits this only after the native turn/start response. A + # launched process, callback ACK or model assertion cannot issue it. + for locator in continuations: + try: + agent_operation_action( + runtime_root, registry_path, proposal_id=locator["operation_id"], + actor={**route, "model": model, "reasoning_effort": reasoning_effort, + "host_turn_id": event.get("upstream_turn_id")}, + action="observe_host_start", turn_key=request["turn_key"], + ) + except (ValueError, KeyError, TypeError, RuntimeError) as exc: + # Stop before dispatching operation tools if confirmation, + # Goal or binding changed while native start was in flight. + raise BuiltInHostError( + "codex_operation_host_start_unrecorded", failure_kind="unknown", + recovery_kind="resume_session", + ) from exc + return session.send( _prompt(request) + "\nUse loopx_operation for context/pending/prepare/inspect/consume/report. " @@ -365,8 +437,13 @@ def store_binding(): "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.", + "Already consumed/unknown effects require evidence reconciliation, never retry. Never treat final-answer prose as an outcome receipt." + + ("\nCanonical confirmed operation continuations for this exact binding: " + + json.dumps(continuations, sort_keys=True) + + ". Inspect these locators through loopx_operation and consume once before any effect." + if continuations else ""), output_schema=codex_cli_result_schema(request), + on_event=on_event, ) except CodexChatAgentError as exc: raise BuiltInHostError( diff --git a/loopx/control_plane/work_items/operation_agent_handoff.ts b/loopx/control_plane/work_items/operation_agent_handoff.ts index 3b93c87da2..aae367e96e 100644 --- a/loopx/control_plane/work_items/operation_agent_handoff.ts +++ b/loopx/control_plane/work_items/operation_agent_handoff.ts @@ -2,7 +2,7 @@ * second approval store. Python supplies locked storage and registry facts; * this owner decides admission, one-shot consumption and result binding. */ import type {JsonObject} from "../effect_program.ts"; -import {BARE_SHA256_PATTERN} from "../content_digest.ts"; +import {BARE_SHA256_PATTERN, ENVELOPED_SHA256_PATTERN} from "../content_digest.ts"; import {EffectRuntimeConflictError, EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; import {requireJsonObject, requireNonEmptyString} from "../runtime_decode.ts"; @@ -175,6 +175,7 @@ export function planAgentOperationHandoff(input: JsonObject): JsonObject { const confirmation = operation.confirmation == null ? null : requireJsonObject(operation.confirmation, "operation confirmation"); const claim = operation.claim == null ? null : requireJsonObject(operation.claim, "operation claim"); + const hostStart = operation.host_start == null ? null : requireJsonObject(operation.host_start, "native host start"); const now = timestamp(input.now); const expires = timestamp(operation.expires_at); const managed = executor.kind === "managed_turn"; @@ -186,7 +187,8 @@ export function planAgentOperationHandoff(input: JsonObject): JsonObject { payload_digest: operation.payload_digest, confirmation_digest: operation.confirmation_digest, claim_id: claim?.claim_id ?? null, executor_revision: executor.revision, expires_at: operation.expires_at, route, authorization_source: "canonical_typed_operation", execution_allowed: false, - host_delivery: "not_attempted", external_write_performed: false, + host_delivery: hostStart ? "native_start_accepted" : "not_attempted", host_start: hostStart, + external_write_performed: false, executor_kind: executor.kind, source_route: parameters.source_route ?? null, host_authentication_required: !managed}; const handoff = operation.agent_handoff == null ? null @@ -206,6 +208,52 @@ export function planAgentOperationHandoff(input: JsonObject): JsonObject { requireThat(confirmation?.decision === "confirm" && confirmation.confirmation_digest === operation.confirmation_digest && claim, "agent execution requires authenticated confirmation"); + if (action === "wake") { + // A launch fence, never authentication or first-consumption authority. + // The existing delegation owner supplies its operator grant; the native + // host rechecks the complete effective profile immediately before resume. + const launch = requireJsonObject(input.launch_context, "operation wake launch context"); + const selected = requireJsonObject(input.executor_route, "operation wake executor route"); + requireThat(managed && input.binding_current === true + && launch.host === "codex-cli" && launch.operation_tools === true + && launch.iteration_context !== "fresh" + && Object.entries(route).every(([key, value]) => selected[key] === value) + && selected.model === executor.model && selected.reasoning_effort === executor.reasoning_effort, + "operation wake must resume the original managed session and profile"); + requireThat(!hostStart && !handoff && !observed + && operation.lifecycle_state === "claimed" && proposal.status === "applying", + "operation wake requires an unstarted, unconsumed confirmed operation"); + requireThat(now < expires && timestamp(confirmation.confirmed_at) <= now, + "operation wake is outside the confirmation lifetime"); + return {...base, status: "wake_admitted", wake_allowed: true}; + } + if (action === "observe_host_start") { + const actor = requireJsonObject(input.actor, "native start actor"); + requireThat(managed && Object.entries(route).every(([key, value]) => actor[key] === value) + && actor.model === executor.model && actor.reasoning_effort === executor.reasoning_effort + && input.binding_current === true, "native start is not the original managed binding"); + const hostTurnId = id(actor.host_turn_id, "native host Turn"); + const turnKey = requireNonEmptyString(input.turn_key, "LoopX Turn key"); + requireThat(ENVELOPED_SHA256_PATTERN.test(turnKey), "native start requires a bound LoopX Turn key"); + // This is first-start evidence, not a launch lock or execution permit. + // A retry must preserve the original time and causal identity, never relabel + // a later scheduled Turn as the first confirmation-triggered continuation. + if (hostStart) return {...base, status: "native_start_accepted", recorded: false}; + requireThat(!handoff && !observed && operation.lifecycle_state === "claimed" && proposal.status === "applying", + "native start cannot manufacture continuation evidence after consumption or outcome"); + requireThat(now < expires && timestamp(confirmation.confirmed_at) <= now, + "native continuation is outside the confirmation lifetime"); + const eventId = requireNonEmptyString(confirmation.event_id, "confirmation event"); + requireThat(eventId.length <= 512 && !/[\x00-\x1f]/.test(eventId), "confirmation event is invalid"); + const receipt: JsonObject = {schema_version: "loopx_operation_host_start_v0", + operation_id: operation.operation_id, payload_digest: operation.payload_digest, + confirmation_digest: operation.confirmation_digest, + confirmation_event_id: eventId, confirmed_at: confirmation.confirmed_at, + claim_id: id(claim.claim_id, "operation claim"), route, turn_key: turnKey, + host_turn_id: hostTurnId, accepted_at: input.now, trigger_kind: "canonical_operation_inbox", + execution_allowed: false, external_write_performed: false}; + return {...base, status: "native_start_accepted", recorded: true, write_host_start: receipt}; + } if (action === "consume") { const actor = requireJsonObject(input.actor, "execution actor"); requireThat(Object.entries(route).every(([key, value]) => actor[key] === value), @@ -264,7 +312,10 @@ export function planAgentOperationHandoff(input: JsonObject): JsonObject { * overflow locator can be inspected directly in the canonical action store. */ export function projectAgentOperationInbox(input: JsonObject): JsonObject { if (!Array.isArray(input.items)) throw new EffectRuntimeRequestError("handoff items must be an array"); - const items = input.items.map(value => requireJsonObject(value, "handoff item")); + const executorRoute = input.executor_route == null ? null : requireJsonObject(input.executor_route, "executor route"); + const items = input.items.map(value => requireJsonObject(value, "handoff item")) + .filter(item => executorRoute === null || Object.entries(executorRoute) + .every(([key, value]) => requireJsonObject(item.route, "operation route")[key] === value)); const rank = (item: JsonObject) => item.needs_reconciliation === true ? 0 : 1; items.sort((a, b) => rank(a) - rank(b) || String(a.operation_id).localeCompare(String(b.operation_id), "en")); diff --git a/loopx/extensions/lark/event_collector.py b/loopx/extensions/lark/event_collector.py index b7313e48fc..41de5dfc1a 100644 --- a/loopx/extensions/lark/event_collector.py +++ b/loopx/extensions/lark/event_collector.py @@ -139,12 +139,29 @@ def load_lark_event_collector_config( raw_operation_callbacks = ( raw_operation_callbacks if isinstance(raw_operation_callbacks, Mapping) else {} ) - unknown_operation_callback_fields = set(raw_operation_callbacks) - {"enabled"} + unknown_operation_callback_fields = set(raw_operation_callbacks) - {"enabled", "managed_turn_wake"} if unknown_operation_callback_fields: raise ValueError("collector operation_callbacks contains unsupported fields") operation_callbacks_enabled = raw_operation_callbacks.get("enabled") is True if operation_callbacks_enabled and schema_version == CONFIG_SCHEMA_VERSION_V0: raise ValueError("collector operation_callbacks requires config v1") + managed_turn_wake = None + raw_wake = raw_operation_callbacks.get("managed_turn_wake") + if raw_wake is not None: + fields = {"registry_path", "goal_id", "requester_agent_id", "execution_config", "binding_id"} + if not isinstance(raw_wake, Mapping) or set(raw_wake) != fields: + raise ValueError("collector managed_turn_wake requires one explicit operator binding") + if not operation_callbacks_enabled: + raise ValueError("collector managed_turn_wake requires enabled operation callbacks") + if any(not isinstance(raw_wake[key], str) or not raw_wake[key] + or len(raw_wake[key]) > 4096 or any(ord(c) < 32 for c in raw_wake[key]) for key in fields): + raise ValueError("collector managed_turn_wake requires bounded string fields") + registry_path = Path(raw_wake["registry_path"]).expanduser() + if not registry_path.is_absolute(): + raise ValueError("collector managed_turn_wake registry_path must be absolute") + config_ref, _ = _relative_project_path(root, raw_wake["execution_config"], "wake execution_config") + managed_turn_wake = {**dict(raw_wake), "registry_path": str(registry_path), + "project": str(root), "execution_config": config_ref} for label, value, lower, upper in ( ( "initial_lookback_seconds", @@ -325,6 +342,7 @@ def load_lark_event_collector_config( "operation_callbacks": { "enabled": operation_callbacks_enabled, "event_key": OPERATION_CALLBACK_EVENT_KEY, + "managed_turn_wake": managed_turn_wake, }, "routes": routes, } @@ -475,6 +493,7 @@ def _plan( "route_count": len(config["routes"]), "multi_chat_routing": len(config["routes"]) > 1, "operation_callbacks_enabled": config["operation_callbacks"]["enabled"], + "operation_callback_managed_wake_configured": config["operation_callbacks"]["managed_turn_wake"] is not None, "operation_callback_event_key": ( config["operation_callbacks"]["event_key"] if config["operation_callbacks"]["enabled"] @@ -725,6 +744,7 @@ def inspect_lark_event_collector( routes_with_event_evidence == len(config["routes"]) ), "operation_callbacks_enabled": callbacks_enabled, + "operation_callback_managed_wake_configured": config["operation_callbacks"]["managed_turn_wake"] is not None, "operation_callback_listener_active": callback_listener_active, "operation_callback_listener_ready": callback_listener_ready, "operation_callback_delivery_verified": callback_delivery_verified, diff --git a/loopx/extensions/lark/event_collector_runtime.py b/loopx/extensions/lark/event_collector_runtime.py index 41fc01b658..e9c9318f80 100644 --- a/loopx/extensions/lark/event_collector_runtime.py +++ b/loopx/extensions/lark/event_collector_runtime.py @@ -853,6 +853,12 @@ def consume_operation_callbacks() -> None: cli_bin=lark_cli_executable, profile=str(config["profile"]), runner=transport_runner, + # Read the original operator owner at the event + # boundary; removing a wake grant takes effect + # without restarting a long-lived collector. + managed_turn_wake=load_lark_event_collector_config( + project=project, config_path=config_path + )["operation_callbacks"]["managed_turn_wake"], ) if receipt.get("ok") is not True: raise RuntimeError( diff --git a/loopx/extensions/lark/goal_channel_operation.py b/loopx/extensions/lark/goal_channel_operation.py index 704e2ed4bb..5b40cf1129 100644 --- a/loopx/extensions/lark/goal_channel_operation.py +++ b/loopx/extensions/lark/goal_channel_operation.py @@ -382,6 +382,8 @@ def build_goal_channel_operation_result_card( result_label = "执行授权已消费,等待真实结果" elif pending and frame.get("executionState") == "managed_turn_pending": result_label = "已确认,等待绑定的受管回合;尚未执行" + elif pending and frame.get("executionState") == "managed_turn_started": + result_label = "原生续接已接受;授权仍待消费,尚无执行结果" summary = str(frame.get("summary") or result_label) return { "schema": "2.0", @@ -964,6 +966,7 @@ def handle_goal_channel_operation_callback( profile: str, runner: CommandRunner = default_subprocess_runner, executor: Callable[[Mapping[str, Any]], Mapping[str, Any]] | None = None, + managed_turn_wake: Mapping[str, Any] | None = None, ) -> dict[str, Any]: action = _callback_action(event) callback_token = str(event.get("token") or "").strip() @@ -1055,6 +1058,7 @@ def handle_goal_channel_operation_callback( }, ) dispatch_lock = store.root / f"{action['operation_id']}.dispatch.lock" + wake_receipt = None with exclusive_file_lock( dispatch_lock, agent_id="loopx-lark-operation", @@ -1071,6 +1075,11 @@ def handle_goal_channel_operation_callback( # existing Inbox. No simulator, host resume or financial effect is # run in the callback process, and no outcome is manufactured. store._agent_operation_plan(current, action="project") + from ...control_plane.collaboration.operation_wake import dispatch_confirmed_operation_wake + + wake_receipt = dispatch_confirmed_operation_wake( + current, runtime_root=runtime_root, configuration=managed_turn_wake + ) elif current_operation.get("lifecycle_state") == "claimed": outcome = dict( executor(current) @@ -1141,6 +1150,7 @@ def handle_goal_channel_operation_callback( else None ), "callback_ack_is_execution_receipt": False, + "managed_turn_wake": wake_receipt, "card_update_verified": update_verified, "status": ( "authorization_pending" diff --git a/tests/control_plane_ts/action_review_plan.test.ts b/tests/control_plane_ts/action_review_plan.test.ts index cc0653b650..662980a011 100644 --- a/tests/control_plane_ts/action_review_plan.test.ts +++ b/tests/control_plane_ts/action_review_plan.test.ts @@ -162,6 +162,11 @@ test("managed executor and source context use the same frame without turning app assert.equal(JSON.stringify(frame).includes("private-source-thread"), false); assert.equal(compileActionReviewPlan(proposal).canApply, false); assert.equal(compileActionReviewPlan(proposal).reason, "operation_authorization_pending"); + proposal.operation.host_start = {schema_version: "loopx_operation_host_start_v0", host_turn_id: "native-turn"}; + frame = compileOperationReviewFrame(proposal); + assert.equal(frame?.kind === "pending" && frame.executionState, "managed_turn_started"); + assert.equal(compileActionReviewPlan(proposal).canApply, false); + assert.equal(compileActionReviewPlan(proposal).reason, "operation_authorization_pending"); proposal.operation.agent_handoff = {consumption_id: "managed-attempt"}; frame = compileOperationReviewFrame(proposal); assert.equal(frame?.kind === "pending" && frame.executionState, "consumed_outcome_pending"); diff --git a/tests/control_plane_ts/operation_agent_handoff.test.ts b/tests/control_plane_ts/operation_agent_handoff.test.ts index a61cabefd2..a8f78e8697 100644 --- a/tests/control_plane_ts/operation_agent_handoff.test.ts +++ b/tests/control_plane_ts/operation_agent_handoff.test.ts @@ -98,6 +98,109 @@ function managedInput(): JsonObject { return value; } +function nativeStartInput(): JsonObject { + const value = managedInput(); + value.action = "observe_host_start"; + value.turn_key = "sha256:" + "d".repeat(64); + Object.assign(value.actor as JsonObject, {model: "test-model", reasoning_effort: "xhigh"}); + Object.assign(operation(value).confirmation as JsonObject, + {event_id: "authenticated-confirmation-event", confirmed_at: "2029-12-31T23:59:59Z"}); + return value; +} + +test("callback wake is an exact-session launch fence, never domain execution authority", () => { + const value = {...nativeStartInput(), action: "wake", + launch_context: {host: "codex-cli", operation_tools: true, iteration_context: "resume"}}; + value.executor_route = value.actor; + const plan = planAgentOperationHandoff(value); + assert.equal(plan.wake_allowed, true); + assert.equal(plan.execution_allowed, false); + assert.equal(plan.host_delivery, "not_attempted"); + assert.equal(plan.write_host_start, undefined); + assert.equal(plan.write_handoff, undefined); + const changes: Array<(row: JsonObject) => void> = [ + row => {row.binding_current = false;}, + row => {(row.executor_route as JsonObject).todo_id = "other";}, + row => {(row.executor_route as JsonObject).thread_id = "replacement";}, + row => {(row.executor_route as JsonObject).profile_digest = "b".repeat(64);}, + row => {(row.executor_route as JsonObject).model = "different";}, + row => {(row.executor_route as JsonObject).reasoning_effort = "high";}, + row => {(row.launch_context as JsonObject).host = "dsh";}, + row => {(row.launch_context as JsonObject).operation_tools = false;}, + row => {(row.launch_context as JsonObject).iteration_context = "fresh";}, + row => {operation(row).confirmation = null;}, + row => {operation(row).host_start = {host_turn_id: "already-started"};}, + row => {operation(row).agent_handoff = {consumption_id: "consumed"};}, + row => {operation(row).outcome = {outcome: "submission_unknown"};}, + row => {row.now = "2030-01-01T01:00:00Z";}, + ]; + for (const change of changes) { + const row = structuredClone(value); change(row); + assert.throws(() => planAgentOperationHandoff(row)); + } + assert.throws(() => planAgentOperationHandoff({...input(), action: "wake", + executor_route: value.executor_route, launch_context: value.launch_context})); +}); + +test("transport accepted start joins the original confirmation and claim without granting execution", () => { + const value = nativeStartInput(); + const plan = planAgentOperationHandoff(value); + const receipt = plan.write_host_start as JsonObject; + assert.equal(plan.execution_allowed, false); + assert.equal(plan.recorded, true); + assert.equal(receipt.confirmation_event_id, "authenticated-confirmation-event"); + assert.equal(receipt.claim_id, "claim-1"); + assert.equal(receipt.host_turn_id, "native-turn"); + assert.equal(receipt.turn_key, value.turn_key); + assert.equal(receipt.trigger_kind, "canonical_operation_inbox"); + assert.equal(receipt.accepted_at, value.now); + assert.equal(receipt.external_write_performed, false); + assert.equal(operation(value).agent_handoff, undefined); + operation(value).host_start = receipt; + // First observation is immutable even after another normally admitted Turn. + const replay = planAgentOperationHandoff({...value, now: "2030-01-01T00:10:00Z", + actor: {...value.actor as JsonObject, host_turn_id: "later-native-turn"}}); + assert.equal(replay.recorded, false); + assert.equal(replay.write_host_start, undefined); + assert.deepEqual(replay.host_start, receipt); + assert.equal(replay.host_delivery, "native_start_accepted"); + assert.equal(planAgentOperationHandoff({...value, action: "consume"}).execution_allowed, true); +}); + +test("native start refuses drift, expiry, missing acceptance identity and effect/reconciliation state", () => { + const changes: Array<(value: JsonObject) => void> = [ + value => {value.binding_current = false;}, + value => {(value.actor as JsonObject).todo_id = "other-todo";}, + value => {(value.actor as JsonObject).thread_id = "source-thread";}, + value => {(value.actor as JsonObject).profile_digest = "b".repeat(64);}, + value => {(value.actor as JsonObject).model = "other-model";}, + value => {(value.actor as JsonObject).reasoning_effort = "high";}, + value => {(value.actor as JsonObject).host_turn_id = null;}, + value => {value.turn_key = "unaccepted-process-launch";}, + value => {value.now = "2030-01-01T01:00:00Z";}, + value => {(operation(value).confirmation as JsonObject).confirmed_at = "2030-01-01T00:01:00Z";}, + value => {operation(value).confirmation = null;}, + value => {operation(value).claim = null;}, + value => {operation(value).agent_handoff = {consumption_id: "already-consumed"};}, + value => {operation(value).outcome = {outcome: "submission_unknown"};}, + ]; + for (const change of changes) { + const value = nativeStartInput(); change(value); + assert.throws(() => planAgentOperationHandoff(value)); + } + assert.throws(() => planAgentOperationHandoff({...input(), action: "observe_host_start"})); +}); + +test("exact managed scope is filtered before bounded inbox pagination", () => { + const route = nativeStartInput().actor as JsonObject; + const items = Array.from({length: 25}, (_, index) => ({operation_id: `operation-${index}`, route: {...route, todo_id: "other-todo"}})); + items.push({operation_id: "operation-matching", route: {...route}}); + const projected = projectAgentOperationInbox({items, executor_route: route, cursor_scope: "a".repeat(64)}); + assert.equal(projected.pending_count, 1); + assert.equal((projected.items as JsonObject[])[0].operation_id, "operation-matching"); + assert.equal(projected.next_cursor, null); +}); + test("source context is not managed execution identity and old approvals never migrate", () => { const value = managedInput(); const plan = planAgentOperationHandoff(value); diff --git a/tests/extensions/test_lark_event_collector_runtime.py b/tests/extensions/test_lark_event_collector_runtime.py index 12a5f0a128..3ef9fa8c3c 100644 --- a/tests/extensions/test_lark_event_collector_runtime.py +++ b/tests/extensions/test_lark_event_collector_runtime.py @@ -228,6 +228,31 @@ def test_operation_callback_plan_requires_pinned_runtime(tmp_path: Path) -> None assert plan["operation_callback_console_configuration_preflighted"] is False +def test_callback_managed_wake_configuration_is_explicit_and_read_back(tmp_path: Path) -> None: + project, collector = _operation_callback_project(tmp_path) + data = json.loads(collector.read_text()) + config = {"registry_path": str(tmp_path / "registry.json"), "goal_id": "goal-fixture", + "requester_agent_id": "requester", "execution_config": ".loopx/config/delegations.json", + "binding_id": "operation-worker"} + data["operation_callbacks"]["managed_turn_wake"] = config + collector.write_text(json.dumps(data)) + loaded = event_collector.load_lark_event_collector_config(project=project, config_path=collector) + assert loaded["operation_callbacks"]["managed_turn_wake"] == {**config, "project": str(project)} + plan = plan_lark_event_collector(project=project, config_path=collector, runtime_root=tmp_path / "runtime") + assert plan["operation_callback_managed_wake_configured"] is True + assert "registry_path" not in json.dumps(plan) + for patch in [{"registry_path": "relative.json"}, {"execution_config": "../outside.json"}, + {"binding_id": "bad\nvalue"}, {"requester_agent_id": None}, {"unexpected": True}]: + data["operation_callbacks"]["managed_turn_wake"] = {**config, **patch} + collector.write_text(json.dumps(data)) + with pytest.raises(ValueError): + event_collector.load_lark_event_collector_config(project=project, config_path=collector) + data["operation_callbacks"] = {"enabled": False, "managed_turn_wake": config} + collector.write_text(json.dumps(data)) + with pytest.raises(ValueError, match="enabled"): + event_collector.load_lark_event_collector_config(project=project, config_path=collector) + + def test_operation_callback_status_separates_readiness_from_qualification( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, @@ -300,11 +325,21 @@ def runner(*_args: object, **_kwargs: object) -> subprocess.CompletedProcess[str assert qualified["operation_callback_qualification_state"] == "callback_qualified" +@pytest.mark.parametrize("wake_revocation", [False, True]) def test_collector_runs_independent_operation_callback_consumer( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, + wake_revocation: bool, ) -> None: project, collector = _operation_callback_project(tmp_path) + if wake_revocation: + config = json.loads(collector.read_text()) + config["operation_callbacks"]["managed_turn_wake"] = { + "registry_path": str(tmp_path / "registry.json"), "goal_id": "fixture-goal", + "requester_agent_id": "coordinator", "execution_config": "delegations.json", + "binding_id": "confirmed-operation", + } + collector.write_text(json.dumps(config)) runtime_root = tmp_path / "runtime" cli = tmp_path / "lark-cli-fixture" cli.write_text( @@ -327,6 +362,10 @@ def test_collector_runs_independent_operation_callback_consumer( def handle(payload: dict[str, object], **kwargs: object) -> dict[str, object]: captured.append({"payload": payload, **kwargs}) + if wake_revocation and len(captured) == 1: + config = json.loads(collector.read_text()) + del config["operation_callbacks"]["managed_turn_wake"] + collector.write_text(json.dumps(config)) return { "ok": len(captured) > 1, "schema_version": "lark_operation_callback_receipt_v0", @@ -364,6 +403,8 @@ def runner(argv: list[str], **_kwargs: object) -> subprocess.CompletedProcess[st assert result["operation_callback_verified_count"] == 1 assert captured[0]["runtime_root"] == runtime_root.resolve() assert captured[0]["action_store_root"] == runtime_root / "chat" / "actions" + assert (captured[0]["managed_turn_wake"] is not None) is wake_revocation + assert captured[1]["managed_turn_wake"] is None status = json.loads( ( project / ".loopx/runtime/lark-collector/operation-callback-status.json" diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py index 69e1300714..9da218f7c5 100644 --- a/tests/extensions/test_lark_goal_channel_operation.py +++ b/tests/extensions/test_lark_goal_channel_operation.py @@ -117,6 +117,76 @@ def _prepare_agent_handoff( ) +@pytest.mark.parametrize("fault", [None, "grant", "todo", "profile", "fresh", "stopped", "dispatch"]) +def test_callback_wake_uses_one_operator_granted_delegation_without_consuming( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, fault: str | None, +) -> None: + from concurrent.futures import ThreadPoolExecutor + from loopx.collaboration_mcp import Delegations + + store, registry, runtime, channel, target = _fixture(tmp_path) + proposal = _prepare_agent_handoff(store, registry, managed=True) + data = json.loads(registry.read_text()) + data["goals"][0]["coordination"]["registered_agents"].append("callback-requester") + if fault == "stopped": + data["goals"][0]["status"] = "stopped" + registry.write_text(json.dumps(data)) + config = registry.parent / "delegations.json" + worker = tmp_path / "worker" + worker.mkdir() + binding = {"id": "confirmed-operation", "agent_id": AGENT_ID, + "todo_id": "other" if fault == "todo" else "todo-managed", + "requesters": [] if fault == "grant" else ["callback-requester"], + "workspace": str(worker), "timeout_seconds": 30, "output_refs": ["result.json"], + "host_args": ["--host", "codex-cli", "--codex-operation-tools", "--codex-model", + "other" if fault == "profile" else "test-model", "--codex-reasoning-effort", "xhigh"]} + if fault == "fresh": + binding["host_args"] += ["--iteration-context", "fresh"] + config.write_text(json.dumps({"schema_version": "loopx_local_delegation_v0", "bindings": [binding]})) + launched = [] + def spawn(self, operation_id): + launched.append(operation_id) + if fault == "dispatch": + raise RuntimeError("private argv must not be returned") + monkeypatch.setattr(Delegations, "_spawn", spawn) + cards = {} + runner = _runner([], cards) + deliver_goal_channel_operation_card(proposal_id=proposal["proposal_id"], action_store_root=store.root, + runtime_root=runtime, binding_path=channel, target_path=target, execute=True, runner=runner) + delivered = store.load(proposal["proposal_id"]) + event = _event(delivered, cards[delivered["operation"]["delivery"]["message_id"]]) + configuration = {"registry_path": str(registry), "goal_id": GOAL_ID, + "requester_agent_id": "callback-requester", "binding_id": binding["id"], + "project": str(registry.parent.parent), "execution_config": ".loopx/delegations.json"} + kwargs = dict(runtime_root=runtime, action_store_root=store.root, profile_app_id=APP_ID, + cli_bin="lark-cli", profile="operation-bot", runner=runner, managed_turn_wake=configuration) + first = handle_goal_channel_operation_callback(event, **kwargs) + # Concurrent callback replay/lost ACK reuses the actual canonical journal. + with ThreadPoolExecutor(max_workers=2) as pool: + replays = list(pool.map(lambda _: handle_goal_channel_operation_callback(event, **kwargs), range(2))) + assert first["ok"] and first["status"] == "authorization_pending" + assert first["managed_turn_wake"]["execution_allowed"] is False + assert first["managed_turn_wake"]["native_start_verified"] is False + assert "private argv" not in json.dumps(first) + stored = store.load(proposal["proposal_id"]) + assert stored["operation"]["lifecycle_state"] == "claimed" + assert stored["operation"].get("agent_handoff") is None + assert stored["operation"].get("host_start") is None + assert stored["operation"]["outcome"] is None + if fault in {None, "dispatch"}: + assert len(launched) == 1 + assert all(row["managed_turn_wake"]["state"] == "existing_delegation" for row in replays) + service = Delegations(runtime, registry, GOAL_ID, "callback-requester", config) + row = json.loads(service.path(launched[0]).read_text()) + assert row["identity"]["confirmed_operation_id"] == proposal["proposal_id"] + argv = service._execution_arguments(binding, launched[0]) + assert argv[-2:] == ["--codex-confirmed-operation-id", proposal["proposal_id"]] + assert first["managed_turn_wake"]["state"] == ("blocked" if fault else "delegation_requested") + else: + assert launched == [] + assert first["managed_turn_wake"]["state"] == "blocked" + + @pytest.mark.parametrize("managed", [False, True]) def test_authenticated_callback_hands_off_without_calling_any_executor_and_reconciles_original_result( tmp_path: Path, @@ -164,6 +234,10 @@ def no_executor(_proposal): assert first["status"] == replay["status"] == "authorization_pending" assert first["outcome"] is None and not first["domain_external_write_performed"] assert first["callback_ack_is_execution_receipt"] is False + for receipt in [first, replay]: + assert receipt["managed_turn_wake"]["state"] == "not_configured" + assert receipt["managed_turn_wake"]["execution_allowed"] is False + assert receipt["managed_turn_wake"]["native_start_verified"] is False claimed = store.load(proposal["proposal_id"]) assert claimed["operation"]["result_delivery"] is None assert ( @@ -186,6 +260,13 @@ def no_executor(_proposal): model="test-model", reasoning_effort="xhigh", ) + start = agent_operation_action( + runtime, registry, proposal_id=proposal["proposal_id"], actor=actor, + action="observe_host_start", turn_key="sha256:" + "b" * 64, + ) + assert start["recorded"] is True and start["execution_allowed"] is False + started_card = build_goal_channel_operation_result_card(store.load(proposal["proposal_id"])) + assert "原生续接已接受;授权仍待消费,尚无执行结果" in normalized_card_text(started_card) args = Namespace( goal_channel_command="consume-operation", goal_id=GOAL_ID, diff --git a/tests/test_codex_operation_host.py b/tests/test_codex_operation_host.py index 1d6faa56d0..c4aaf5e52d 100644 --- a/tests/test_codex_operation_host.py +++ b/tests/test_codex_operation_host.py @@ -29,7 +29,8 @@ FAKE_SERVER = """#!/usr/bin/env python3 import json, sys thread = "owned-app-server-thread" -turn = "native-app-server-turn" +import os +turn = os.environ.get("FAKE_OPERATION_TURN_ID", "native-app-server-turn") key = None def emit(value): print(json.dumps(value), flush=True) @@ -58,6 +59,8 @@ def emit(value): if os.environ.get("FAKE_OPERATION_HANG") == "1": while True: time.sleep(.1) text = row["params"]["input"][0]["text"] + if os.environ.get("FAKE_OPERATION_PROMPT_FILE"): + pathlib.Path(os.environ["FAKE_OPERATION_PROMPT_FILE"]).write_text(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 @@ -85,6 +88,221 @@ def emit(value): """ +def _claimed_native_fixture(tmp_path: Path): + from examples import operation_action_fixtures as fixtures + from loopx.control_plane.turn_driver.codex_operation_host import operation_tool_handler + + service, store = fixtures.service(tmp_path, goal_id="fixture-goal") + executable = tmp_path / "fake-codex-operation" + executable.write_text(FAKE_SERVER) + executable.chmod(0o700) + request = _request() + request["turn_envelope"]["agent_id"] = "finance-fixture-agent" + request["turn_envelope"]["action"]["selected_todo"]["todo_id"] = "todo-managed" + options = dict(runtime_root=store.root.parent.parent, registry_path=service.registry_path, + project=service.registry_path.parent.parent, codex_bin=str(executable), + model="test-model", reasoning_effort="xhigh", timeout_seconds=5) + run_codex_operation_host(request, **options) + binding = load_codex_cli_session(options["runtime_root"], lineage=_lineage(request)) + handler = operation_tool_handler( + runtime_root=options["runtime_root"], registry_path=service.registry_path, + lineage=_lineage(request), session_id=binding["session_id"], + profile_digest=binding["operation_profile_digest"], model="test-model", reasoning_effort="xhigh", + ) + operation_request = fixtures.request(goal_id="fixture-goal", + payload={"schema_version": "qualification_v0", "marker": "synthetic"}) + terms = operation_request["normalized_parameters"] + terms.pop("executor") + terms.update(domain="qualification", operation_kind="qualification.observe", operation_schema="qualification_v0") + terms["projection"].update(simulated=False, title="Synthetic qualification", focus="No external effect") + prepared = handler("loopx_operation", {"action": "prepare", "request": operation_request}, + {"thread_id": binding["session_id"], "host_turn_id": "preparing-native-turn"}) + assert prepared["ok"], prepared + proposal = prepared["proposal"] + delivered = store.record_operation_delivery(proposal["proposal_id"], delivery=fixtures.delivery(proposal)) + claimed = store.decide_operation(proposal["proposal_id"], decision="confirm", confirmation=fixtures.confirmation(delivered)) + request["session"]["action"] = "resume" + request["turn_key"] = "sha256:" + "b" * 64 + return store, claimed, request, options + + +def test_confirmed_operation_is_automatically_carried_into_accepted_native_turn( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + store, claimed, request, options = _claimed_native_fixture(tmp_path) + prompt = tmp_path / "native-prompt.txt" + monkeypatch.setenv("FAKE_OPERATION_PROMPT_FILE", str(prompt)) + result = run_codex_operation_host(request, **options) + assert result["result_kind"] == "wait" + stored = store.load(claimed["proposal_id"]) + receipt = stored["operation"]["host_start"] + assert claimed["proposal_id"] in prompt.read_text() + assert claimed["operation"]["claim"]["claim_id"] in prompt.read_text() + assert receipt["confirmation_event_id"] == claimed["operation"]["confirmation"]["event_id"] + assert receipt["claim_id"] == claimed["operation"]["claim"]["claim_id"] + assert receipt["host_turn_id"] == "native-app-server-turn" + assert receipt["turn_key"] == request["turn_key"] + assert receipt["confirmed_at"] <= receipt["accepted_at"] + assert receipt["execution_allowed"] is False + assert stored["operation"].get("agent_handoff") is None + assert stored["operation"]["outcome"] is None + monkeypatch.setenv("FAKE_OPERATION_TURN_ID", "later-native-turn") + request["turn_key"] = "sha256:" + "c" * 64 + run_codex_operation_host(request, **options) + assert store.load(claimed["proposal_id"])["operation"]["host_start"] == receipt + + +def test_confirmed_callback_launch_fence_resumes_only_the_original_native_session(tmp_path: Path) -> None: + store, claimed, request, options = _claimed_native_fixture(tmp_path) + run_codex_operation_host(request, confirmed_operation_id=claimed["proposal_id"], **options) + assert store.load(claimed["proposal_id"])["operation"]["host_start"]["turn_key"] == request["turn_key"] + assert store.load(claimed["proposal_id"])["operation"].get("agent_handoff") is None + # A callback/recovery cannot start another native Turn after acceptance. + with pytest.raises(Exception, match="unstarted"): + run_codex_operation_host(request, confirmed_operation_id=claimed["proposal_id"], **options) + + +@pytest.mark.parametrize("drift", ["session", "todo", "profile", "fresh", "stopped", "missing"]) +def test_callback_launch_fence_rejects_drift_before_native_process_start( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, drift: str, +) -> None: + from loopx.control_plane.turn_driver import codex_operation_host as owner + from loopx.control_plane.turn_driver.codex_cli import _store_codex_cli_session + + store, claimed, request, options = _claimed_native_fixture(tmp_path) + if drift == "session": + original = load_codex_cli_session(options["runtime_root"], lineage=_lineage(request)) + _store_codex_cli_session(options["runtime_root"], lineage=_lineage(request), session_id="replacement", + operation_profile_digest=original["operation_profile_digest"], operation_model="test-model", + operation_reasoning_effort="xhigh") + elif drift == "todo": + request["turn_envelope"]["action"]["selected_todo"]["todo_id"] = "other" + elif drift == "profile": + options["reasoning_effort"] = "high" + elif drift == "fresh": + request["session"]["context_policy"] = {"mode": "fresh"} + elif drift == "stopped": + data = json.loads(options["registry_path"].read_text()) + data["goals"][0]["status"] = "stopped" + options["registry_path"].write_text(json.dumps(data)) + def no_start(*args, **kwargs): + pytest.fail("the fenced callback must not start/rebind a native process") + monkeypatch.setattr(owner.CodexChatAgentSession, "start", no_start) + before = store.path.read_bytes() + with pytest.raises((ValueError, RuntimeError)): + run_codex_operation_host(request, confirmed_operation_id="missing" if drift == "missing" else claimed["proposal_id"], **options) + assert store.path.read_bytes() == before + + +def test_dynamic_pending_filters_other_managed_tasks_before_pagination( + tmp_path: Path, +) -> None: + from examples import operation_action_fixtures as fixtures + + service, store = fixtures.service(tmp_path, goal_id="fixture-goal") + current = fixtures.managed_handler(service, store, goal_id="fixture-goal") + other = fixtures.managed_handler( + service, store, goal_id="fixture-goal", todo_id="todo-other", + session_id="other-managed-thread", + ) + + def prepare_and_confirm(handler, index: int, thread_id: str): + request = fixtures.request( + goal_id="fixture-goal", + payload={"schema_version": "qualification_v0", "marker": "synthetic"}, + ) + request["idempotency_key"] = f"pending-fixture-{index}" + terms = request["normalized_parameters"] + terms.pop("executor") + terms.update(domain="qualification", operation_kind="qualification.observe", operation_schema="qualification_v0") + terms["projection"].update(simulated=False, title="Synthetic qualification", focus="No external effect") + prepared = handler( + "loopx_operation", {"action": "prepare", "request": request}, + {"thread_id": thread_id, "host_turn_id": "fixture-native-turn"}, + ) + assert prepared["ok"], prepared + proposal = prepared["proposal"] + delivered = store.record_operation_delivery( + proposal["proposal_id"], delivery=fixtures.delivery(proposal), + ) + return store.decide_operation( + proposal["proposal_id"], decision="confirm", + confirmation=fixtures.confirmation(delivered), + )["proposal_id"] + + # All these are current, approved canonical operations for the SAME Agent. + # They must neither leak into this task nor exhaust its first page. + other_ids = { + prepare_and_confirm(other, index, "other-managed-thread") + for index in range(25) + } + own_id = prepare_and_confirm(current, 25, "owned-managed-thread") + pending = current( + "loopx_operation", {"action": "pending"}, + {"thread_id": "owned-managed-thread", "host_turn_id": "fixture-native-turn"}, + ) + assert pending["ok"], pending + assert pending["pending_count"] == 1 + assert {item["operation_id"] for item in pending["items"]} == {own_id} + assert pending["next_cursor"] is None + assert not other_ids.intersection(item["operation_id"] for item in pending["items"]) + + # A task-local page token is not reusable by another execution subject. + for index in range(26, 50): + prepare_and_confirm(current, index, "owned-managed-thread") + page = current( + "loopx_operation", {"action": "pending"}, + {"thread_id": "owned-managed-thread", "host_turn_id": "fixture-native-turn"}, + ) + assert page["next_cursor"] + rejected = other( + "loopx_operation", {"action": "pending", "cursor": page["next_cursor"]}, + {"thread_id": "other-managed-thread", "host_turn_id": "fixture-native-turn"}, + ) + assert rejected["ok"] is False + assert rejected["execution_allowed"] is False + + +def test_process_launch_without_native_acceptance_cannot_record_operation_start( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + store, claimed, request, options = _claimed_native_fixture(tmp_path) + monkeypatch.setenv("FAKE_OPERATION_HANG", "1") + with pytest.raises(BuiltInHostError): + run_codex_operation_host(request, **{**options, "timeout_seconds": 0.5}) + assert store.load(claimed["proposal_id"])["operation"].get("host_start") is None + + +def test_start_receipt_failure_aborts_before_operation_tool_dispatch( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, +) -> None: + from loopx.chat_action_store import ChatActionStore + from loopx.control_plane.turn_driver import codex_operation_host + store, claimed, request, options = _claimed_native_fixture(tmp_path) + before = store.path.read_bytes() + calls = [] + original = codex_operation_host.operation_tool_handler + + def observe(**options): + handler = original(**options) + + def record(*args): + calls.append(args) + return handler(*args) + + return record + + def fail(*args, **kwargs): + raise RuntimeError("synthetic receipt write failure") + + monkeypatch.setattr(codex_operation_host, "operation_tool_handler", observe) + monkeypatch.setattr(ChatActionStore, "record_agent_operation_host_start", fail) + with pytest.raises(BuiltInHostError) as error: + run_codex_operation_host(request, **options) + assert str(error.value) == "codex_operation_host_start_unrecorded" + assert calls == [] + assert store.path.read_bytes() == before + assert store.load(claimed["proposal_id"])["operation"].get("agent_handoff") is None SOURCE_PREPARE_SERVER = """#!/usr/bin/env python3 import json, os, sys thread, turn = "owned-source-test-thread", "native-source-test-turn" diff --git a/tests/test_local_delegation.py b/tests/test_local_delegation.py index d29b13fd7e..4ece5df6a1 100644 --- a/tests/test_local_delegation.py +++ b/tests/test_local_delegation.py @@ -149,6 +149,108 @@ def wait(service, operation="analysis-1"): pytest.fail(str(service.read(operation))) +@pytest.mark.parametrize("mcp_profile", ["delegation-default", "explicit-null", "unmatched-null"]) +def test_confirmed_wake_reaches_native_acceptance_through_real_delegation_and_turn(service, mcp_profile): + """Real detached worker/CLI/Turn; synthetic native transport, no model/effect.""" + from examples import operation_action_fixtures as fixtures + from loopx.chat_action_store import ChatActionStore + from loopx.cli import build_parser + from loopx.control_plane.collaboration.operation_wake import dispatch_confirmed_operation_wake + from loopx.control_plane.turn_driver.codex_cli import load_codex_cli_session + from loopx.control_plane.turn_driver.codex_operation_host import run_codex_operation_host, operation_tool_handler + from test_codex_operation_host import FAKE_SERVER + from test_loopx_turn_codex_cli import _request + + root, runner = service + executable = root / "fixture-codex" + native_starts = root / "native-starts" + startup_probe = ( + "from pathlib import Path\n" + f"counter = Path({str(native_starts)!r})\n" + "counter.write_text(str(int(counter.read_text()) + 1 if counter.exists() else 1))\n" + ) + executable.write_text(FAKE_SERVER.replace('thread = "owned-app-server-thread"', + startup_probe + 'thread = "owned-app-server-thread"')) + executable.chmod(0o700) + config = json.loads(runner.config.read_text()) + binding = config["bindings"][0] + binding["host_args"] = ["--host", "codex-cli", "--codex-bin", str(executable), + "--codex-operation-tools", "--codex-model", "test-model", "--codex-reasoning-effort", "xhigh"] + if mcp_profile == "explicit-null": + # Existing operator option, not a new profile or a worker-selected grant. + binding["host_args"] += ["--codex-mcp-server-json", "null"] + runner.config.write_text(json.dumps(config)) + argv = runner._execution_arguments(binding, "preparation") + parsed = build_parser().parse_args(["turn", "run-once", "--goal-id", runner.goal_id, + "--agent-id", binding["agent_id"], "--todo-id", binding["todo_id"], *argv]) + if mcp_profile == "explicit-null": + assert argv.count("--codex-mcp-server-json") == 2 + assert parsed.codex_mcp_server_json is None + else: + assert parsed.codex_mcp_server_json["name"] == "loopx_delegation" + # Reproduce an original standalone operation Session: do NOT prepare it + # with the injected delegation MCP just to make the later resume match. + mcp_server = parsed.codex_mcp_server_json if mcp_profile == "delegation-default" else None + lineage = {"goal_id": runner.goal_id, "agent_id": binding["agent_id"], "todo_id": binding["todo_id"]} + request = _request() + request["turn_envelope"].update(goal_id=runner.goal_id, agent_id=binding["agent_id"]) + request["turn_envelope"]["action"]["selected_todo"]["todo_id"] = binding["todo_id"] + run_codex_operation_host(request, runtime_root=runner.root, registry_path=runner.registry, + project=Path(binding["workspace"]), codex_bin=str(executable), model="test-model", + reasoning_effort="xhigh", mcp_server=mcp_server, timeout_seconds=5) + assert native_starts.read_text() == "1" + session = load_codex_cli_session(runner.root, lineage=lineage) + handler = operation_tool_handler(runtime_root=runner.root, registry_path=runner.registry, + lineage=lineage, session_id=session["session_id"], profile_digest=session["operation_profile_digest"], + model="test-model", reasoning_effort="xhigh") + intent = fixtures.request(goal_id=runner.goal_id, + payload={"schema_version": "qualification_v0", "marker": "synthetic"}) + terms = intent["normalized_parameters"] + terms.pop("executor") + terms.update(agent_id=binding["agent_id"], domain="qualification", + operation_kind="qualification.observe", operation_schema="qualification_v0") + terms["projection"].update(simulated=False, title="Synthetic native wake qualification") + prepared = handler("loopx_operation", {"action": "prepare", "request": intent}, + {"thread_id": session["session_id"], "host_turn_id": "preparing-turn"}) + assert prepared["ok"], prepared + store = ChatActionStore(runner.root / "chat" / "actions") + proposal = prepared["proposal"] + delivered = store.record_operation_delivery(proposal["proposal_id"], delivery=fixtures.delivery(proposal)) + confirmed = store.decide_operation(proposal["proposal_id"], decision="confirm", + confirmation=fixtures.confirmation(delivered)) + wake_config = {"registry_path": str(runner.registry), "goal_id": runner.goal_id, + "requester_agent_id": runner.agent_id, "project": str(root), + "execution_config": "delegations.json", "binding_id": binding["id"]} + receipt = dispatch_confirmed_operation_wake(confirmed, runtime_root=runner.root, configuration=wake_config) + assert receipt["state"] == "delegation_requested", receipt + result = wait(runner, receipt["operation_id"]) + stored = store.load(proposal["proposal_id"]) + resumed = load_codex_cli_session(runner.root, lineage=lineage) + assert resumed["session_id"] == session["session_id"] + assert resumed["operation_profile_digest"] == session["operation_profile_digest"] + assert stored["operation"].get("agent_handoff") is None + assert stored["operation"]["outcome"] is None + # Omitting the original null override must still fail closed, rather than + # changing/replacing the original Session to bypass the profile fence. + if mcp_profile == "unmatched-null": + assert result["status"] == "rejected", result + assert stored["operation"].get("host_start") is None, result + assert native_starts.read_text() == "1" # refused before native process start + assert stored["status"] == confirmed["status"] + assert stored["operation"]["confirmation"] == confirmed["operation"]["confirmation"] + return + assert native_starts.read_text() == "2" + assert stored["operation"]["host_start"] is not None, result + assert stored["operation"]["host_start"]["route"]["thread_id"] == session["session_id"] + # Fixture deliberately waits: startup is not validated domain completion. + assert result["status"] == "rejected" + journal = json.loads(runner.path(receipt["operation_id"]).read_text()) + assert stored["operation"]["host_start"]["turn_key"] == journal["turn_key"] + assert dispatch_confirmed_operation_wake(stored, runtime_root=runner.root, + configuration=wake_config)["state"] == "existing_delegation" + assert native_starts.read_text() == "2" + + def test_detached_result_reconnects_without_duplicate_execution(service): root, original = service (root / "hold").touch()