diff --git a/apps/presentation/dashboard/src/data/status.ts b/apps/presentation/dashboard/src/data/status.ts index 61d91e2d41..a835dd9388 100644 --- a/apps/presentation/dashboard/src/data/status.ts +++ b/apps/presentation/dashboard/src/data/status.ts @@ -363,8 +363,8 @@ export const projectAssetTodoProjectionGapSchema = z.object({ export const nativeChildActivitySchema = z.object({ schema_version: z.literal("native_subagent_activity_v0"), - observation: z.enum(["unknown", "coordinator_reported"]), - host_attested: z.literal(false), + observation: z.enum(["unknown", "coordinator_reported", "host_observed", "mixed"]), + host_attested: z.boolean(), configured_limit: z.number().int().nonnegative(), launched_count: z.number().int().nonnegative(), skipped_count: z.number().int().nonnegative(), 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 6415d75441..42c79930cd 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/context-drawer.tsx @@ -801,10 +801,12 @@ export function ContextDrawer({ agents, attentionHistory = [], onSelectAttention : null} - {selection.item.nativeChildActivity?.observation === "coordinator_reported" ? ( + {selection.item.nativeChildActivity && selection.item.nativeChildActivity.observation !== "unknown" ? (

{t("drawer.subagentReportTitle")}

-

{t("drawer.subagentReportedActivity", { +

{t(selection.item.nativeChildActivity.observation === "host_observed" + ? "drawer.subagentHostActivity" : selection.item.nativeChildActivity.observation === "mixed" + ? "drawer.subagentMixedActivity" : "drawer.subagentReportedActivity", { started: selection.item.nativeChildActivity.launched_count, skipped: selection.item.nativeChildActivity.skipped_count, rejected: selection.item.nativeChildActivity.capacity_rejected_count, diff --git a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx index bb592e78a2..1f0ecc3a84 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx @@ -296,6 +296,8 @@ const en = { "drawer.subagentCurrentBoundary": "Current task-domain restriction", "drawer.subagentDescription": "Allows the runtime to create temporary child agents for independent tasks only after Todo, quota, capability, and write-scope gates pass. It does not force parallel work or grant durable authority.", "drawer.subagentReportTitle": "Child activity", + "drawer.subagentHostActivity": "Host-observed activity: {started} starts, {rejected} capacity rejections, {failed} host failures, {accepted} parent-accepted results.", + "drawer.subagentMixedActivity": "Host observations and coordinator reports: {started} starts, {skipped} skips, {rejected} capacity rejections, {failed} host failures, {accepted} parent-accepted results. Some decisions are unverified.", "drawer.subagentReportedActivity": "Latest coordinator report: {started} starts, {skipped} skips, {rejected} capacity rejections, {failed} host failures, {accepted} parent-accepted results. Host verification is unavailable.", "drawer.subagentDisable": "Preview turning off sub-agent execution", "drawer.subagentDisableSummary": "New child-agent execution will be disabled for this Goal. Existing Todo ownership and execution records stay unchanged.", @@ -1574,6 +1576,8 @@ const zhCN: Record = { "drawer.subagentCurrentBoundary": "当前任务领域限制", "drawer.subagentDescription": "仅在 Todo、配额、能力和写入范围门禁全部通过后,允许运行时为相互独立的任务临时创建子代理;不会强制并行,也不会授予持久权限。", "drawer.subagentReportTitle": "子代理活动", + "drawer.subagentHostActivity": "宿主已观察:启动 {started} 次、容量拒绝 {rejected} 次、宿主失败 {failed} 次、主 Agent 验收 {accepted} 项。", + "drawer.subagentMixedActivity": "宿主观察与主 Agent 回报:启动 {started} 次、跳过 {skipped} 次、容量拒绝 {rejected} 次、宿主失败 {failed} 次、主 Agent 验收 {accepted} 项;部分决策未经宿主核验。", "drawer.subagentReportedActivity": "最近一轮主 Agent 回报:启动 {started} 次、跳过 {skipped} 次、容量拒绝 {rejected} 次、宿主失败 {failed} 次、主 Agent 验收 {accepted} 项;目前没有宿主核验。", "drawer.subagentDisable": "预览关闭子代理执行", "drawer.subagentDisableSummary": "这个 Goal 将不再创建新的子代理;现有 Todo 归属和执行记录不受影响。", diff --git a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-model.ts b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-model.ts index def13b066e..03f703fd50 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-model.ts +++ b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-model.ts @@ -159,8 +159,8 @@ export type WorkspaceGoal = { subagentExecution?: WorkspaceGoalSubagentConfiguration; nativeChildActivity?: { turn_instance_id: string; - observation: "unknown" | "coordinator_reported"; - host_attested: false; + observation: "unknown" | "coordinator_reported" | "host_observed" | "mixed"; + host_attested: boolean; launched_count: number; skipped_count: number; capacity_rejected_count: number; diff --git a/apps/presentation/dashboard/src/views/dashboard-page.tsx b/apps/presentation/dashboard/src/views/dashboard-page.tsx index f39a930c7a..2e8fd7679d 100644 --- a/apps/presentation/dashboard/src/views/dashboard-page.tsx +++ b/apps/presentation/dashboard/src/views/dashboard-page.tsx @@ -468,8 +468,8 @@ type PersonalGoalItem = { hasRunObservation: boolean; nativeChildActivity?: { turn_instance_id: string; - observation: "unknown" | "coordinator_reported"; - host_attested: false; + observation: "unknown" | "coordinator_reported" | "host_observed" | "mixed"; + host_attested: boolean; launched_count: number; skipped_count: number; capacity_rejected_count: number; diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index b30eeaf878..597555b7a5 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -16,20 +16,32 @@ Goal 事件流,以 Turn ID、稳定操作 ID 和不限定宿主的 `entrypoint 记录动作不会启动子代理、调度新 Turn、授予写入权限或消耗配额。 `max_children` 只是配置上限,不代表当前可用槽位,也不是必须启动的数量。 -`native-child record` currently accepts the coordinator's typed report. Its -`observation` is `coordinator_reported` and `host_attested` is always `false`. -LoopX cannot intercept an arbitrary external host's native tool call. A future -host adapter must observe that call at its own boundary and extend this event -and read-model contract with a separately verified provenance variant; this -v0 recorder cannot claim host attestation. Missing records remain -`unknown`; neither a missing record nor `max_children > 0` proves that a child -was created or deliberately skipped. - -目前 `native-child record` 接受主 Agent 的类型化上报,因此 `observation` 为 -`coordinator_reported`,`host_attested` 始终为 `false`。LoopX 无法拦截任意 -外部宿主的原生工具调用。后续宿主适配器须在自己的边界观察调用,给事件与 -读模型扩展单独核验的来源类型;当前 v0 上报器不能声称宿主核验。缺少回执 -就是 `unknown`;没有回执或配置上限大于零,都不能证明已启动或主动跳过。 +`native-child record` accepts a coordinator report and cannot select host +provenance. Managed Codex CLI and operation-equipped app-server Turns also +observe native collaboration items directly on their owned connection. A +successful spawn, failed call and observed child completion become durable +`host_observed` records before Turn settlement. A parent review is a separate +explicit record; a host completion does not adopt the evidence. + +The shared projection distinguishes `host_observed`, `coordinator_reported`, +`mixed` and `unknown`. `host_attested` is true only when every decision in the +Turn came from the host. Configured capacity remains an upper bound. Failed +Codex collaboration items do not carry a typed capacity error, so the adapter +records `host_failed` and forbids report-based same-Turn retry rather than +classifying provider prose. Missing native events stay unknown. Persisted Codex +history is not used to reconstruct native activity because supported host +versions may omit collaboration items from that history. + +`native-child record` 仍是主 Agent 上报入口,不能指定宿主来源。托管 Codex CLI +与启用操作工具的 app-server Turn 会在自身连接上直接观察原生协作事件, +把实际启动、宿主失败和观察到的结果写入同一事件流。结果完成之后仍须由主 +Agent 明确记录验收;宿主完成不代表证据已被采纳。 + +共享投影区分宿主观察、主 Agent 上报、混合来源和未知;仅全部决策均来自 +宿主时 `host_attested` 为真。Codex 的失败协作项没有容量错误码,适配器保留 +通用 `host_failed` 和禁止同 Turn 上报重试的规则,不从错误文字推断容量。 +缺少原生事件仍是未知。部分宿主版本的持久化历史会遗漏协作项,因此不用于 +重建活动。第三方宿主上报仍不会自动获得宿主核验标记。 ## Lifecycle / 生命周期 @@ -106,3 +118,65 @@ for the parent validation of the underlying work. 已绑定的 LoopX delegation 仍以自身操作回执为权威。原生子代理上报不能代替 delegation 回执,也不能代替主 Agent 对工作结果的实际核验。 + +## Codex host qualification / Codex 宿主验证 + +The adapter belongs to the built-in Codex Turn host and the existing +`multi_subagent` receipt owner. It installs no scheduler and makes no additional +provider request. Feature-off Turns create no observer and retain their host +request/result contract. Ordinary CLI and Lark status use the same shared +projection as the dashboard; Lark has no native child configuration to change. + +A resumed CLI invocation can receive only a new `wait` completion. The adapter +resolves its hashed child reference against the latest started host-observed +spawn or followup in the +same admitted Goal instance, coordinator and LoopX Turn, including receipts +outside the bounded status window. A spawn/followup item's terminal snapshot +belongs to its own stable decision and receivers; replaying an old spawn never +completes a later followup. A terminal wait first resolves the immutable binding +of its session, invocation, native item and child identity. Only its first +observation uses the latest child association. Each distinct wait retains a +hashed scalar reference on the existing result event; exact replay preserves +that original operation across restart, and changed outcomes or reassignment +are rejected. Multiple waits observing one result do not add operations, +launches, parent acceptance or quota. No raw host identifiers or content enter +the public activity projection. +It does not adopt coordinator reports or scan external host history. Spawn IDs +retain their existing child binding. +CLI followup and failed-call IDs use the Turn journal's durable `host_attempt` +plus the owned parent session and native item ID; app-server calls use their +native Turn ID. A real retry advances the journal attempt before launch; +replaying the same binding and item remains idempotent. Direct enabled CLI +adapter calls require that attempt; feature-off calls retain their original +request. Only compact hashed child references are retained for correlation; +they do not enter public activity rows or replace independent parent review. +No prompt or raw result is added to a receipt. + +恢复 CLI 时可能只收到新的 `wait` 完成事件。适配器从同一已准入 Goal 实例、 +主 Agent、LoopX Turn 的最近一次已启动宿主 spawn 或 followup 回执恢复关联, +覆盖状态窗口外的操作。spawn/followup 的完成快照只归属自身稳定决策和接收者; +重放旧启动事件不会完成后来的跟进任务。wait 首先恢复其会话、调用、原生项和 +接收子 Agent 身份对应的首次结果关联;只有首次观察才使用最近操作。哈希关联 +作为标量留在现有结果事件中,重启重放仍归原操作,结果冲突或重新归属会被拒绝。 +多个 wait 观察同一结果不会增加操作、启动、主 Agent 验收或配额;公开活动投影 +不暴露原始宿主标识或内容。 +不采纳主 Agent 上报或扫描外部宿主历史。启动保留既有子代理绑定;CLI 跟进和 +失败调用复用 Turn 日志持久化的 `host_attempt`、父会话和原生工具 ID, +app-server 使用其原生 Turn ID。真实重试在启动前递增尝试次数;同一绑定与 +事件的重放仍幂等。直接调用已启用的 CLI 适配器须提供该尝试次数,关闭能力 +时请求不变。关联只保留紧凑的子代理哈希引用,不进入公开活动行, +也不替代独立父任务验收。回执不增加原始提示或结果内容。 + +A live isolated Codex 0.142.5 test observed one successful spawn, a second failed +spawn at `agents.max_threads=1`, child completion and independent parent +acceptance. Durable readback preserved one launch and one accepted result. The +native failure subtype remains unqualified: Codex emitted only `failed`, not +`agent_thread_limit_reached`. Synthetic typed capacity-report tests cover the +existing no-same-Turn retry rule without pretending that this host supplies that +error code. No live production Goal or raw child output is part of this evidence. + +适配器复用内置 Codex Turn 宿主与 `multi_subagent` 回执所有者,不新增调度器 +或模型调用。关闭能力时不创建观察器。CLI、Lark 状态和仪表板共用同一投影。 +隔离的真实 Codex 0.142.5 验证观察到了一个成功启动、上限为 1 时第二次启动 +失败、首个子任务完成以及独立的父任务验收。持久回读保留一个启动和一个 +采纳结果。失败子类型仍有宿主协议缺口,不能声称已核验容量错误码。 diff --git a/examples/personal-workspace-browser-smoke.mjs b/examples/personal-workspace-browser-smoke.mjs index 3f6291729b..ee9cf4830d 100644 --- a/examples/personal-workspace-browser-smoke.mjs +++ b/examples/personal-workspace-browser-smoke.mjs @@ -1,4 +1,5 @@ #!/usr/bin/env node +import {nativeChildActivityScenario} from "./personal-workspace-browser/native-child-activity.mjs"; import {conversationImageRequestScenario} from "./personal-workspace-browser/conversation-image-request.mjs"; // Isolated browser acceptance scenarios for the personal Agent workspace. @@ -72,6 +73,7 @@ scenarioCatalog.push(turnStepsScenario); scenarioCatalog.push(goalWorkMapScenario); scenarioCatalog.push(performanceDiagnosisScenario); scenarioCatalog.push(blockedNoticeSettingsScenario); +scenarioCatalog.push(nativeChildActivityScenario); const requestedScenario = process.env.LOOPX_PERSONAL_WORKSPACE_SCENARIO; const scenarios = requestedScenario ? scenarioCatalog.filter((scenario) => scenario.id === requestedScenario) diff --git a/examples/personal-workspace-browser/fixture.mjs b/examples/personal-workspace-browser/fixture.mjs index 829db0433f..7326fa376d 100644 --- a/examples/personal-workspace-browser/fixture.mjs +++ b/examples/personal-workspace-browser/fixture.mjs @@ -558,6 +558,10 @@ export async function installApi(page, { goalSubagentConfigurationEnabled = true }); } const first = fixture.attention_queue?.items?.[0]; + if (first && state.nativeChildActivity) { + first.project_asset ??= {owner: "codex", gate: "ready", next_action: first.recommended_action ?? "Review the fixture", stop_condition: "Fixture accepted"}; + first.project_asset.native_child_activity = state.nativeChildActivity; + } if (first) { first.waiting_on = "user_or_controller"; const gateDecided = state.decidedGateTodoIds.has("todo-browser-user-gate"); diff --git a/examples/personal-workspace-browser/native-child-activity.mjs b/examples/personal-workspace-browser/native-child-activity.mjs new file mode 100644 index 0000000000..3c13d0ac8e --- /dev/null +++ b/examples/personal-workspace-browser/native-child-activity.mjs @@ -0,0 +1,55 @@ +import assert from "node:assert/strict"; +import {resolve} from "node:path"; +import {outputDir} from "./fixture.mjs"; +import {openWorkspacePage} from "./scenario-context.mjs"; + +export const nativeChildActivityScenario = { + id: "native-child-activity", + async run({browser, collectCoverage, url}) { + const coverageEntries = []; + for (const observation of ["host_observed", "coordinator_reported", "mixed", "unknown"]) { + const context = await openWorkspacePage(browser, url, {collectCoverage, + beforeGoto(api, page) { + api.goalSubagentConfigurationEnabled = true; + page.__loopxRuntime.goalSubagentConfigurations.set("loopx-meta", + {mode: "multi_subagent", spawn_allowed: true, max_children: 3, allowed_domains: []}); + api.nativeChildActivity = {schema_version: "native_subagent_activity_v0", + observation, host_attested: observation === "host_observed", configured_limit: 3, + launched_count: 1, skipped_count: 0, capacity_rejected_count: 0, + host_failed_count: 1, parent_accepted_count: 1, turn_instance_id: "turn-browser-native"}; + }, + }); + try { + const {page} = context; + await page.locator(".personal-goal-link").filter({hasText: "LoopX meta"}).click(); + await page.getByRole("navigation", {name: "Goal 视图"}).getByRole("button", {name: "概览", exact: true}).click(); + await page.getByRole("button", {name: "Goal 信息", exact: true}).click(); + const drawer = page.locator('.personal-context-drawer[data-context-kind="goal"]'); + await drawer.waitFor(); + const activity = drawer.locator(".personal-native-child-activity"); + if (observation === "unknown") { + assert.equal(await activity.count(), 0); + } else { + await activity.waitFor(); + const text = await activity.innerText(); + assert.match(text, /启动 1 次/); + assert.match(text, /主 Agent 验收 1 项/); + assert.match(text, observation === "host_observed" ? /宿主已观察/ + : observation === "mixed" ? /部分决策未经宿主核验/ : /目前没有宿主核验/); + if (observation === "host_observed") { + await activity.scrollIntoViewIfNeeded(); + await page.screenshot({path: resolve(outputDir, "native-child-desktop.png"), animations: "disabled"}); + await page.setViewportSize({width: 390, height: 844}); + await activity.scrollIntoViewIfNeeded(); + assert(await activity.evaluate(el => el.scrollWidth <= el.clientWidth)); + await page.screenshot({path: resolve(outputDir, "native-child-mobile.png"), animations: "disabled"}); + } + } + assert.deepEqual(context.errors, []); + } finally { + coverageEntries.push(...await context.close()); + } + } + return {coverageEntries, note: "Host, coordinator, mixed and unknown native child activity preserve provenance in the packaged Goal drawer."}; + }, +}; diff --git a/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index 1044b845b5..e823a2859d 100644 --- a/loopx/capabilities/multi_subagent/native_child_receipts.py +++ b/loopx/capabilities/multi_subagent/native_child_receipts.py @@ -1,10 +1,4 @@ -"""Turn-bound reports of host-native child-tool decisions. - -The native tool belongs to the host. LoopX can durably reconcile the -coordinator's typed report of its result, but cannot attest that a host call -occurred unless the host itself supplies an integration. This distinction is -part of the projection, not an implicit promise of configured capacity. -""" +"""Turn-bound native child decisions, with explicit report/host provenance.""" from __future__ import annotations @@ -119,13 +113,17 @@ def native_child_activity( attempted = sum(row.get("operation") in {"spawn", "followup"} for row in ordered) rejected = sum(row.get("outcome") == "capacity_rejected" for row in ordered) host_failed = sum(row.get("outcome") == "host_failed" for row in ordered) + sources = {row.get("observation_source") for row in ordered} + observation = ("unknown" if not sources else "host_observed" + if sources == {"host_observed"} else "mixed" + if "host_observed" in sources else "coordinator_reported") return { "schema_version": NATIVE_SUBAGENT_ACTIVITY_SCHEMA_VERSION, "goal_id": goal_id, "agent_id": agent_id, "turn_instance_id": turn_instance_id, "entrypoint_scope": "host_native_child_tools", - "observation": "coordinator_reported" if ordered else "unknown", - "host_attested": False, + "observation": observation, + "host_attested": observation == "host_observed", "configured_limit_kind": "upper_bound", "configured_limit": configured_limit, "observed_capacity": "capacity_rejection_reported" if rejected else @@ -145,12 +143,12 @@ def native_child_activity( } -def load_native_child_activity( +def _load_native_child_events( runtime_root: Path, *, goal_id: str, agent_id: str, - turn_instance_id: str, configured_limit: int, + turn_instance_id: str, goal_ref: Mapping[str, Any] | None = None, registry_path: Path | None = None, -) -> dict[str, Any]: +) -> list[dict[str, Any]]: with quota_accounting_admission( runtime_root=runtime_root, registry_path=registry_path, @@ -183,6 +181,7 @@ def load_native_child_activity( event for event in source if event.get("event_kind") in EVENT_KINDS.values() + and event.get("goal_id") == goal_id and event.get("agent_id") == agent_id and event.get("run_id") == turn_instance_id and ( @@ -191,14 +190,25 @@ def load_native_child_activity( else "goal_ref" not in event ) ] - return native_child_activity( - events, - goal_id=goal_id, - agent_id=agent_id, - turn_instance_id=turn_instance_id, - configured_limit=configured_limit, - goal_ref=goal_ref, - ) + return events + + +def load_native_child_activity( + runtime_root: Path, *, goal_id: str, agent_id: str, + turn_instance_id: str, configured_limit: int, + goal_ref: Mapping[str, Any] | None = None, + registry_path: Path | None = None, +) -> dict[str, Any]: + events = _load_native_child_events( + runtime_root, goal_id=goal_id, agent_id=agent_id, + turn_instance_id=turn_instance_id, goal_ref=goal_ref, + registry_path=registry_path, + ) + return native_child_activity( + events, goal_id=goal_id, agent_id=agent_id, + turn_instance_id=turn_instance_id, configured_limit=configured_limit, + goal_ref=goal_ref, + ) def latest_native_child_activity( @@ -284,6 +294,9 @@ def _record_native_child( registry_path: Path | None = None, goal_ref: Mapping[str, Any] | None = None, source_admission: Mapping[str, Any] | None = None, + _host_observed: bool = False, + _host_child_refs: Sequence[str] | None = None, + _host_wait_ref: str | None = None, ) -> dict[str, Any]: """Preview or append a typed report; never launch a child or spend quota.""" goal_id = _id(goal_id, field="goal_id") @@ -292,10 +305,26 @@ def _record_native_child( operation_id = _id(operation_id, field="operation_id") if isinstance(configured_limit, bool) or not isinstance(configured_limit, int) or configured_limit < 1: raise ValueError("enabled multi_subagent configured_limit must be positive") - fields = _normalized_fields( + fields: dict[str, Any] = _normalized_fields( stage=stage, operation=operation, outcome=outcome, entrypoint_id=entrypoint_id, reason_code=reason_code, evidence_ref=evidence_ref, validation_ref=validation_ref, ) + if _host_observed: + if stage == "review" or operation == "skip": + raise ValueError("host observation cannot attest a parent review or skip") + fields["observation_source"] = "host_observed" + if _host_child_refs is not None: + if (not _host_observed or stage != "decision" or outcome != "started" + or isinstance(_host_child_refs, (str, bytes)) or not _host_child_refs): + raise ValueError("child correlation requires a started host-observed decision") + # Rollout details are scalar; each opaque binding remains exact and + # cannot collide with the decision fields or be text-truncated. + fields.update({"host_child_ref_" + _id(ref, field="host_child_ref"): True + for ref in sorted(set(_host_child_refs))}) + if _host_wait_ref is not None: + if not _host_observed or stage != "result": + raise ValueError("wait correlation requires a host-observed result") + fields["host_wait_ref"] = _id(_host_wait_ref, field="host_wait_ref") log_path = rollout_event_log_path(runtime_root, goal_id) events = load_rollout_events(log_path) prior = _events_for_turn( @@ -307,10 +336,35 @@ def _record_native_child( ) existing = next((event for event in prior if event.get("case_id") == operation_id - and event.get("event_kind") == EVENT_KINDS[stage]), None) + and event.get("event_kind") == EVENT_KINDS[stage] + and (stage != "result" or _host_wait_ref is None + or _details(event).get("host_wait_ref") == _host_wait_ref)), None) + if stage == "result" and _host_wait_ref is None and existing is not None: + # A later decision snapshot can confirm the same typed result without + # replacing the causal wait binding already stored on that result. + wait_ref = _details(existing).get("host_wait_ref") + if wait_ref is not None: + fields["host_wait_ref"] = wait_ref if existing is not None and _details(existing) != fields: raise ValueError("operation identity already has a conflicting native child report") + def validate_result_identity(current: Sequence[Mapping[str, Any]]) -> None: + if stage != "result": + return + for result in current: + if result.get("event_kind") != EVENT_KINDS["result"]: + continue + details = _details(result) + if (_host_wait_ref is not None and details.get("host_wait_ref") == _host_wait_ref + and result.get("case_id") != operation_id): + raise ValueError("wait identity already has a conflicting native child binding") + if (result.get("case_id") == operation_id + and any(details.get(key) != fields.get(key) + for key in ("outcome", "observation_source"))): + raise ValueError("operation identity already has a conflicting native child report") + + validate_result_identity(prior) + def report_admission() -> Mapping[str, Any]: readback = read_heartbeat_settlement( runtime_root, goal_id=goal_id, agent_id=agent_id, todo_id=None, @@ -337,12 +391,15 @@ def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: current = _events_for_turn(observed, goal_id=goal_id, agent_id=agent_id, turn_instance_id=turn_instance_id, goal_ref=goal_ref) + validate_result_identity(current) decisions = {str(item.get("case_id")): item for item in current if item.get("event_kind") == EVENT_KINDS["decision"]} if stage == "decision": if admission["report_permission"] != "new_operation": raise ValueError("native child decision requires an open, work-admitted Turn guard") - if fields["operation"] in {"spawn", "followup"} and any( + # Host observations record calls that already happened, including + # violations; recording one cannot authorize another host call. + if not _host_observed and fields["operation"] in {"spawn", "followup"} and any( _details(item).get("outcome") in {"capacity_rejected", "host_failed"} for item in decisions.values() ): @@ -381,6 +438,7 @@ def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: "run_id", "case_id", *(("goal_ref",) if goal_ref is not None else ()), + *(("details",) if _host_wait_ref is not None else ()), ), precondition=lambda: validate_transition(load_rollout_events(log_path)), ) @@ -414,6 +472,9 @@ def record_native_child( evidence_ref: str | None = None, validation_ref: str | None = None, execute: bool = False, registry_path: Path | None = None, goal_ref: Mapping[str, Any] | None = None, + _host_observed: bool = False, + _host_child_refs: Sequence[str] | None = None, + _host_wait_ref: str | None = None, ) -> dict[str, Any]: """Preview or append one report under its exact quota owner.""" @@ -443,4 +504,7 @@ def record_native_child( registry_path=registry_path, goal_ref=goal_ref, source_admission=source_admission, + _host_observed=_host_observed, + _host_child_refs=_host_child_refs, + _host_wait_ref=_host_wait_ref, ) diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index 97720aff59..917af87fb0 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -933,6 +933,7 @@ def send( attachments: list[dict[str, Any]] | None = None, on_event: Callable[[str, dict[str, Any]], None] | None = None, output_schema: dict[str, Any] | None = None, + on_native_item: Callable[[dict[str, Any]], None] | None = None, ) -> dict[str, Any]: text = " ".join(str(user_message or "").split()) if not text: @@ -1030,6 +1031,10 @@ def send( self.current_turn_id = turn_id method = str(message.get("method") or "") params = message.get("params") + if method == "item/completed" and isinstance(params, dict) and on_native_item: + native_item = params.get("item") + if isinstance(native_item, dict) and native_item.get("type") == "collabAgentToolCall": + on_native_item(native_item) if on_event: phase = { "turn/started": "Agent 已开始处理", diff --git a/loopx/cli_commands/agent_context.py b/loopx/cli_commands/agent_context.py index 1465b050d7..4d96a51027 100644 --- a/loopx/cli_commands/agent_context.py +++ b/loopx/cli_commands/agent_context.py @@ -203,7 +203,8 @@ def handle_agent_context(args, registry_path, runtime_root, print_payload, outpu ) if native_activity is not None: payload["native_child_activity"] = native_activity - payload["host_receipts_scope"] = "turn_bound_coordinator_report" + payload["host_receipts_observed"] = native_activity["host_attested"] + payload["host_receipts_scope"] = "turn_bound_native_child_receipts" print_payload(payload, output_format(args), render_agent_context) return 0 diff --git a/loopx/control_plane/subagent_context.ts b/loopx/control_plane/subagent_context.ts index 72f4d23aae..7ca3df64be 100644 --- a/loopx/control_plane/subagent_context.ts +++ b/loopx/control_plane/subagent_context.ts @@ -109,14 +109,14 @@ function boundedNativeChildActivity(value: unknown): JsonObject | null { if (!source || source.schema_version !== "native_subagent_activity_v0" || source.entrypoint_scope !== "host_native_child_tools") return null; const observation = String(source.observation ?? ""); - if (!["unknown", "coordinator_reported"].includes(observation)) return null; + if (!["unknown", "coordinator_reported", "host_observed", "mixed"].includes(observation)) return null; const count = (key: string) => Number.isInteger(source[key]) && Number(source[key]) >= 0 ? Math.min(Number(source[key]), 10_000) : 0; const result: JsonObject = { schema_version: "native_subagent_activity_v0", entrypoint_scope: "host_native_child_tools", observation, - host_attested: false, + host_attested: observation === "host_observed", configured_limit_kind: "upper_bound", configured_limit: count("configured_limit"), attempted_count: count("attempted_count"), diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index d2e5ab370c..93224f246c 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -15,7 +15,7 @@ from .subagent_execution_topology import ( child_execution_receipts_json_schema, ) -from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES +from ...extensions.codex_native_child import native_child_observer from .codex_sessions import ( CODEX_CLI_SESSION_SCHEMA_VERSION as CODEX_CLI_SESSION_SCHEMA_VERSION, _discard_codex_cli_session, @@ -28,6 +28,7 @@ codex_session_profile_digest, require_codex_session_profile, ) +from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES from .executor import ( HOST_AGENT_VISION_JSON_MAX_CHARS, HOST_REWARD_MEMORY_REFLECTION_JSON_MAX_CHARS, @@ -721,6 +722,15 @@ def commit() -> None: else: goal_admission.accept_result(commit) + child_observer = native_child_observer(request, runtime_root=runtime_root, lineage=lineage, + registry_path=goal_admission.registry_path if goal_admission is not None else None) + invocation_id = "" + if child_observer is not None: + attempt = request.get("host_attempt") + if isinstance(attempt, bool) or not isinstance(attempt, int) or attempt < 1: + raise ValueError("native child CLI observation requires the durable host attempt") + invocation_id = f"exec:{request['turn_key']}:{attempt}" + with tempfile.TemporaryDirectory(prefix="loopx-turn-codex-") as directory: temporary = Path(directory) schema_path = temporary / "result-schema.json" @@ -757,6 +767,16 @@ def observe_event(line: str) -> None: candidate = codex_cli_event_session_id(event) if candidate and candidate not in observed_session: observed_session.append(candidate) + item = event.get("item") + if (child_observer is not None and observed_session + and event.get("type") == "item.completed" and isinstance(item, Mapping)): + def record_child() -> None: + child_observer.observe(item, session_id=observed_session[0], invocation_id=invocation_id) + + if goal_admission is None: + record_child() + else: + goal_admission.accept_result(record_child) structured, diagnostic = _event_failure_categories(event) if structured: structured_failure_categories.add(structured) diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 5135d93d5d..1df534c970 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -37,6 +37,7 @@ _store_codex_cli_session, load_codex_cli_session, ) +from ...extensions.codex_native_child import native_child_observer from .executor import LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION from .host_failure import BuiltInHostError @@ -430,7 +431,22 @@ def on_event(kind: str, event: dict[str, Any]) -> None: recovery_kind="resume_session", ) from exc - return session.send( + child_observer = native_child_observer(request, runtime_root=runtime_root, lineage=lineage, + registry_path=goal_admission.registry_path if goal_admission is not None else None) + + def observe_child(item: Mapping[str, Any]) -> None: + if child_observer is None: + return + def record_child() -> None: + child_observer.observe(item, session_id=session.thread_id, + invocation_id=session.current_turn_id) + + if goal_admission is None: + record_child() + else: + goal_admission.accept_result(record_child) + + result = session.send( _prompt(request) + "\nUse loopx_operation for context/pending/prepare/inspect/consume/report. " "Source conversations are not executor identity. context/pending/inspect do not require consumption. " @@ -446,7 +462,9 @@ def on_event(kind: str, event: dict[str, Any]) -> None: if continuations else ""), output_schema=codex_cli_result_schema(request), on_event=on_event, + **({"on_native_item": observe_child} if child_observer is not None else {}), ) + return result except CodexChatAgentError as exc: raise BuiltInHostError( "codex_operation_host_" + exc.error_code, diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index dc83e8eee1..4c2903f14a 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -161,6 +161,10 @@ def build_loopx_turn_host_request(plan: Mapping[str, Any]) -> dict[str, Any]: if isinstance(reward_memory_recall, Mapping): request["reward_memory_recall"] = dict(reward_memory_recall) request.update(subagent.subagent_host_request_projection(plan)) + from ...extensions.codex_native_child import configured_native_child_limit + + if configured_native_child_limit(request) is not None: + request["turn_instance_id"] = transaction.get("turn_instance_id") or turn_key return request @@ -797,6 +801,11 @@ def _host_result_stage( if "typed_result" not in completed_phases: journal["host_attempt_count"] = int(journal.get("host_attempt_count") or 0) + 1 persist_journal(journal) + from ...extensions.codex_native_child import configured_native_child_limit + if configured_native_child_limit(request) is not None: + # Reuse the journal's durable attempt identity; replay never creates + # another identity, and a real host retry always advances it. + request = {**request, "host_attempt": journal["host_attempt_count"]} # The attempt is durable now, so a later restart must not resume this # reservation. Confirmation failure stops before the host starts. if confirm_start is not None: diff --git a/loopx/extensions/codex_native_child.py b/loopx/extensions/codex_native_child.py new file mode 100644 index 0000000000..1dff414114 --- /dev/null +++ b/loopx/extensions/codex_native_child.py @@ -0,0 +1,170 @@ +"""Transient Codex host events adapted to native-child receipts. + +Only opaque identities and typed outcomes reach the existing multi_subagent +log. Prompts, child messages and raw tool output are never retained. +""" +from __future__ import annotations + +import hashlib +import json +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from ..capabilities.multi_subagent.native_child_receipts import ( + _load_native_child_events, record_native_child, +) + + +def configured_native_child_limit(request: Mapping[str, Any]) -> int | None: + envelope = request.get("turn_envelope") + context = envelope.get("agent_context") if isinstance(envelope, Mapping) else None + contributions = context.get("contributions") if isinstance(context, Mapping) else None + if not isinstance(contributions, list): + return None + for contribution in contributions: + if not isinstance(contribution, Mapping) or contribution.get("capability_id") != "multi_subagent": + continue + facts = contribution.get("facts") + count = facts.get("max_children") if isinstance(facts, Mapping) else None + if isinstance(count, int) and not isinstance(count, bool) and count > 0: + return count + return None + + +class CodexNativeChildObserver: + """Bound to one owned host invocation and one admitted LoopX Turn. + + The public recorder cannot select host provenance. Codex exec and app-server + use different field casing; both are normalized here at the provider seam. + Failed collab items lack a typed capacity error, so they remain host_failed. + A host retry is recorded as observed fact, never authorized by this adapter. + """ + + def __init__(self, *, runtime_root: Path, lineage: Mapping[str, str], + turn_instance_id: str, configured_limit: int, + goal_ref: Mapping[str, Any] | None = None, registry_path: Path | None = None): + self.runtime_root = runtime_root + self.lineage = lineage + self.turn_instance_id = turn_instance_id + self.configured_limit = configured_limit + self.goal_ref = goal_ref + self.registry_path = registry_path + + def _record(self, *, stage: str, **record: Any) -> None: + record_native_child( + runtime_root=self.runtime_root, goal_id=self.lineage["goal_id"], + agent_id=self.lineage["agent_id"], turn_instance_id=self.turn_instance_id, + configured_limit=self.configured_limit, stage=stage, + entrypoint_id="codex_native_tools" if stage == "decision" else None, + execute=True, _host_observed=True, goal_ref=self.goal_ref, + registry_path=self.registry_path, **record, + ) + + def _restore_operation(self, child: str, *, wait_ref: str) -> str | None: + # Restore the latest successful host decision for this opaque child, + # including followups and rows outside the presentation window. The + # existing log supplies order and exact Goal/agent/Turn ownership. + child_ref = "codex-child-" + hashlib.sha256(child.encode()).hexdigest()[:32] + legacy_spawn = "codex-" + hashlib.sha256(child.encode()).hexdigest()[:32] + events = _load_native_child_events( + self.runtime_root, goal_id=self.lineage["goal_id"], agent_id=self.lineage["agent_id"], + turn_instance_id=self.turn_instance_id, goal_ref=self.goal_ref, + registry_path=self.registry_path, + ) + # The first terminal observation owns this wait identity permanently. + # Replayed waits must never reinterpret a later child association. + for event in events: + details = event.get("details") + if (event.get("event_kind") == "native_child_result" + and isinstance(details, Mapping) + and details.get("observation_source") == "host_observed" + and details.get("host_wait_ref") == wait_ref): + return str(event["case_id"]) + for event in reversed(events): + details = event.get("details") + if not isinstance(details, Mapping): + continue + if (event.get("event_kind") == "native_child_decision" + and details.get("operation") in {"spawn", "followup"} + and details.get("outcome") == "started" + and details.get("observation_source") == "host_observed" + and details.get("entrypoint_id") == "codex_native_tools"): + if (details.get("host_child_ref_" + child_ref) is True + or (details.get("operation") == "spawn" + and event.get("case_id") == legacy_spawn)): + return str(event["case_id"]) + return None + + def observe(self, item: Mapping[str, Any], *, session_id: str, invocation_id: str) -> None: + item_type = item.get("type") + if item_type not in {"collab_tool_call", "collabAgentToolCall"}: + return + snake = item_type == "collab_tool_call" + sender = item.get("sender_thread_id" if snake else "senderThreadId") + status = item.get("status") + native_id = item.get("id") + if sender != session_id or not isinstance(native_id, str) or not native_id or status not in {"completed", "failed"}: + return + tool = {"spawn_agent": "spawn", "spawnAgent": "spawn", "send_input": "followup", + "sendInput": "followup", "resumeAgent": "followup", "wait": "wait"}.get(item.get("tool")) + if tool is None: + return + receivers = item.get("receiver_thread_ids" if snake else "receiverThreadIds") + decision_id: str | None = None + if tool != "wait": + started = status == "completed" and isinstance(receivers, list) and bool(receivers) + if started and any(not isinstance(child, str) or not child for child in receivers): + return + # Successful spawn identity survives host event replay/restart; host + # exec display-item counters alone are not globally unique. + if not invocation_id: + raise ValueError("native child decision requires its owned host invocation") + identity = (receivers[0] if tool == "spawn" and started else + json.dumps([session_id, invocation_id, native_id], separators=(",", ":"))) + operation_id = "codex-" + hashlib.sha256(identity.encode()).hexdigest()[:32] + self._record(stage="decision", operation_id=operation_id, operation=tool, + outcome="started" if started else "host_failed", + **({"reason_code": "host_failed"} if not started else { + "_host_child_refs": sorted({"codex-child-" + hashlib.sha256(child.encode()).hexdigest()[:32] + for child in receivers})})) + if started: + decision_id = operation_id + states = item.get("agents_states" if snake else "agentsStates") + if not isinstance(states, Mapping): + return + for child, state in states.items(): + if not isinstance(child, str) or not child or not isinstance(state, Mapping): + continue + # A spawn/followup snapshot belongs to that stable decision, even + # when replayed after later work. A wait first restores its own + # consumed binding, then falls back to the latest child operation. + if tool != "wait" and (not isinstance(receivers, list) or child not in receivers): + continue + outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) + if outcome is None: + continue + wait_ref = None + if tool == "wait": + if not invocation_id: + raise ValueError("native child wait requires its owned host invocation") + identity = json.dumps([session_id, invocation_id, native_id, child], separators=(",", ":")) + wait_ref = "codex-wait-" + hashlib.sha256(identity.encode()).hexdigest()[:32] + operation_id = self._restore_operation(child, wait_ref=wait_ref) if wait_ref else decision_id + if operation_id is not None: + self._record(stage="result", operation_id=operation_id, outcome=outcome, + _host_wait_ref=wait_ref) + + +def native_child_observer(request: Mapping[str, Any], *, runtime_root: Path, + lineage: Mapping[str, str], + registry_path: Path | None = None) -> CodexNativeChildObserver | None: + limit = configured_native_child_limit(request) + turn = request.get("turn_instance_id") + if limit is None or not isinstance(turn, str) or not turn: + return None + goal_ref = request.get("goal_ref") + return CodexNativeChildObserver(runtime_root=runtime_root, lineage=lineage, + turn_instance_id=turn, configured_limit=limit, + goal_ref=goal_ref if isinstance(goal_ref, Mapping) else None, + registry_path=registry_path) diff --git a/loopx/presentation/renderers/status_markdown.py b/loopx/presentation/renderers/status_markdown.py index 7ab021b9a5..768a5588c0 100644 --- a/loopx/presentation/renderers/status_markdown.py +++ b/loopx/presentation/renderers/status_markdown.py @@ -1275,11 +1275,12 @@ def _append_project_asset_runtime_policy_markdown( if isinstance(project_asset.get("native_child_activity"), dict) else {} ) - if native_child_activity.get("observation") == "coordinator_reported": + if native_child_activity.get("observation") in {"coordinator_reported", "host_observed", "mixed"}: lines.append( " - native_child_activity: " f"turn={markdown_scalar(native_child_activity.get('turn_instance_id'))} " - "source=coordinator_reported host_attested=false " + f"source={native_child_activity.get('observation')} " + f"host_attested={str(native_child_activity.get('host_attested')).lower()} " f"configured_max={native_child_activity.get('configured_limit')} " f"starts={native_child_activity.get('launched_count')} " f"skips={native_child_activity.get('skipped_count')} " diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py new file mode 100644 index 0000000000..178f866434 --- /dev/null +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -0,0 +1,431 @@ +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from loopx.extensions.codex_native_child import ( + configured_native_child_limit, CodexNativeChildObserver, native_child_observer, +) +from loopx.capabilities.multi_subagent.native_child_receipts import load_native_child_activity, record_native_child +from loopx.rollout_event_log import load_rollout_events, rollout_event_log_path +from tests.capabilities.test_native_child_receipts import _admit, GOAL, AGENT, TURN + + +def _turn(): + return {"id": "host-turn-1", "itemsView": "full", "status": "completed", "items": [ + {"type": "collabAgentToolCall", "id": "call-1", "senderThreadId": "parent-1", + "tool": "spawnAgent", "status": "completed", "receiverThreadIds": ["child-1"], + "prompt": "private child instructions", "agentsStates": {"child-1": {"status": "running"}}}, + {"type": "collabAgentToolCall", "id": "wait-1", "senderThreadId": "parent-1", + "tool": "wait", "status": "completed", "agentsStates": { + "child-1": {"status": "completed", "message": "private child result"}}}, + ]} + + +def _observe(root, items, session_id="parent-1"): + observer = CodexNativeChildObserver(runtime_root=root, + lineage={"goal_id": GOAL, "agent_id": AGENT}, + turn_instance_id=TURN, configured_limit=3) + for item in items: + observer.observe(item, session_id=session_id, invocation_id="host-turn-1") + + +def test_host_spawn_result_parent_review_and_restart(tmp_path: Path): + _admit(tmp_path) + _observe(tmp_path, _turn()["items"]) + _observe(tmp_path, _turn()["items"]) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["observation"] == "host_observed" + assert activity["host_attested"] is True + assert activity["launched_count"] == 1 + assert activity["parent_accepted_count"] == 0 + [operation] = activity["operations"] + assert operation["result"] == "completed" + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id=operation["operation_id"], + stage="review", outcome="accepted", evidence_ref="evidence-1", + validation_ref="validation-1", execute=True) + readback = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert readback["parent_accepted_count"] == 1 + assert readback["host_attested"] is True + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert len(events) == 4 + assert "private child" not in json.dumps(events) + assert all(event.get("event_kind") != "quota_spend" for event in events) + + +@pytest.mark.parametrize("items", [[], [{"type": "agentMessage", "text": "I spawned three children"}], + [{**_turn()["items"][0], "status": "inProgress"}], + [{**_turn()["items"][0], "senderThreadId": "historical-parent"}]]) +def test_missing_or_unrelated_host_events_stay_unknown(tmp_path: Path, items): + _admit(tmp_path) + _observe(tmp_path, items) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["observation"] == "unknown" + assert activity["launched_count"] == 0 + + +def test_exec_casing_and_display_counter_replay_do_not_duplicate_spawn(tmp_path: Path): + _admit(tmp_path) + snake = {"type": "collab_tool_call", "id": "item_1", "sender_thread_id": "parent-1", + "tool": "spawn_agent", "status": "completed", "receiver_thread_ids": ["child-1"], + "agents_states": {"child-1": {"status": "completed", "message": "private content"}}} + _observe(tmp_path, [snake, {**snake, "id": "item_7"}]) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["host_attested"] is True + assert activity["launched_count"] == 1 + assert activity["operations"][0]["result"] == "completed" + assert len(load_rollout_events(rollout_event_log_path(tmp_path, GOAL))) == 3 + + +def test_host_failure_does_not_infer_capacity_from_prose(tmp_path: Path): + _admit(tmp_path) + failed = {**_turn()["items"][0], "status": "failed", "receiverThreadIds": [], + "agentsStates": {"child-1": {"status": "errored", "message": "agent_thread_limit_reached"}}} + _observe(tmp_path, [failed]) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["host_attested"] is True + assert activity["host_failed_count"] == 1 + assert activity["capacity_rejected_count"] == 0 + assert activity["retry_same_turn"] is False + + +def test_feature_off_has_no_observer_and_reports_cannot_attest_reviews(tmp_path: Path): + assert native_child_observer({"turn_envelope": {}}, runtime_root=tmp_path, + lineage={"goal_id": GOAL, "agent_id": AGENT}) is None + assert configured_native_child_limit({"turn_envelope": {}}) is None + assert configured_native_child_limit({"turn_envelope": {"agent_context": { + "contributions": [{"capability_id": "multi_subagent", "facts": {"max_children": 0}}]}}}) is None + with pytest.raises(ValueError, match="cannot attest"): + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id="review-1", stage="review", + outcome="accepted", evidence_ref="evidence-1", validation_ref="validation-1", + execute=True, _host_observed=True) + + +def test_cli_host_collects_native_items_before_returning_parent_result(tmp_path: Path, monkeypatch): + import sys + from loopx.control_plane.turn_driver import codex_cli + from tests.test_loopx_turn_codex_cli import _request + + _admit(tmp_path) + request = _request() + request["turn_instance_id"] = TURN + request["host_attempt"] = 1 + request["turn_envelope"].update(goal_id=GOAL, agent_id=AGENT, agent_context={ + "contributions": [{"capability_id": "multi_subagent", "facts": {"max_children": 3}}]}) + request["turn_envelope"]["action"]["selected_todo"]["todo_id"] = "todo_native_1" + + def host(command, **kwargs): + kwargs["on_stdout"](json.dumps({"type": "thread.started", "thread_id": "parent-1"}) + "\n") + for item in _turn()["items"]: + kwargs["on_stdout"](json.dumps({"type": "item.completed", "item": item}) + "\n") + Path(command[command.index("--output-last-message") + 1]).write_text(json.dumps({"parent_work": "preserved"})) + return {"returncode": 0, "outcome": "exited", "output_complete": True} + + monkeypatch.setattr(codex_cli, "run_host_process", host) + assert codex_cli.run_codex_cli_host(request, runtime_root=tmp_path, project=tmp_path, + codex_bin=sys.executable) == {"parent_work": "preserved"} + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["host_attested"] is True + assert activity["launched_count"] == 1 + assert activity["operations"][0]["result"] == "completed" + + +def _real_cli_calls(tmp_path, monkeypatch, batches, *, host_attempts=None): + import sys + from loopx.control_plane.turn_driver import codex_cli + from tests.test_loopx_turn_codex_cli import _request + + script = tmp_path / "host.py" + script.write_text(""" +import json, sys +from pathlib import Path +sys.stdin.read() +print(json.dumps({"type": "thread.started", "thread_id": "parent-1"}), flush=True) +for item in json.loads(Path(sys.argv[1]).read_text(encoding="utf-8")): + print(json.dumps({"type": "item.completed", "item": item}), flush=True) +Path(sys.argv[2]).write_text(json.dumps({"parent_work": "preserved"}), encoding="utf-8") +""", encoding="utf-8") + event_file = tmp_path / "events.json" + sessions = [] + + def command(**kwargs): + sessions.append(kwargs["session_id"]) + return [sys.executable, str(script), str(event_file), str(kwargs["output_path"])] + + # Replace only executable selection; run the real process, stream parser, + # session store, receipt admission and durable readback. + monkeypatch.setattr(codex_cli, "_codex_command", command) + for attempt, items in enumerate(batches, 1): + request = _request(session_action="start_new" if attempt == 1 else "resume") + request["turn_instance_id"] = TURN + request["host_attempt"] = host_attempts[attempt - 1] if host_attempts else attempt + request["turn_envelope"].update(goal_id=GOAL, agent_id=AGENT, agent_context={ + "contributions": [{"capability_id": "multi_subagent", "facts": {"max_children": 3}}]}) + request["turn_envelope"]["action"]["selected_todo"]["todo_id"] = "todo_native_1" + event_file.write_text(json.dumps(items), encoding="utf-8") + assert codex_cli.run_codex_cli_host(request, runtime_root=tmp_path, project=tmp_path, + codex_bin=sys.executable) == {"parent_work": "preserved"} + assert sessions == [None] + ["parent-1"] * (len(batches) - 1) + return load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + + +def test_real_cli_resume_wait_only_completes_original_spawn(tmp_path, monkeypatch): + _admit(tmp_path) + spawn, wait = _turn()["items"] + activity = _real_cli_calls(tmp_path, monkeypatch, [[spawn], [wait, wait]]) + assert activity["launched_count"] == 1 + assert activity["operation_count"] == 1 + assert activity["operations"][0]["result"] == "completed" + assert activity["quota_spend_slots"] == 0 + assert len(load_rollout_events(rollout_event_log_path(tmp_path, GOAL))) == 3 + + +def test_real_cli_resume_counter_reuse_and_exact_replay(tmp_path, monkeypatch): + _admit(tmp_path) + followup = {"type": "collab_tool_call", "id": "item_0", "sender_thread_id": "parent-1", + "tool": "send_input", "status": "completed", "receiver_thread_ids": ["child-1"]} + other_child = {**followup, "receiver_thread_ids": ["child-2"]} + failed = {**followup, "tool": "spawn_agent", "status": "failed", "receiver_thread_ids": []} + activity = _real_cli_calls(tmp_path, monkeypatch, [ + [followup, followup], [other_child, other_child], [followup, followup], + [failed, failed], [failed, failed]]) + assert activity["operation_count"] == 5 + assert activity["attempted_count"] == 5 + assert activity["host_failed_count"] == 2 + assert activity["launched_count"] == 0 + assert activity["quota_spend_slots"] == 0 + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert len(events) == 6 + assert "private child" not in json.dumps(events) + + +def test_resume_restores_spawn_outside_the_visible_window(tmp_path): + _admit(tmp_path) + spawn, wait = _turn()["items"] + items = [{**spawn, "receiverThreadIds": [f"child-{i}"], "agentsStates": {}} + for i in range(1, 11)] + _observe(tmp_path, items) + _observe(tmp_path, [wait]) + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert sum(event["event_kind"] == "native_child_decision" for event in events) == 10 + assert sum(event["event_kind"] == "native_child_result" for event in events) == 1 + + +def test_resume_does_not_attest_coordinator_reported_or_unknown_children(tmp_path): + import hashlib + _admit(tmp_path) + operation = "codex-" + hashlib.sha256(b"child-1").hexdigest()[:32] + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id=operation, stage="decision", + operation="spawn", outcome="started", entrypoint_id="codex_native_tools", execute=True) + _observe(tmp_path, [_turn()["items"][1]]) + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert not any(event["event_kind"] == "native_child_result" for event in events) + + +def test_restart_replay_uses_the_original_invocation_binding(tmp_path): + _admit(tmp_path) + followup = {**_turn()["items"][0], "tool": "sendInput", "agentsStates": {}} + _observe(tmp_path, [followup]) + _observe(tmp_path, [followup]) + assert len(load_rollout_events(rollout_event_log_path(tmp_path, GOAL))) == 2 + + +@pytest.mark.parametrize('with_spawn', [False, True]) +@pytest.mark.parametrize('terminal_status,expected', [('completed', 'completed'), ('errored', 'failed')]) +def test_real_cli_resumed_followup_owns_its_result(tmp_path, monkeypatch, with_spawn, terminal_status, expected): + _admit(tmp_path) + spawn, wait = _turn()['items'] + followup = {**spawn, 'id': 'item_0', 'tool': 'sendInput', 'agentsStates': {}} + terminal = {**wait, 'agentsStates': {'child-1': {'status': terminal_status}}} + batches = ([[spawn, wait]] if with_spawn else []) + [[followup], [terminal, terminal], [terminal]] + activity = _real_cli_calls(tmp_path, monkeypatch, batches, + host_attempts=[*range(1, len(batches)), len(batches) - 1]) + followups = [row for row in activity['operations'] if row['operation'] == 'followup'] + assert len(followups) == 1 and followups[0]['result'] == expected + assert activity['operation_count'] == 1 + int(with_spawn) + assert activity['launched_count'] == int(with_spawn) + assert activity['quota_spend_slots'] == 0 + if with_spawn: + original = next(row for row in activity['operations'] if row['operation'] == 'spawn') + assert original['result'] == 'completed' + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert sum(event['event_kind'] == 'native_child_result' for event in events) == 1 + int(with_spawn) + assert '"child-1"' not in json.dumps(events) and 'private child' not in json.dumps(events) + + +def test_real_cli_consecutive_followups_restore_the_latest_owned_operation(tmp_path, monkeypatch): + _admit(tmp_path) + spawn, wait = _turn()['items'] + followup = {**spawn, 'id': 'item_0', 'tool': 'sendInput', 'agentsStates': {}} + failed = {**wait, 'agentsStates': {'child-1': {'status': 'errored'}}} + activity = _real_cli_calls(tmp_path, monkeypatch, + [[spawn, wait], [followup], [wait], [followup], [failed, failed], [failed]]) + followups = [row for row in activity['operations'] if row['operation'] == 'followup'] + assert len(followups) == 2 + assert {row['result'] for row in followups} == {'completed', 'failed'} + assert activity['operation_count'] == 3 and activity['launched_count'] == 1 + assert activity['quota_spend_slots'] == 0 + + +@pytest.mark.parametrize('host_observed,stage,outcome', [(False, 'decision', 'started'), (True, 'result', 'completed')]) +def test_child_correlation_cannot_attest_a_report_or_result(tmp_path, host_observed, stage, outcome): + _admit(tmp_path) + with pytest.raises(ValueError, match='child correlation requires'): + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id='correlation-1', + stage=stage, outcome=outcome, operation='followup' if stage == 'decision' else None, + entrypoint_id='codex_native_tools' if stage == 'decision' else None, + execute=True, _host_observed=host_observed, _host_child_refs=['codex-child-opaque']) + assert not any(event['event_kind'] == 'native_child_decision' + for event in load_rollout_events(rollout_event_log_path(tmp_path, GOAL))) + + + +@pytest.mark.parametrize("snake", [False, True]) +@pytest.mark.parametrize("new_wait", [False, True]) +def test_real_cli_old_spawn_snapshot_cannot_complete_a_later_followup( + tmp_path, monkeypatch, snake, new_wait, +): + _admit(tmp_path) + spawn, wait = _turn()["items"] + spawn = {**spawn, "agentsStates": {"child-1": {"status": "completed"}}} + followup = {**spawn, "id": "item_0", "tool": "sendInput", + "agentsStates": {"child-1": {"status": "running"}}} + if snake: + def exec_item(item): + return {"type": "collab_tool_call", "id": item["id"], + "sender_thread_id": item["senderThreadId"], + "tool": {"spawnAgent": "spawn_agent", "sendInput": "send_input", "wait": "wait"}[item["tool"]], + "status": item["status"], "receiver_thread_ids": item.get("receiverThreadIds", []), + "agents_states": item["agentsStates"]} + spawn, followup, wait = map(exec_item, (spawn, followup, wait)) + # The third invocation replays only the old spawn's terminal snapshot. It + # contains no new wait or followup result and cannot prove future work done. + batches = [[spawn], [followup], [spawn]] + ([[wait, wait]] if new_wait else []) + activity = _real_cli_calls(tmp_path, monkeypatch, batches) + assert activity["operation_count"] == 2 and activity["launched_count"] == 1 + original = next(row for row in activity["operations"] if row["operation"] == "spawn") + later = next(row for row in activity["operations"] if row["operation"] == "followup") + assert original["result"] == "completed" + assert later.get("result") == ("completed" if new_wait else None) + assert activity["parent_accepted_count"] == 0 and activity["quota_spend_slots"] == 0 + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert sum(row["event_kind"] == "native_child_result" for row in events) == 1 + int(new_wait) + + +@pytest.mark.parametrize("tool", ["spawnAgent", "sendInput"]) +def test_nonwait_terminal_snapshot_ignores_unrelated_receivers(tmp_path, tool): + _admit(tmp_path) + spawn = _turn()["items"][0] + other = {**spawn, "receiverThreadIds": ["child-2"], "agentsStates": {}} + unrelated_snapshot = {**spawn, "tool": tool, + "agentsStates": {"child-2": {"status": "completed"}}} + _observe(tmp_path, [other, unrelated_snapshot]) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert activity["operation_count"] == 2 + assert all(row.get("result") is None for row in activity["operations"]) + + +@pytest.mark.parametrize("snake", [False, True]) +@pytest.mark.parametrize("restart", [False, True]) +def test_real_cli_consumed_wait_replay_cannot_complete_a_later_followup( + tmp_path, monkeypatch, snake, restart, +): + _admit(tmp_path) + spawn, wait = _turn()["items"] + followup = {**spawn, "id": "followup-1", "tool": "sendInput", + "agentsStates": {"child-1": {"status": "running"}}} + if snake: + def exec_item(item): + return {"type": "collab_tool_call", "id": item["id"], + "sender_thread_id": item["senderThreadId"], + "tool": {"spawnAgent": "spawn_agent", "sendInput": "send_input", "wait": "wait"}[item["tool"]], + "status": item["status"], "receiver_thread_ids": item.get("receiverThreadIds", []), + "agents_states": item["agentsStates"]} + spawn, followup, wait = map(exec_item, (spawn, followup, wait)) + # Distinct native IDs within one invocation exclude counter reuse. Restart + # must preserve the original attempt for an exact replay of the old wait. + batches = [[spawn, wait], [followup], [wait]] if restart else [[spawn, wait, followup, wait]] + activity = _real_cli_calls(tmp_path, monkeypatch, batches, + host_attempts=[1, 2, 1] if restart else [1]) + original = next(row for row in activity["operations"] if row["operation"] == "spawn") + later = next(row for row in activity["operations"] if row["operation"] == "followup") + assert original["result"] == "completed" + assert later.get("result") is None + assert activity["launched_count"] == 1 and activity["operation_count"] == 2 + assert activity["parent_accepted_count"] == 0 and activity["quota_spend_slots"] == 0 + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + assert sum(event["event_kind"] == "native_child_result" for event in events) == 1 + assert '"child-1"' not in json.dumps(events) and "private child" not in json.dumps(events) + + +@pytest.mark.parametrize("spawn_completed", [False, True]) +def test_distinct_waits_keep_their_first_binding_and_reject_changed_outcomes(tmp_path, spawn_completed): + _admit(tmp_path) + spawn, wait = _turn()["items"] + if spawn_completed: + spawn = {**spawn, "agentsStates": {"child-1": {"status": "completed"}}} + followup = {**spawn, "id": "followup-1", "tool": "sendInput", "agentsStates": {}} + fresh_wait = {**wait, "id": "wait-2"} + snapshot = {**spawn, "agentsStates": {"child-1": {"status": "completed"}}} + _observe(tmp_path, [spawn, wait, followup, snapshot, wait]) + pending = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + later = next(row for row in pending["operations"] if row["operation"] == "followup") + assert later.get("result") is None + _observe(tmp_path, [fresh_wait, fresh_wait]) + completed = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert all(row["result"] == "completed" for row in completed["operations"]) + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + with pytest.raises(ValueError, match="conflicting"): + _observe(tmp_path, [{**wait, "agentsStates": {"child-1": {"status": "errored"}}}]) + assert load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) == events + assert completed["parent_accepted_count"] == 0 and completed["quota_spend_slots"] == 0 + + +def test_wait_binding_is_per_child_and_cannot_be_reassigned(tmp_path): + _admit(tmp_path) + spawn, wait = _turn()["items"] + other = {**spawn, "id": "spawn-2", "receiverThreadIds": ["child-2"], "agentsStates": {}} + both = {**wait, "agentsStates": {"child-1": {"status": "completed"}, + "child-2": {"status": "shutdown"}}} + followup = {**spawn, "id": "followup-1", "tool": "sendInput", "agentsStates": {}} + _observe(tmp_path, [spawn, other, both, followup, both]) + activity = load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3) + assert [row.get("result") for row in activity["operations"] if row["operation"] == "followup"] == [None] + events = load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) + result = next(row for row in events if row["event_kind"] == "native_child_result") + later = next(row for row in activity["operations"] if row["operation"] == "followup") + with pytest.raises(ValueError, match="conflicting native child binding"): + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id=later["operation_id"], + stage="result", outcome="completed", execute=True, _host_observed=True, + _host_wait_ref=result["details"]["host_wait_ref"]) + assert load_rollout_events(rollout_event_log_path(tmp_path, GOAL)) == events + + +@pytest.mark.parametrize("host_observed,stage", [(False, "result"), (True, "decision")]) +def test_wait_correlation_requires_an_observed_result(tmp_path, host_observed, stage): + _admit(tmp_path) + with pytest.raises(ValueError, match="wait correlation requires"): + record_native_child(runtime_root=tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=3, operation_id="op-1", stage=stage, + outcome="completed" if stage == "result" else "started", + operation="spawn" if stage == "decision" else None, + entrypoint_id="codex_native_tools" if stage == "decision" else None, + execute=True, _host_observed=host_observed, _host_wait_ref="codex-wait-known") diff --git a/tests/control_plane_ts/agent_context.test.ts b/tests/control_plane_ts/agent_context.test.ts index 324dfbb445..36b385f214 100644 --- a/tests/control_plane_ts/agent_context.test.ts +++ b/tests/control_plane_ts/agent_context.test.ts @@ -310,3 +310,20 @@ test("durable native child report is bounded and does not claim host attestation assert.equal(facts.native_receipt_observation, "coordinator_reported"); assert.ok(!JSON.stringify(packet).includes("private result")); }); + + +test("host receipt provenance survives projection without raw child content", () => { + for (const observation of ["host_observed", "mixed"]) { + const packet = evaluateSubagentContext({ phase: "after_delegate_result", scope, + orchestration: policy, observations: { native_child_activity: { + schema_version: "native_subagent_activity_v0", entrypoint_scope: "host_native_child_tools", + observation, host_attested: true, configured_limit: 3, launched_count: 1, + attempted_count: 1, parent_accepted_count: 1, raw_host_result: "private child result", + } } })!; + const facts = (packet.contributions as Record[])[0].facts; + assert.equal(facts.native_child_activity.observation, observation); + assert.equal(facts.native_child_activity.host_attested, observation === "host_observed"); + assert.equal(facts.native_child_activity.parent_accepted_count, 1); + assert.ok(!JSON.stringify(packet).includes("private child result")); + } +}); diff --git a/tests/test_chat_agent.py b/tests/test_chat_agent.py index 5ff1be8ddf..b918a7182b 100644 --- a/tests/test_chat_agent.py +++ b/tests/test_chat_agent.py @@ -610,3 +610,21 @@ def test_retry_and_unrelated_policy_events_do_not_terminate_current_turn( assert result["message"] == "Recovered." assert any(k == "agent.phase" and p["label"] == "Codex 正在重试" for k, p in events) assert sum(k == "answer.final" for k, p in events) == 1 + + +def test_native_child_callback_is_scoped_to_owned_thread_and_turn(monkeypatch, tmp_path): + session = chat_agent.CodexChatAgentSession(process=_FakeAppServerProcess(), + messages=queue.Queue(), thread_id="thread-fixture", work_dir=tmp_path) + item = {"type": "collabAgentToolCall", "id": "call-1", "tool": "spawnAgent"} + events = iter([ + {"method": "item/completed", "params": {"threadId": "other", "turnId": "turn-fixture", "item": item}}, + {"method": "item/completed", "params": {"threadId": "thread-fixture", "turnId": "old-turn", "item": item}}, + {"method": "item/completed", "params": {"threadId": "thread-fixture", "turnId": "turn-fixture", "item": item}}, + {"method": "item/agentMessage/delta", "params": {"delta": "Ready."}}, + {"method": "turn/completed", "params": {"turn": {"status": "completed"}}}, + ]) + monkeypatch.setattr(session, "_request", lambda *a, **kw: {"turn": {"id": "turn-fixture"}}) + monkeypatch.setattr(session, "_next_event", lambda **kw: next(events)) + observed = [] + session.send("Reply briefly.", on_native_item=observed.append) + assert observed == [item] diff --git a/tests/test_loopx_turn_executor.py b/tests/test_loopx_turn_executor.py index 76b7ddd6c7..4fd2bae845 100644 --- a/tests/test_loopx_turn_executor.py +++ b/tests/test_loopx_turn_executor.py @@ -1574,16 +1574,25 @@ def host(_request: dict[str, object]) -> dict[str, object]: assert calls == {"host": 1, "writeback": 0, "spend": 0, "scheduler": 0} +@pytest.mark.parametrize("native_children", [False, True]) def test_run_once_resumes_session_observed_by_recoverable_failed_turn( - tmp_path: Path, + tmp_path: Path, native_children: bool, ) -> None: plan = _codex_plan() + if native_children: + plan["turn_envelope"]["agent_context"] = {"contributions": [ + {"capability_id": "multi_subagent", "facts": {"max_children": 3}}]} calls = {"host": 0, "writeback": 0, "spend": 0, "scheduler": 0} session_actions: list[str] = [] writeback, spend, scheduler = _callbacks(calls) def host(request: dict[str, object]) -> dict[str, object]: calls["host"] += 1 + if native_children: + assert request["host_attempt"] == calls["host"] + assert request["host_attempt"] == _journal(tmp_path / "runtime")["host_attempt_count"] + else: + assert "host_attempt" not in request session = request["session"] assert isinstance(session, dict) session_actions.append(str(session["action"])) @@ -1638,6 +1647,9 @@ def session_binding( assert recovered["recovery"]["planned"] == inspected["recovery_decision"] assert recovered["status"] == "committed" assert session_actions == ["start_new", "resume"] + replay = run_loopx_turn_once(plan, **common) + assert replay["replayed"] is True + assert _journal(tmp_path / "runtime")["host_attempt_count"] == 2 assert calls == {"host": 2, "writeback": 1, "spend": 1, "scheduler": 1}