From dbefcf741d35d7bff9fbfadfe90dc815f4d8dc78 Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Fri, 2 Oct 2026 16:08:53 +0800 Subject: [PATCH 01/12] fix(subagents): observe native Codex child events on owned Turns Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../presentation/dashboard/src/data/status.ts | 4 +- .../personal-workspace/context-drawer.tsx | 6 +- .../src/features/personal-workspace/i18n.tsx | 4 + .../personal-workspace-model.ts | 4 +- .../dashboard/src/views/dashboard-page.tsx | 4 +- .../multi_subagent/native_child_receipts.py | 25 +++-- loopx/chat_agent.py | 5 + loopx/cli_commands/agent_context.py | 3 +- loopx/control_plane/subagent_context.ts | 4 +- loopx/control_plane/turn_driver/codex_cli.py | 15 ++- .../turn_driver/codex_native_child.py | 105 ++++++++++++++++++ .../turn_driver/codex_operation_host.py | 18 ++- loopx/control_plane/turn_driver/executor.py | 4 + .../presentation/renderers/status_markdown.py | 5 +- 14 files changed, 181 insertions(+), 25 deletions(-) create mode 100644 loopx/control_plane/turn_driver/codex_native_child.py 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/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index 1044b845b5..707247c364 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 @@ -284,6 +282,7 @@ 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, ) -> dict[str, Any]: """Preview or append a typed report; never launch a child or spend quota.""" goal_id = _id(goal_id, field="goal_id") @@ -296,6 +295,10 @@ def _record_native_child( 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" log_path = rollout_event_log_path(runtime_root, goal_id) events = load_rollout_events(log_path) prior = _events_for_turn( @@ -342,7 +345,9 @@ def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: 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() ): 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..7c586dc11d 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -15,7 +15,8 @@ from .subagent_execution_topology import ( child_execution_receipts_json_schema, ) -from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES +from .codex_native_child import native_child_observer +from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES, selected_turn_todo from .codex_sessions import ( CODEX_CLI_SESSION_SCHEMA_VERSION as CODEX_CLI_SESSION_SCHEMA_VERSION, _discard_codex_cli_session, @@ -721,6 +722,8 @@ def commit() -> None: else: goal_admission.accept_result(commit) + child_observer = native_child_observer(request, runtime_root=runtime_root, lineage=lineage) + with tempfile.TemporaryDirectory(prefix="loopx-turn-codex-") as directory: temporary = Path(directory) schema_path = temporary / "result-schema.json" @@ -757,6 +760,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]) + + 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_native_child.py b/loopx/control_plane/turn_driver/codex_native_child.py new file mode 100644 index 0000000000..6b7c219f65 --- /dev/null +++ b/loopx/control_plane/turn_driver/codex_native_child.py @@ -0,0 +1,105 @@ +"""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 +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +from ...capabilities.multi_subagent.native_child_receipts import 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): + self.runtime_root = runtime_root + self.lineage = lineage + self.turn_instance_id = turn_instance_id + self.configured_limit = configured_limit + self.children: dict[str, str] = {} + + def _record(self, *, stage: str, **record: str) -> 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, **record, + ) + + def observe(self, item: Mapping[str, Any], *, session_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") + 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. + identity = receivers[0] if tool == "spawn" and started else session_id + ":" + native_id + 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 {})) + if started: + for child in receivers: + self.children[child] = operation_id + states = item.get("agents_states" if snake else "agentsStates") + if not isinstance(states, Mapping): + return + for child, state in states.items(): + if child not in self.children or not isinstance(state, Mapping): + continue + outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) + if outcome: + self._record(stage="result", operation_id=self.children[child], outcome=outcome) + + +def native_child_observer(request: Mapping[str, Any], *, runtime_root: Path, + lineage: Mapping[str, str]) -> 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 + return CodexNativeChildObserver(runtime_root=runtime_root, lineage=lineage, + turn_instance_id=turn, configured_limit=limit) diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 5135d93d5d..34e2810a18 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 .codex_native_child import native_child_observer from .executor import LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION from .host_failure import BuiltInHostError @@ -430,7 +431,20 @@ 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) + + 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) + + 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 +460,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..59d504f484 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 .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 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')} " From 49c2c32f73f8405217a4ef6fd653e20cae1d1e58 Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Fri, 2 Oct 2026 16:08:53 +0800 Subject: [PATCH 02/12] docs(subagents): explain native host provenance and protocol limits Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../host-native-child-receipts.md | 62 ++++++++++++++----- 1 file changed, 48 insertions(+), 14 deletions(-) diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index b30eeaf878..120218ed09 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,25 @@ 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 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 时第二次启动 +失败、首个子任务完成以及独立的父任务验收。持久回读保留一个启动和一个 +采纳结果。失败子类型仍有宿主协议缺口,不能声称已核验容量错误码。 From 50b6b8c3022a5f9bbb9d1397e3413e39b3b5d7df Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Fri, 2 Oct 2026 16:08:53 +0800 Subject: [PATCH 03/12] test(subagents): cover native receipt replay and product readback Signed-off-by: jackie-cqz <2557911191@qq.com> --- examples/personal-workspace-browser-smoke.mjs | 2 + .../personal-workspace-browser/fixture.mjs | 4 + .../native-child-activity.mjs | 55 +++++++ .../test_codex_native_child_receipts.py | 139 ++++++++++++++++++ tests/control_plane_ts/agent_context.test.ts | 17 +++ tests/test_chat_agent.py | 18 +++ 6 files changed, 235 insertions(+) create mode 100644 examples/personal-workspace-browser/native-child-activity.mjs create mode 100644 tests/capabilities/test_codex_native_child_receipts.py 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/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py new file mode 100644 index 0000000000..82b3f8f686 --- /dev/null +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -0,0 +1,139 @@ +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from loopx.control_plane.turn_driver.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) + + +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["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" 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] From 4d5865543ac93fcd1d4b838bb33af0a8d9a1c14b Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 02:02:43 +0800 Subject: [PATCH 04/12] fix(subagents): preserve receipt correlation across host resumes Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../host-native-child-receipts.md | 20 ++++ .../multi_subagent/native_child_receipts.py | 36 ++++-- loopx/control_plane/turn_driver/codex_cli.py | 11 +- .../turn_driver/codex_native_child.py | 56 ++++++++-- .../turn_driver/codex_operation_host.py | 6 +- loopx/control_plane/turn_driver/executor.py | 5 + .../test_codex_native_child_receipts.py | 105 +++++++++++++++++- tests/test_loopx_turn_executor.py | 14 ++- 8 files changed, 228 insertions(+), 25 deletions(-) diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index 120218ed09..12da0d7938 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -127,6 +127,26 @@ 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 opaque child ID against the original host-observed spawn in the +same admitted Goal instance, coordinator and LoopX Turn, including receipts +outside the bounded status window. 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. No prompt or raw result is added to a receipt. + +恢复 CLI 时可能只收到新的 `wait` 完成事件。适配器从同一已准入 Goal 实例、 +主 Agent、LoopX Turn 的原始宿主启动回执恢复关联,覆盖状态窗口外的操作, +不采纳主 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 diff --git a/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index 707247c364..1d038eee8d 100644 --- a/loopx/capabilities/multi_subagent/native_child_receipts.py +++ b/loopx/capabilities/multi_subagent/native_child_receipts.py @@ -143,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, @@ -181,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 ( @@ -189,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( @@ -419,6 +431,7 @@ 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, ) -> dict[str, Any]: """Preview or append one report under its exact quota owner.""" @@ -448,4 +461,5 @@ def record_native_child( registry_path=registry_path, goal_ref=goal_ref, source_admission=source_admission, + _host_observed=_host_observed, ) diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index 7c586dc11d..a9d3ccc4f6 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -722,7 +722,14 @@ def commit() -> None: else: goal_admission.accept_result(commit) - child_observer = native_child_observer(request, runtime_root=runtime_root, lineage=lineage) + 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) @@ -764,7 +771,7 @@ def observe_event(line: str) -> None: 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]) + child_observer.observe(item, session_id=observed_session[0], invocation_id=invocation_id) if goal_admission is None: record_child() diff --git a/loopx/control_plane/turn_driver/codex_native_child.py b/loopx/control_plane/turn_driver/codex_native_child.py index 6b7c219f65..c799b34b60 100644 --- a/loopx/control_plane/turn_driver/codex_native_child.py +++ b/loopx/control_plane/turn_driver/codex_native_child.py @@ -6,11 +6,14 @@ 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 record_native_child +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: @@ -39,11 +42,14 @@ class CodexNativeChildObserver: """ def __init__(self, *, runtime_root: Path, lineage: Mapping[str, str], - turn_instance_id: str, configured_limit: int): + 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 self.children: dict[str, str] = {} def _record(self, *, stage: str, **record: str) -> None: @@ -52,10 +58,33 @@ def _record(self, *, stage: str, **record: str) -> None: 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, **record, + execute=True, _host_observed=True, goal_ref=self.goal_ref, + registry_path=self.registry_path, **record, ) - def observe(self, item: Mapping[str, Any], *, session_id: str) -> None: + def _restore_spawn(self, child: str) -> None: + # Spawn IDs already bind the opaque child identity. Restore only an + # admitted, host-observed spawn from this exact Goal/agent/Turn owner, + # including operations outside the bounded presentation window. + operation_id = "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, + ) + for event in events: + details = event.get("details") + if not isinstance(details, Mapping): + continue + if (event.get("event_kind") == "native_child_decision" + and event.get("case_id") == operation_id + and details.get("operation") == "spawn" and details.get("outcome") == "started" + and details.get("observation_source") == "host_observed" + and details.get("entrypoint_id") == "codex_native_tools"): + self.children[child] = operation_id + return + + 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 @@ -76,7 +105,10 @@ def observe(self, item: Mapping[str, Any], *, session_id: str) -> None: return # Successful spawn identity survives host event replay/restart; host # exec display-item counters alone are not globally unique. - identity = receivers[0] if tool == "spawn" and started else session_id + ":" + native_id + 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", @@ -88,7 +120,11 @@ def observe(self, item: Mapping[str, Any], *, session_id: str) -> None: if not isinstance(states, Mapping): return for child, state in states.items(): - if child not in self.children or not isinstance(state, Mapping): + if not isinstance(child, str) or not child or not isinstance(state, Mapping): + continue + if child not in self.children: + self._restore_spawn(child) + if child not in self.children: continue outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) if outcome: @@ -96,10 +132,14 @@ def observe(self, item: Mapping[str, Any], *, session_id: str) -> None: def native_child_observer(request: Mapping[str, Any], *, runtime_root: Path, - lineage: Mapping[str, str]) -> CodexNativeChildObserver | None: + 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) + 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/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index 34e2810a18..a3267baa2b 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -431,13 +431,15 @@ def on_event(kind: str, event: dict[str, Any]) -> None: recovery_kind="resume_session", ) from exc - child_observer = native_child_observer(request, runtime_root=runtime_root, lineage=lineage) + 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) + child_observer.observe(item, session_id=session.thread_id, + invocation_id=session.current_turn_id) if goal_admission is None: record_child() diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 59d504f484..89a404b139 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -801,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 .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/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index 82b3f8f686..7f5489584e 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -29,7 +29,7 @@ def _observe(root, items, session_id="parent-1"): 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) + observer.observe(item, session_id=session_id, invocation_id="host-turn-1") def test_host_spawn_result_parent_review_and_restart(tmp_path: Path): @@ -118,6 +118,7 @@ def test_cli_host_collects_native_items_before_returning_parent_result(tmp_path: _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" @@ -137,3 +138,105 @@ def host(command, **kwargs): 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): + 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"] = 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 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} From 326322d859173a5c6c3f56248229e5d24b36f7ec Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 13:33:24 +0800 Subject: [PATCH 05/12] refactor(codex): place native child receipts at the provider seam Signed-off-by: jackie-cqz <2557911191@qq.com> --- loopx/control_plane/turn_driver/codex_cli.py | 4 ++-- loopx/control_plane/turn_driver/codex_operation_host.py | 2 +- loopx/control_plane/turn_driver/executor.py | 4 ++-- .../turn_driver => extensions}/codex_native_child.py | 2 +- tests/capabilities/test_codex_native_child_receipts.py | 2 +- 5 files changed, 7 insertions(+), 7 deletions(-) rename loopx/{control_plane/turn_driver => extensions}/codex_native_child.py (99%) diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index a9d3ccc4f6..456222d5ec 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -15,8 +15,7 @@ from .subagent_execution_topology import ( child_execution_receipts_json_schema, ) -from .codex_native_child import native_child_observer -from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES, selected_turn_todo +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, @@ -29,6 +28,7 @@ codex_session_profile_digest, require_codex_session_profile, ) +from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES, selected_turn_todo from .executor import ( HOST_AGENT_VISION_JSON_MAX_CHARS, HOST_REWARD_MEMORY_REFLECTION_JSON_MAX_CHARS, diff --git a/loopx/control_plane/turn_driver/codex_operation_host.py b/loopx/control_plane/turn_driver/codex_operation_host.py index a3267baa2b..1df534c970 100644 --- a/loopx/control_plane/turn_driver/codex_operation_host.py +++ b/loopx/control_plane/turn_driver/codex_operation_host.py @@ -37,7 +37,7 @@ _store_codex_cli_session, load_codex_cli_session, ) -from .codex_native_child import native_child_observer +from ...extensions.codex_native_child import native_child_observer from .executor import LOOPX_TURN_HOST_REQUEST_SCHEMA_VERSION from .host_failure import BuiltInHostError diff --git a/loopx/control_plane/turn_driver/executor.py b/loopx/control_plane/turn_driver/executor.py index 89a404b139..4c2903f14a 100644 --- a/loopx/control_plane/turn_driver/executor.py +++ b/loopx/control_plane/turn_driver/executor.py @@ -161,7 +161,7 @@ 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 .codex_native_child import configured_native_child_limit + 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 @@ -801,7 +801,7 @@ 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 .codex_native_child import configured_native_child_limit + 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. diff --git a/loopx/control_plane/turn_driver/codex_native_child.py b/loopx/extensions/codex_native_child.py similarity index 99% rename from loopx/control_plane/turn_driver/codex_native_child.py rename to loopx/extensions/codex_native_child.py index c799b34b60..2cfab2f2a6 100644 --- a/loopx/control_plane/turn_driver/codex_native_child.py +++ b/loopx/extensions/codex_native_child.py @@ -11,7 +11,7 @@ from pathlib import Path from typing import Any -from ...capabilities.multi_subagent.native_child_receipts import ( +from ..capabilities.multi_subagent.native_child_receipts import ( _load_native_child_events, record_native_child, ) diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index 7f5489584e..d82e8ebe94 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -5,7 +5,7 @@ import pytest -from loopx.control_plane.turn_driver.codex_native_child import ( +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 bf6753d206d16cca02fd2714b5ffa160c8c531a0 Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 13:40:06 +0800 Subject: [PATCH 06/12] test(codex): reproduce resumed followup result ownership Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../test_codex_native_child_receipts.py | 36 +++++++++++++++++++ 1 file changed, 36 insertions(+) diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index d82e8ebe94..0f6f9aa5bf 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -240,3 +240,39 @@ def test_restart_replay_uses_the_original_invocation_binding(tmp_path): _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) + 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 From 29668f841d70d9313060aa6e7565c4eb18c73024 Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 13:46:08 +0800 Subject: [PATCH 07/12] fix(codex): recover the latest native followup result binding Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../multi_subagent/native_child_receipts.py | 13 +++++- loopx/extensions/codex_native_child.py | 41 ++++++++++--------- .../test_codex_native_child_receipts.py | 25 ++++++++++- 3 files changed, 57 insertions(+), 22 deletions(-) diff --git a/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index 1d038eee8d..c38ddc1221 100644 --- a/loopx/capabilities/multi_subagent/native_child_receipts.py +++ b/loopx/capabilities/multi_subagent/native_child_receipts.py @@ -295,6 +295,7 @@ def _record_native_child( 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, ) -> dict[str, Any]: """Preview or append a typed report; never launch a child or spend quota.""" goal_id = _id(goal_id, field="goal_id") @@ -303,7 +304,7 @@ 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, ) @@ -311,6 +312,14 @@ def _record_native_child( 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))}) log_path = rollout_event_log_path(runtime_root, goal_id) events = load_rollout_events(log_path) prior = _events_for_turn( @@ -432,6 +441,7 @@ def record_native_child( 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, ) -> dict[str, Any]: """Preview or append one report under its exact quota owner.""" @@ -462,4 +472,5 @@ def record_native_child( goal_ref=goal_ref, source_admission=source_admission, _host_observed=_host_observed, + _host_child_refs=_host_child_refs, ) diff --git a/loopx/extensions/codex_native_child.py b/loopx/extensions/codex_native_child.py index 2cfab2f2a6..64d3ad7595 100644 --- a/loopx/extensions/codex_native_child.py +++ b/loopx/extensions/codex_native_child.py @@ -50,9 +50,8 @@ def __init__(self, *, runtime_root: Path, lineage: Mapping[str, str], self.configured_limit = configured_limit self.goal_ref = goal_ref self.registry_path = registry_path - self.children: dict[str, str] = {} - def _record(self, *, stage: str, **record: str) -> None: + 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, @@ -62,27 +61,31 @@ def _record(self, *, stage: str, **record: str) -> None: registry_path=self.registry_path, **record, ) - def _restore_spawn(self, child: str) -> None: - # Spawn IDs already bind the opaque child identity. Restore only an - # admitted, host-observed spawn from this exact Goal/agent/Turn owner, - # including operations outside the bounded presentation window. - operation_id = "codex-" + hashlib.sha256(child.encode()).hexdigest()[:32] + def _restore_operation(self, child: 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, ) - for event in events: + for event in reversed(events): details = event.get("details") if not isinstance(details, Mapping): continue if (event.get("event_kind") == "native_child_decision" - and event.get("case_id") == operation_id - and details.get("operation") == "spawn" and details.get("outcome") == "started" + 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"): - self.children[child] = operation_id - return + 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") @@ -112,23 +115,21 @@ def observe(self, item: Mapping[str, Any], *, session_id: str, invocation_id: st 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 {})) - if started: - for child in receivers: - self.children[child] = operation_id + **({"reason_code": "host_failed"} if not started else { + "_host_child_refs": sorted({"codex-child-" + hashlib.sha256(child.encode()).hexdigest()[:32] + for child in receivers})})) 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 - if child not in self.children: - self._restore_spawn(child) - if child not in self.children: + operation_id = self._restore_operation(child) + if operation_id is None: continue outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) if outcome: - self._record(stage="result", operation_id=self.children[child], outcome=outcome) + self._record(stage="result", operation_id=operation_id, outcome=outcome) def native_child_observer(request: Mapping[str, Any], *, runtime_root: Path, diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index 0f6f9aa5bf..76fa5007b2 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -261,7 +261,7 @@ def test_real_cli_resumed_followup_owns_its_result(tmp_path, monkeypatch, with_s 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) + 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): @@ -276,3 +276,26 @@ def test_real_cli_consecutive_followups_restore_the_latest_owned_operation(tmp_p 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))) + + + +def test_real_cli_old_spawn_replay_does_not_replace_a_later_followup(tmp_path, monkeypatch): + _admit(tmp_path) + spawn, wait = _turn()['items'] + followup = {**spawn, 'id': 'item_0', 'tool': 'sendInput', 'agentsStates': {}} + activity = _real_cli_calls(tmp_path, monkeypatch, [[spawn, wait], [followup], [spawn, wait]]) + assert activity['operation_count'] == 2 and activity['launched_count'] == 1 + assert all(row['result'] == 'completed' for row in activity['operations']) From c15b9ba5e211c9a8958c0f01124eefcebeba81cc Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 13:46:08 +0800 Subject: [PATCH 08/12] docs(codex): describe durable followup result recovery Signed-off-by: jackie-cqz <2557911191@qq.com> --- docs/integrations/host-native-child-receipts.md | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index 12da0d7938..8b49fc004c 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -128,7 +128,8 @@ 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 opaque child ID against the original host-observed spawn in the +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. It does not adopt coordinator reports or scan external host history. Spawn IDs retain their existing child binding. @@ -137,15 +138,19 @@ 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. No prompt or raw result is added to a receipt. +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 的原始宿主启动回执恢复关联,覆盖状态窗口外的操作, +主 Agent、LoopX Turn 的最近一次已启动宿主 spawn 或 followup 回执恢复关联, +覆盖状态窗口外的操作, 不采纳主 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 From cdfa279b28e69b9696bf5ba2cf83c983b582d434 Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 15:35:27 +0800 Subject: [PATCH 09/12] fix(codex): bind terminal snapshots to their own decision Signed-off-by: jackie-cqz <2557911191@qq.com> --- loopx/extensions/codex_native_child.py | 10 +++- .../test_codex_native_child_receipts.py | 49 ++++++++++++++++--- 2 files changed, 52 insertions(+), 7 deletions(-) diff --git a/loopx/extensions/codex_native_child.py b/loopx/extensions/codex_native_child.py index 64d3ad7595..b05296b51d 100644 --- a/loopx/extensions/codex_native_child.py +++ b/loopx/extensions/codex_native_child.py @@ -102,6 +102,7 @@ def observe(self, item: Mapping[str, Any], *, session_id: str, invocation_id: st 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): @@ -118,13 +119,20 @@ def observe(self, item: Mapping[str, Any], *, session_id: str, invocation_id: st **({"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 - operation_id = self._restore_operation(child) + # A spawn/followup snapshot belongs to that stable decision, even + # when replayed after later work. Only a new wait needs to recover + # the latest operation from the durable child association. + if tool != "wait" and (not isinstance(receivers, list) or child not in receivers): + continue + operation_id = self._restore_operation(child) if tool == "wait" else decision_id if operation_id is None: continue outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index 76fa5007b2..4407e26fcc 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -292,10 +292,47 @@ def test_child_correlation_cannot_attest_a_report_or_result(tmp_path, host_obser -def test_real_cli_old_spawn_replay_does_not_replace_a_later_followup(tmp_path, monkeypatch): +@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'] - followup = {**spawn, 'id': 'item_0', 'tool': 'sendInput', 'agentsStates': {}} - activity = _real_cli_calls(tmp_path, monkeypatch, [[spawn, wait], [followup], [spawn, wait]]) - assert activity['operation_count'] == 2 and activity['launched_count'] == 1 - assert all(row['result'] == 'completed' for row in activity['operations']) + 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"]) From ddac46075cb8b9790bfc78312e7afca021b6413a Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 15:35:27 +0800 Subject: [PATCH 10/12] docs(native-child): clarify terminal replay ownership Signed-off-by: jackie-cqz <2557911191@qq.com> --- docs/integrations/host-native-child-receipts.md | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index 8b49fc004c..f3a30958ef 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -131,8 +131,11 @@ 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. It does not adopt coordinator reports or -scan external host history. Spawn IDs retain their existing child binding. +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. Only a new wait uses the latest child association. +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; @@ -144,7 +147,8 @@ No prompt or raw result is added to a receipt. 恢复 CLI 时可能只收到新的 `wait` 完成事件。适配器从同一已准入 Goal 实例、 主 Agent、LoopX Turn 的最近一次已启动宿主 spawn 或 followup 回执恢复关联, -覆盖状态窗口外的操作, +覆盖状态窗口外的操作。spawn/followup 的完成快照只归属自身稳定决策和接收者; +重放旧启动事件不会完成后来的跟进任务,只有新的 wait 使用最近关联。 不采纳主 Agent 上报或扫描外部宿主历史。启动保留既有子代理绑定;CLI 跟进和 失败调用复用 Turn 日志持久化的 `host_attempt`、父会话和原生工具 ID, app-server 使用其原生 Turn ID。真实重试在启动前递增尝试次数;同一绑定与 From f241c562e7f0b56f8b9fb4c8506dd93ee5e74a5f Mon Sep 17 00:00:00 2001 From: jackie-cqz <2557911191@qq.com> Date: Sat, 3 Oct 2026 19:11:29 +0800 Subject: [PATCH 11/12] fix(codex): retain consumed wait result ownership across replay Signed-off-by: jackie-cqz <2557911191@qq.com> --- .../host-native-child-receipts.md | 15 ++- .../multi_subagent/native_child_receipts.py | 36 ++++++- loopx/extensions/codex_native_child.py | 32 ++++-- .../test_codex_native_child_receipts.py | 99 ++++++++++++++++++- 4 files changed, 168 insertions(+), 14 deletions(-) diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index f3a30958ef..597555b7a5 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -133,7 +133,14 @@ 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. Only a new wait uses the latest child association. +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` @@ -148,7 +155,11 @@ No prompt or raw result is added to a receipt. 恢复 CLI 时可能只收到新的 `wait` 完成事件。适配器从同一已准入 Goal 实例、 主 Agent、LoopX Turn 的最近一次已启动宿主 spawn 或 followup 回执恢复关联, 覆盖状态窗口外的操作。spawn/followup 的完成快照只归属自身稳定决策和接收者; -重放旧启动事件不会完成后来的跟进任务,只有新的 wait 使用最近关联。 +重放旧启动事件不会完成后来的跟进任务。wait 首先恢复其会话、调用、原生项和 +接收子 Agent 身份对应的首次结果关联;只有首次观察才使用最近操作。哈希关联 +作为标量留在现有结果事件中,重启重放仍归原操作,结果冲突或重新归属会被拒绝。 +多个 wait 观察同一结果不会增加操作、启动、主 Agent 验收或配额;公开活动投影 +不暴露原始宿主标识或内容。 不采纳主 Agent 上报或扫描外部宿主历史。启动保留既有子代理绑定;CLI 跟进和 失败调用复用 Turn 日志持久化的 `host_attempt`、父会话和原生工具 ID, app-server 使用其原生 Turn ID。真实重试在启动前递增尝试次数;同一绑定与 diff --git a/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index c38ddc1221..e823a2859d 100644 --- a/loopx/capabilities/multi_subagent/native_child_receipts.py +++ b/loopx/capabilities/multi_subagent/native_child_receipts.py @@ -296,6 +296,7 @@ def _record_native_child( 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") @@ -320,6 +321,10 @@ def _record_native_child( # 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( @@ -331,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, @@ -361,6 +391,7 @@ 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": @@ -407,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)), ) @@ -442,6 +474,7 @@ def record_native_child( 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.""" @@ -473,4 +506,5 @@ def record_native_child( source_admission=source_admission, _host_observed=_host_observed, _host_child_refs=_host_child_refs, + _host_wait_ref=_host_wait_ref, ) diff --git a/loopx/extensions/codex_native_child.py b/loopx/extensions/codex_native_child.py index b05296b51d..1dff414114 100644 --- a/loopx/extensions/codex_native_child.py +++ b/loopx/extensions/codex_native_child.py @@ -61,7 +61,7 @@ def _record(self, *, stage: str, **record: Any) -> None: registry_path=self.registry_path, **record, ) - def _restore_operation(self, child: str) -> str | None: + 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. @@ -72,6 +72,15 @@ def _restore_operation(self, child: str) -> str | None: 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): @@ -128,16 +137,23 @@ def observe(self, item: Mapping[str, Any], *, session_id: str, invocation_id: st 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. Only a new wait needs to recover - # the latest operation from the durable child association. + # 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 - operation_id = self._restore_operation(child) if tool == "wait" else decision_id - if operation_id is None: - continue outcome = {"completed": "completed", "errored": "failed", "shutdown": "cancelled"}.get(state.get("status")) - if outcome: - self._record(stage="result", operation_id=operation_id, outcome=outcome) + 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, diff --git a/tests/capabilities/test_codex_native_child_receipts.py b/tests/capabilities/test_codex_native_child_receipts.py index 4407e26fcc..178f866434 100644 --- a/tests/capabilities/test_codex_native_child_receipts.py +++ b/tests/capabilities/test_codex_native_child_receipts.py @@ -140,7 +140,7 @@ def host(command, **kwargs): assert activity["operations"][0]["result"] == "completed" -def _real_cli_calls(tmp_path, monkeypatch, batches): +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 @@ -168,7 +168,7 @@ def command(**kwargs): 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"] = attempt + 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" @@ -250,7 +250,8 @@ def test_real_cli_resumed_followup_owns_its_result(tmp_path, monkeypatch, with_s 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) + 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) @@ -336,3 +337,95 @@ def test_nonwait_terminal_snapshot_ignores_unrelated_receivers(tmp_path, tool): 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") From c2ea347d6156c860a2e0dd07e8682abfd9fa2d78 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 00:08:26 +0800 Subject: [PATCH 12/12] chore(codex): drop the lineage import the session refactor moved The rebase onto main moved _lineage into codex_sessions.py, which keeps the selected_turn_todo lookup, so codex_cli.py no longer imports it. Behavior is unchanged; this only satisfies the required Ruff check. Signed-off-by: huangruiteng --- loopx/control_plane/turn_driver/codex_cli.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/loopx/control_plane/turn_driver/codex_cli.py b/loopx/control_plane/turn_driver/codex_cli.py index 456222d5ec..93224f246c 100644 --- a/loopx/control_plane/turn_driver/codex_cli.py +++ b/loopx/control_plane/turn_driver/codex_cli.py @@ -28,7 +28,7 @@ codex_session_profile_digest, require_codex_session_profile, ) -from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES, selected_turn_todo +from .driver import SUPPORTED_ITERATION_CONTEXT_POLICIES from .executor import ( HOST_AGENT_VISION_JSON_MAX_CHARS, HOST_REWARD_MEMORY_REFLECTION_JSON_MAX_CHARS,