From 2cbe77ec50ebbab0b0700a6df5d6b595e7abf23b Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Mon, 28 Sep 2026 22:50:24 +0800 Subject: [PATCH] fix(subagents): admit replan reports through TS settlement readback Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../host-native-child-receipts.md | 26 +++- loopx/capabilities/multi_subagent/cli.py | 2 +- .../multi_subagent/native_child_receipts.py | 53 +++++-- .../capabilities/native_child_admission.ts | 47 ++++++ loopx/control_plane/quota/settlement.py | 4 + .../quota/settlement_readback.ts | 63 ++++++-- .../test_native_child_receipts.py | 109 ++++++++++++- .../test_native_child_closeout_cli.py | 80 ++++++++++ .../test_native_child_replan_guard_cli.py | 120 ++++++++++++++ .../native_child_admission.test.ts | 146 ++++++++++++++++++ 10 files changed, 618 insertions(+), 32 deletions(-) create mode 100644 loopx/control_plane/capabilities/native_child_admission.ts create mode 100644 tests/control_plane/test_native_child_closeout_cli.py create mode 100644 tests/control_plane/test_native_child_replan_guard_cli.py create mode 100644 tests/control_plane_ts/native_child_admission.test.ts diff --git a/docs/integrations/host-native-child-receipts.md b/docs/integrations/host-native-child-receipts.md index 7ccc04b921..b30eeaf878 100644 --- a/docs/integrations/host-native-child-receipts.md +++ b/docs/integrations/host-native-child-receipts.md @@ -33,8 +33,13 @@ was created or deliberately skipped. ## Lifecycle / 生命周期 -1. The admitted Turn guard must already have a settlement binding. `record` - rejects unregistered coordinators and disabled policy before writing. +1. The admitted Turn guard must already have a settlement binding. The existing + TS settlement readback verifies its exact original identity and reporting + admission. Committed work-admission facts also admit lawful replan Turns; + the recorder does not whitelist their status labels. Old ordinary guards + without work-projection fields retain bounded compatibility, never a fallback + for partial or negative facts. `record` rejects unregistered coordinators and + disabled policy before writing. 2. Use `stage=decision` with a stable `operation-id` for `spawn`, `followup`, or a bounded `skip` reason. A host capacity rejection maps to the generic `host_capacity_exhausted` reason. Capacity rejection and typed host failure @@ -50,8 +55,18 @@ was created or deliberately skipped. the latest reported Turn only when the capability is enabled and a decision exists. The dashboard uses that status projection, and omits the activity line for unconfigured or unrelated Goals. +5. A new decision requires an open, work-admitted Turn and no begun closeout. + Once closeout is pending or settled, an already-recorded started operation + may still receive its result and parent review. This records late facts; it + does not reopen the Turn. An absent decision cannot be backfilled through a + closed Turn. Exact duplicates remain reads; conflicts remain errors. For a + new append the readback is checked inside the existing event-log lock, with + one coarse TS call per report, not a second Python phase rule. -1. Turn 须先有已提交的结算绑定;未注册主 Agent 或未启用策略不能写入。 +1. Turn 须先有已提交的结算绑定;既有 TS 结算读回核验精确原身份及报告准入。 + 已提交的工作准入事实同样覆盖合法重规划,报告入口不按状态名称建立白名单。 + 没有工作投影字段的旧普通 guard 保留有界兼容,不能用来绕过不完整或否定的 + 准入事实。未注册主 Agent 或未启用策略不能写入。 2. 用稳定 `operation-id` 写 `decision`,区分 `spawn`、`followup` 和有界理由的 `skip`。宿主容量拒绝映射为通用 `host_capacity_exhausted`;容量拒绝和 类型化宿主失败都停止同一 Turn 的启动或跟进重试,主 Agent 仍可继续工作。 @@ -61,6 +76,11 @@ was created or deliberately skipped. 4. 同一身份和内容重放幂等,内容冲突会被拒绝。`read` 与带 Turn ID 的 `agent-context` 读取同一模型。Goal 状态的 JSON 与 Markdown 只在能力启用且确有决策时投影 最近一轮;仪表板读取该投影,未配置或无关 Goal 不显示活动行。 +5. 新决策要求 Turn 已准入、仍开放且尚未开始结算。开始结算或结清后,已经登记 + 的 started 操作仍可接收结果及主 Agent 验收;这是迟到事实登记,不会重开 Turn。 + 不能借已关闭的 Turn 首次补建缺失决策。精确重复仍是读取,内容冲突仍拒绝。 + 新写入在既有事件流锁内重新核对读回,每次报告只调用一个粗粒度 TS 边界, + 不增加第二套 Python 阶段判断。 Example / 示例: diff --git a/loopx/capabilities/multi_subagent/cli.py b/loopx/capabilities/multi_subagent/cli.py index d411ac8f6d..a9fcf0f2eb 100644 --- a/loopx/capabilities/multi_subagent/cli.py +++ b/loopx/capabilities/multi_subagent/cli.py @@ -63,7 +63,7 @@ def handle_native_child_command(args, registry_path, runtime_root, print_payload evidence_ref=args.evidence_ref, validation_ref=args.validation_ref, execute=args.execute, ) - except (OSError, ValueError, KeyError) as exc: + except (OSError, ValueError, KeyError, RuntimeError) as exc: payload = {"ok": False, "error": str(exc)} print_payload(payload, output_format(args), render_native_child) return 0 if payload["ok"] else 1 diff --git a/loopx/capabilities/multi_subagent/native_child_receipts.py b/loopx/capabilities/multi_subagent/native_child_receipts.py index 698e4982d6..ead59dcd0b 100644 --- a/loopx/capabilities/multi_subagent/native_child_receipts.py +++ b/loopx/capabilities/multi_subagent/native_child_receipts.py @@ -13,7 +13,7 @@ from pathlib import Path from typing import Any -from ...control_plane.quota.heartbeat_receipt import find_heartbeat_receipt +from ...control_plane.quota.settlement import read_heartbeat_settlement from ...control_plane.runtime.public_safety import validate_public_safe_value from ...rollout_event_log import ( append_rollout_event_once, @@ -63,10 +63,11 @@ def _details(event: Mapping[str, Any]) -> dict[str, Any]: def _events_for_turn( - events: Sequence[Mapping[str, Any]], *, agent_id: str, turn_instance_id: str, + events: Sequence[Mapping[str, Any]], *, goal_id: str, agent_id: str, turn_instance_id: str, ) -> list[dict[str, Any]]: return [dict(event) for event in events 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] @@ -76,7 +77,7 @@ def native_child_activity( turn_instance_id: str, configured_limit: int, ) -> dict[str, Any]: """One read model for CLI, agent-context and product projections.""" - rows = _events_for_turn(events, agent_id=agent_id, turn_instance_id=turn_instance_id) + rows = _events_for_turn(events, goal_id=goal_id, agent_id=agent_id, turn_instance_id=turn_instance_id) operations: dict[str, dict[str, Any]] = {} for event in rows: details = _details(event) @@ -230,35 +231,50 @@ 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, ) - guard = find_heartbeat_receipt( - runtime_root, goal_id=goal_id, agent_id=agent_id, - turn_instance_id=turn_instance_id, - ) - if not guard or not _details(guard).get("settlement_effect_id"): - raise ValueError("native child report requires an admitted, settlement-bound Turn guard") log_path = rollout_event_log_path(runtime_root, goal_id) events = load_rollout_events(log_path) - prior = _events_for_turn(events, agent_id=agent_id, turn_instance_id=turn_instance_id) + prior = _events_for_turn(events, goal_id=goal_id, agent_id=agent_id, turn_instance_id=turn_instance_id) existing = next((event for event in prior if event.get("case_id") == operation_id and event.get("event_kind") == EVENT_KINDS[stage]), None) if existing is not None and _details(existing) != fields: raise ValueError("operation identity already has a conflicting native child report") + def report_admission() -> Mapping[str, Any]: + readback = read_heartbeat_settlement( + runtime_root, goal_id=goal_id, agent_id=agent_id, todo_id=None, + turn_instance_id=turn_instance_id, resolve_original_binding=True, + ) + if readback is None or readback.identity.value is None or readback.identity.failure is not None: + reason = (readback.identity.failure.reason + if readback is not None and readback.identity.failure is not None else "guard missing") + raise ValueError("native child report requires an admitted, settlement-bound Turn guard: " + reason) + admission = readback.native_child_admission + if (not isinstance(admission, Mapping) + or admission.get("schema_version") != "native_child_report_admission_v0" + or admission.get("report_permission") not in { + "new_operation", "existing_operation_only", "not_admitted", + }): + raise ValueError("TypeScript native child report admission missing or invalid") + return admission + def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: - current = _events_for_turn(observed, agent_id=agent_id, + admission = report_admission() + current = _events_for_turn(observed, goal_id=goal_id, agent_id=agent_id, turn_instance_id=turn_instance_id) decisions = {str(item.get("case_id")): item for item in current if item.get("event_kind") == EVENT_KINDS["decision"]} if stage == "decision": - if str(guard.get("status") or "") not in {"normal_run", "turn_run_once"}: - raise ValueError("native child decision requires a runnable Turn guard") + 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( _details(item).get("outcome") in {"capacity_rejected", "host_failed"} for item in decisions.values() ): raise ValueError("host failure forbids same-Turn spawn/followup retry") return + if admission["report_permission"] == "not_admitted": + raise ValueError("native child report requires a work-admitted Turn guard") decision = decisions.get(operation_id) if decision is None or _details(decision).get("outcome") != "started": raise ValueError("result/review requires a started native child decision") @@ -268,8 +284,11 @@ def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: if result is None or _details(result).get("outcome") != "completed": raise ValueError("parent review requires a completed native child result") - if existing is None: - validate_transition(prior) + if not execute: + if existing is None: + validate_transition(prior) + else: + report_admission() event = build_rollout_event( goal_id=goal_id, event_kind=EVENT_KINDS[stage], agent_id=agent_id, run_id=turn_instance_id, case_id=operation_id, status=fields["outcome"], @@ -284,6 +303,10 @@ def validate_transition(observed: Sequence[Mapping[str, Any]]) -> None: ) if _details(stored) != fields: raise ValueError("operation identity already has a conflicting native child report") + if not appended: + # Append-once skips its precondition for an exact duplicate. Verify + # the original identity without re-authorizing work or appending. + report_admission() events = load_rollout_events(log_path) else: stored = existing or event diff --git a/loopx/control_plane/capabilities/native_child_admission.ts b/loopx/control_plane/capabilities/native_child_admission.ts new file mode 100644 index 0000000000..03c1d0a1cc --- /dev/null +++ b/loopx/control_plane/capabilities/native_child_admission.ts @@ -0,0 +1,47 @@ +/** Report admission is not permission to launch a child or spend quota. */ +import type {JsonObject} from "../effect_program.ts"; +import type {ReceiptBoundReplayPhase} from "../quota/settlement_phase.ts"; + +export const NATIVE_CHILD_REPORT_ADMISSION_SCHEMA = "native_child_report_admission_v0"; + +export interface NativeChildReportAdmission extends JsonObject { + schema_version: typeof NATIVE_CHILD_REPORT_ADMISSION_SCHEMA; + report_permission: "new_operation" | "existing_operation_only" | "not_admitted"; + reason_code: "turn_work_admitted" | "turn_closeout_started" | "guard_work_not_admitted"; + settlement_effect_id: string; + turn_phase: ReceiptBoundReplayPhase; +} + +/** Consume committed admission facts, not a whitelist of current status labels. + * Identity and phase must first be verified by the settlement readback owner. */ +export function nativeChildReportAdmission( + guardDetails: JsonObject, + guardStatus: unknown, + effectId: string, + phase: ReceiptBoundReplayPhase, + closeoutStarted: boolean, +): NativeChildReportAdmission { + const projectionAbsent = ["must_attempt_work", "delivery_allowed", "quiet_noop_allowed"] + .every((field) => !Object.hasOwn(guardDetails, field)); + // Old ordinary/managed guards predate the work projection. Never use this + // compatibility branch to turn partial or negative modern facts into work. + const legacyAdmission = projectionAbsent && + (guardStatus === "normal_run" || guardStatus === "turn_run_once") && + (guardDetails.ok === undefined || guardDetails.ok === true) && + (guardDetails.should_run === undefined || guardDetails.should_run === true); + const admitted = legacyAdmission || ( + guardDetails.ok === true && guardDetails.should_run === true && + guardDetails.must_attempt_work === true && guardDetails.delivery_allowed === true && + guardDetails.quiet_noop_allowed === false + ); + const open = phase === "open" && !closeoutStarted; + return { + schema_version: NATIVE_CHILD_REPORT_ADMISSION_SCHEMA, + report_permission: !admitted ? "not_admitted" + : open ? "new_operation" : "existing_operation_only", + reason_code: !admitted ? "guard_work_not_admitted" + : open ? "turn_work_admitted" : "turn_closeout_started", + settlement_effect_id: effectId, + turn_phase: phase, + }; +} diff --git a/loopx/control_plane/quota/settlement.py b/loopx/control_plane/quota/settlement.py index 03e045f8c1..f5317506cc 100644 --- a/loopx/control_plane/quota/settlement.py +++ b/loopx/control_plane/quota/settlement.py @@ -163,6 +163,7 @@ class QuotaSettlementReadback: refresh_recovery: dict[str, Any] | None = None external_delivery: dict[str, Any] | None = None progress: dict[str, Any] | None = None + native_child_admission: dict[str, Any] | None = None def attach_settlement_progress( @@ -308,6 +309,7 @@ def read_heartbeat_settlement( replan_obligation_id: str | None = None, infer_turn_instance_id: bool = False, allow_unbound_binding: bool = False, + resolve_original_binding: bool = False, refresh_retry: dict[str, Any] | None = None, ) -> QuotaSettlementReadback | None: """Read one complete heartbeat settlement through the TS domain owner.""" @@ -325,6 +327,7 @@ def read_heartbeat_settlement( "replan_obligation_id": replan_obligation_id, "infer_turn_instance_id": infer_turn_instance_id, "allow_unbound_binding": allow_unbound_binding, + **({"resolve_original_binding": True} if resolve_original_binding else {}), **( {"refresh_retry": refresh_retry} if refresh_retry is not None @@ -376,6 +379,7 @@ def read_heartbeat_settlement( refresh_recovery=_optional_readback_record(payload.get("refresh_recovery")), external_delivery=_optional_readback_record(payload.get("external_delivery")), progress=_optional_readback_record(payload.get("progress")), + native_child_admission=_optional_readback_record(payload.get("native_child_admission")), spend_run=_optional_readback_record(payload.get("spend_run")), heartbeat_receipt=_optional_readback_record(payload.get("heartbeat_receipt")), writeback_event=_optional_readback_record(payload.get("writeback_event")), diff --git a/loopx/control_plane/quota/settlement_readback.ts b/loopx/control_plane/quota/settlement_readback.ts index af0d222e23..8c2dbaab22 100644 --- a/loopx/control_plane/quota/settlement_readback.ts +++ b/loopx/control_plane/quota/settlement_readback.ts @@ -55,6 +55,7 @@ import { import { refreshExternalDelivery } from "./refresh_external_delivery.ts"; import { BLOCKED_WAIT_REQUEST_SCHEMA, prepareBlockedWait } from "./blocked_wait.ts"; +import {nativeChildReportAdmission} from "../capabilities/native_child_admission.ts"; export const QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA = "loopx_quota_settlement_readback_request_v0"; @@ -74,6 +75,7 @@ interface ReadbackRequest { replan_obligation_id: string | null; infer_turn_instance_id: boolean; allow_unbound_binding: boolean; + resolve_original_binding: boolean; refresh_retry: RefreshRetryRequest | null; } @@ -168,6 +170,13 @@ function decodeRequest(value: unknown): ReadbackRequest { if (typeof request.allow_unbound_binding !== "boolean") { throw new EffectRuntimeRequestError("allow_unbound_binding must be a boolean"); } + if (request.resolve_original_binding !== undefined && + typeof request.resolve_original_binding !== "boolean") { + throw new EffectRuntimeRequestError("resolve_original_binding must be a boolean"); + } + if (request.resolve_original_binding === true && request.infer_turn_instance_id) { + throw new EffectRuntimeRequestError("original binding requires an explicit Turn identity"); + } return { runtime_root: runtimeRoot, goal_id: goalId, @@ -183,6 +192,7 @@ function decodeRequest(value: unknown): ReadbackRequest { ), infer_turn_instance_id: request.infer_turn_instance_id, allow_unbound_binding: request.allow_unbound_binding, + resolve_original_binding: request.resolve_original_binding === true, refresh_retry: decodeRefreshRetry(request.refresh_retry), }; } @@ -747,6 +757,32 @@ function resolveIdentity( "invalid_identity", ); } + if (request.resolve_original_binding && agentId && turnInstanceId) { + // Resolve only this exact Turn's committed guard. This is a read, never + // the guard's binder, a current-Todo selection or a latest-run inference. + let guard: JsonObject | null; + try { + guard = effectiveHeartbeatReceipt(events, { + goal_id: request.goal_id, agent_id: agentId, turn_instance_id: turnInstanceId, + }); + } catch (error) { + return failedIdentity((error as Error).message, "identity_mismatch"); + } + if (!guard) return failedIdentity("matching Turn guard is missing", "receipt_missing"); + const fact = heartbeatReceiptFactFromEvent(guard); + const originalTodoId = normalizeTodoId(fact.todo_id); + const originalReplanId = normalizeReplanObligationId(fact.replan_obligation_id); + if (!originalTodoId && !originalReplanId) { + return failedIdentity("the original Turn guard has no settlement binding", "receipt_unbound"); + } + if ((request.todo_id !== null && (todoId === null || todoId !== originalTodoId)) || + (request.replan_obligation_id !== null && + (replanObligationId === null || replanObligationId !== originalReplanId))) { + return failedIdentity("requested binding differs from the original Turn guard", "identity_mismatch"); + } + todoId = originalTodoId; + replanObligationId = originalReplanId; + } if (!agentId || !turnInstanceId || Boolean(todoId) === Boolean(replanObligationId)) { return failedIdentity( "turn-scoped settlement requires agent_id, turn_instance_id, and exactly one todo_id or replan_obligation_id", @@ -894,6 +930,7 @@ function failedReadback( monitor_phase: null, replay_phase: null, refresh_recovery: null, + native_child_admission: null, }; } @@ -1011,6 +1048,16 @@ function readQuotaSettlementFromRequest( semanticReplanGuard.selected_obligation_id !== null; const inFlightWriteback = writeback.failure === null && isAcceptedInFlightWriteback(writebackRun, identity); + const replayPhase = receiptBoundReplayPhase({ + binding_kind: identity.binding_kind, + writeback_completes_binding: todoBoundReplan || blockedNoSpend || inFlightWriteback, + completion_receipt_present: completionEvent !== null, + supersede_receipt_present: supersedeEvent?.status === "done" && + optionalString(supersedeEvent.todo_id) === identity.todo_id, + durable_writeback_present: writeback.failure === null, + quota_spend_present: spend.failure === null, + no_spend_closeout_present: blockedNoSpend, + }); const recovery = request.refresh_retry === null ? null : refreshRecovery( request.refresh_retry, writebackRun, writeback.failure === null, @@ -1054,16 +1101,12 @@ function readQuotaSettlementFromRequest( durable_writeback_present: writeback.failure === null, quota_spend_present: spend.failure === null, }), - replay_phase: receiptBoundReplayPhase({ - binding_kind: identity.binding_kind, - writeback_completes_binding: todoBoundReplan || blockedNoSpend || inFlightWriteback, - completion_receipt_present: completionEvent !== null, - supersede_receipt_present: supersedeEvent?.status === "done" && - optionalString(supersedeEvent.todo_id) === identity.todo_id, - durable_writeback_present: writeback.failure === null, - quota_spend_present: spend.failure === null, - no_spend_closeout_present: blockedNoSpend, - }), + replay_phase: replayPhase, + native_child_admission: nativeChildReportAdmission( + receiptDetails, heartbeatReceipt.status, identity.effect_id, replayPhase, + [writebackRun, writebackEvent, spendRun, spendEvent, completionEvent, supersedeEvent] + .some((receipt) => receipt !== null), + ), }; } diff --git a/tests/capabilities/test_native_child_receipts.py b/tests/capabilities/test_native_child_receipts.py index 59ce71799b..f5131f005b 100644 --- a/tests/capabilities/test_native_child_receipts.py +++ b/tests/capabilities/test_native_child_receipts.py @@ -10,6 +10,7 @@ from loopx.capabilities.multi_subagent.native_child_receipts import ( latest_native_child_activity, load_native_child_activity, + native_child_activity, record_native_child, ) from loopx.rollout_event_log import ( @@ -30,9 +31,9 @@ def _admit(runtime_root: Path) -> None: rollout_event_log_path(runtime_root, GOAL), build_rollout_event( goal_id=GOAL, event_kind="quota_should_run", agent_id=AGENT, - run_id=TURN, todo_id="todo-1", status="normal_run", - details={"todo_id": "todo-1", - "settlement_effect_id": f"{GOAL}:{AGENT}:todo-1:{TURN}"}, + run_id=TURN, todo_id="todo_native_1", status="normal_run", + details={"todo_id": "todo_native_1", + "settlement_effect_id": f"{GOAL}:{AGENT}:todo_native_1:{TURN}"}, ), ) @@ -258,3 +259,105 @@ def test_goal_status_only_attaches_reported_activity_to_enabled_goal( contract={}, history={}, global_registry={}, runtime_root=tmp_path, ) assert "native_child_activity" not in disabled["items"][0]["project_asset"] + + +def test_report_phase_is_checked_once_under_the_append_lock(tmp_path: Path, monkeypatch): + import loopx.capabilities.multi_subagent.native_child_receipts as native + + _admit(tmp_path) + readback, append_once = native.read_heartbeat_settlement, native.append_rollout_event_once + calls = [] + + def read_once(*args, **kwargs): + calls.append(kwargs["turn_instance_id"]) + return readback(*args, **kwargs) + + monkeypatch.setattr(native, "read_heartbeat_settlement", read_once) + first = _record(tmp_path, "op-1", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + assert calls == [TURN] + replay = _record(tmp_path, "op-1", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + assert replay["receipt"]["event_id"] == first["receipt"]["event_id"] + assert calls == [TURN, TURN] + + def close_before_append(*args, **kwargs): + append_rollout_event(rollout_event_log_path(tmp_path, GOAL), build_rollout_event( + goal_id=GOAL, event_kind="todo_complete", agent_id=AGENT, run_id=TURN, + todo_id="todo_native_1", status="done", + details={"settlement_effect_id": f"{GOAL}:{AGENT}:todo_native_1:{TURN}"}, + )) + return append_once(*args, **kwargs) + + monkeypatch.setattr(native, "append_rollout_event_once", close_before_append) + with pytest.raises(ValueError, match="open, work-admitted"): + _record(tmp_path, "op-new", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + assert calls == [TURN, TURN, TURN] + assert load_native_child_activity(tmp_path, goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=6)["operation_count"] == 1 + + +def test_negative_guard_facts_and_binding_conflicts_do_not_reauthorize_reports(tmp_path: Path): + _admit(tmp_path) + _record(tmp_path, "op-1", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + log = rollout_event_log_path(tmp_path, GOAL) + append_rollout_event(log, build_rollout_event( + goal_id=GOAL, event_kind="quota_should_run", agent_id=AGENT, run_id=TURN, + status="normal_run", details={"todo_id": "todo_native_1", + "settlement_effect_id": f"{GOAL}:{AGENT}:todo_native_1:{TURN}", + "must_attempt_work": True, "delivery_allowed": False}, + )) + with pytest.raises(ValueError, match="open, work-admitted"): + _record(tmp_path, "op-new", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + # A duplicate is a fact read, not permission to start the operation again. + assert _record(tmp_path, "op-1", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host")["appended"] is False + append_rollout_event(log, build_rollout_event( + goal_id=GOAL, event_kind="quota_should_run", agent_id=AGENT, run_id=TURN, + status="normal_run", details={"todo_id": "todo_other_binding", + "settlement_effect_id": f"{GOAL}:{AGENT}:todo_other_binding:{TURN}"}, + )) + with pytest.raises(ValueError, match="settlement-bound"): + _record(tmp_path, "op-1", stage="decision", operation="spawn", + outcome="started", entrypoint_id="generic_host") + assert sum(row["event_kind"] == "native_child_decision" for row in load_rollout_events(log)) == 1 + + +def test_foreign_goal_facts_cannot_supply_a_native_stage_prerequisite(tmp_path: Path): + _admit(tmp_path) + foreign = build_rollout_event(goal_id="other-goal", event_kind="native_child_decision", + agent_id=AGENT, run_id=TURN, case_id="foreign-op", status="started", + details={"operation": "spawn", "outcome": "started", "entrypoint_id": "generic_host"}) + append_rollout_event(rollout_event_log_path(tmp_path, GOAL), foreign) + with pytest.raises(ValueError, match="started"): + _record(tmp_path, "foreign-op", stage="result", outcome="completed") + assert native_child_activity([foreign], goal_id=GOAL, agent_id=AGENT, + turn_instance_id=TURN, configured_limit=6)["operation_count"] == 0 + + +def test_cli_runtime_failure_is_visible_and_never_writes_a_report(tmp_path: Path, monkeypatch, capsys): + import loopx.capabilities.multi_subagent.native_child_receipts as native + from loopx.cli import main + + runtime = tmp_path / "runtime" + _admit(runtime) + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"schema_version": 1, "goals": [{ + "id": GOAL, "repo": str(tmp_path), "status": "active", "registered_agents": [AGENT], + "spawn_policy": {"mode": "multi_subagent", "allowed": True, "max_children": 6}, + }]})) + + def unavailable(*_args, **_kwargs): + raise RuntimeError("TypeScript runtime unavailable: synthetic diagnostic") + + monkeypatch.setattr(native, "read_heartbeat_settlement", unavailable) + code = main(["--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "native-child", "record", "--goal-id", GOAL, "--agent-id", AGENT, "--turn-instance-id", TURN, + "--operation-id", "op-1", "--stage", "decision", "--operation", "spawn", + "--outcome", "started", "--entrypoint-id", "generic_host", "--execute"]) + assert code == 1 + assert "TypeScript runtime unavailable" in json.loads(capsys.readouterr().out)["error"] + assert len(load_rollout_events(rollout_event_log_path(runtime, GOAL))) == 1 diff --git a/tests/control_plane/test_native_child_closeout_cli.py b/tests/control_plane/test_native_child_closeout_cli.py new file mode 100644 index 0000000000..bbe2025db0 --- /dev/null +++ b/tests/control_plane/test_native_child_closeout_cli.py @@ -0,0 +1,80 @@ +"""Asynchronous facts may arrive after closeout, but cannot authorize new work.""" +from __future__ import annotations + +import json +import subprocess +import sys +from pathlib import Path + +import pytest + +from test_native_child_replan_guard_cli import AGENT, GOAL, ROOT, TODO, TURN, _fixture + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_closed_replan_only_accepts_existing_native_operations( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, provider: str, +) -> None: + call, runtime, index = _fixture(tmp_path, monkeypatch, provider, True) + guard = call("quota", "should-run", "--codex-app", "--goal-id", GOAL, + "--agent-id", AGENT, "--turn-instance-id", TURN) + original = guard["heartbeat_receipt"]["settlement_identity"] + base = ("native-child", "--goal-id", GOAL, "--agent-id", AGENT, "--turn-instance-id", TURN) + + def reject(*args: str, reason: str) -> None: + result = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", + str(tmp_path / "registry.json"), "--runtime-root", str(runtime), "--format", "json", + *args], cwd=ROOT, capture_output=True, text=True, timeout=60) + assert result.returncode == 1, (result.stdout, result.stderr) + assert reason in json.loads(result.stdout)["error"] + + decision = (*base, "record", "--operation-id", "op-original", "--stage", "decision", + "--operation", "spawn", "--outcome", "started", "--entrypoint-id", "generic_host", + "--execute") + first = call(*decision) + added = call("todo", "add", "--goal-id", GOAL, "--role", "agent", "--claimed-by", AGENT, + "--text", "Independently validate the source artifact", "--task-class", "advancement_task", + "--action-kind", "validate", "--target-key", "independent-source-artifact", + "--operation-id", "native-late-source-successor", + "--replan-obligation-id", guard["autonomous_replan_obligation"]["obligation_id"]) + assert added["replan_transition"]["recorded"] is True + binding = ("--goal-id", GOAL, "--agent-id", AGENT, "--todo-id", TODO, "--turn-instance-id", TURN) + refreshed = call("refresh-state", *binding, "--classification", "bounded_replan_progress", + "--delivery-batch-scale", "single_surface", "--delivery-outcome", "outcome_progress", + "--vision-unchanged-reason", "The source validation remains open; its independent successor changes the path.", + "--no-global-sync", "--suppress-external-sinks") + assert refreshed["settlement_progress"]["state"] == "spend_required" + reject(*base, "record", "--operation-id", "op-new", "--stage", "decision", + "--operation", "spawn", "--outcome", "started", "--entrypoint-id", "generic_host", + "--execute", reason="open, work-admitted") + pending_replay = call(*decision) + assert pending_replay["appended"] is False + assert pending_replay["receipt"]["event_id"] == first["receipt"]["event_id"] + result = (*base, "record", "--operation-id", "op-original", "--stage", "result", + "--outcome", "completed", "--execute") + assert call(*result)["appended"] is True + spent = call("quota", "spend-slot", *binding, "--slots", "1", "--source", "heartbeat", "--execute") + assert spent["settlement_progress"]["state"] == "settled" + closed_index = index.read_bytes() + reject(*base, "record", "--operation-id", "op-new-skip", "--stage", "decision", + "--operation", "skip", "--outcome", "skipped", "--entrypoint-id", "generic_host", + "--reason-code", "parent_work_priority", "--execute", reason="open, work-admitted") + reject(*base, "record", "--operation-id", "op-unknown", "--stage", "result", + "--outcome", "completed", "--execute", reason="started") + review = (*base, "record", "--operation-id", "op-original", "--stage", "review", + "--outcome", "accepted", "--evidence-ref", "evidence-source", "--validation-ref", "validation-source", + "--execute") + accepted = call(*review) + assert accepted["appended"] is True + assert call(*review)["appended"] is False + assert call(*result)["appended"] is False + reject(*base, "record", "--operation-id", "op-original", "--stage", "result", + "--outcome", "failed", "--execute", reason="conflicting") + assert accepted["native_child_activity"]["parent_accepted_count"] == 1 + assert accepted["native_child_activity"]["quota_spend_slots"] == 0 + assert index.read_bytes() == closed_index + rows = [json.loads(line) for line in index.read_text().splitlines()] + assert sum(row.get("classification") == "quota_slot_spent" for row in rows) == 1 + assert call("todo", "list", "--goal-id", GOAL, "--todo-id", TODO)["todo"]["status"] == "open" + assert call("todo", "list", "--goal-id", GOAL, "--todo-id", added["todo_id"])["todo"]["status"] == "open" + assert original.get("todo_id") == TODO diff --git a/tests/control_plane/test_native_child_replan_guard_cli.py b/tests/control_plane/test_native_child_replan_guard_cli.py new file mode 100644 index 0000000000..465e910cad --- /dev/null +++ b/tests/control_plane/test_native_child_replan_guard_cli.py @@ -0,0 +1,120 @@ +"""Public reporting works in a lawful replan without launching or settling work.""" +from __future__ import annotations + +import json +import subprocess +import sys +from pathlib import Path + +import pytest +import loopx + +from canonical_authority_fixture import initialize_canonical_authority, isolate_sqlite_runtime +from test_replan_successor_durable_ack import history + +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.control_plane.goals.goal_vision import compact_goal_vision_packet, normalize_goal_vision_packet +from loopx.control_plane.todos.active_state_todo_parser import parse_active_state_todos +from loopx.rollout_event_log import load_rollout_events, rollout_event_log_path + + +GOAL = "native-child-replan-fixture" +AGENT = "generic-coordinator" +TURN = "turn-native-replan" +TODO = "todo_source_validation" +ROOT = Path(loopx.__file__).resolve().parents[1] + + +def _fixture(tmp_path: Path, monkeypatch: pytest.MonkeyPatch, provider: str, todo_bound: bool): + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + project, runtime = tmp_path / "project", tmp_path / "runtime" + project.mkdir() + state = project / "ACTIVE_GOAL_STATE.md" + state.write_text("---\nstatus: active\n---\n\n# Synthetic Goal\n\n## Agent Todo\n" + ( + "\n- [ ] [P1] Validate the original source.\n" + f" \n" + if todo_bound else "" + )) + index = runtime / "goals" / GOAL / "runs" / "index.jsonl" + index.parent.mkdir(parents=True) + evidence, report = index.parent / "source.json", index.parent / "source.md" + evidence.write_text(json.dumps({"ok": True, "fixture": "public-safe-native-replan"})) + report.write_text("# Synthetic source audit\n") + runs = history() + for row in runs: + row.update(agent_id=AGENT, json_path=str(evidence), markdown_path=str(report)) + runs[0]["agent_vision"] = compact_goal_vision_packet(normalize_goal_vision_packet({ + "goal_id": GOAL, "agent_id": AGENT, "state": "vision_drift_detected", + "vision_patch": {"acceptance_summary": "Independently validate the source.", + "advancement_policy": "repeat_until_closed"}, + }, goal_id=GOAL, agent_id=AGENT)) + index.write_text("".join(json.dumps(row) + "\n" for row in reversed(runs))) + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"common_runtime_root": str(runtime), "goals": [{ + "id": GOAL, "status": "active", "repo": str(project), "state_file": state.name, + "domain": "synthetic-replan", + "adapter": {"kind": "fixture_connected_delivery_v0", "status": "connected-delivery"}, + "quota": {"compute": 1.0, "window_hours": 24}, + "spawn_policy": {"mode": "multi_subagent", "allowed": True, "max_children": 6}, + "coordination": {"agent_model": "peer_v1", "registered_agents": [AGENT]}, + }]})) + todos = parse_active_state_todos(state.read_text(), item_limit=None)["agent_todos"]["items"] + initialize_canonical_authority(runtime, GOAL, build_todo_runtime_shadow_projection( + goal_id=GOAL, todos=todos, handoff_mode="soft_claim", leases=[], + ), state_path=state, provider=provider) + + def call(*args: str) -> dict: + result = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), + "--runtime-root", str(runtime), "--format", "json", *args], cwd=ROOT, + capture_output=True, text=True, timeout=60) + assert result.returncode == 0, (result.stdout, result.stderr) + return json.loads(result.stdout) + + return call, runtime, index + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("todo_bound", [True, False]) +def test_legal_replan_reports_native_child_without_settling_parent( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, provider: str, todo_bound: bool, +) -> None: + call, runtime, index = _fixture(tmp_path, monkeypatch, provider, todo_bound) + guard = call("quota", "should-run", "--codex-app", "--goal-id", GOAL, + "--agent-id", AGENT, "--turn-instance-id", TURN) + assert guard["decision"] == "autonomous_replan_required", guard + identity = guard["heartbeat_receipt"]["settlement_identity"] + assert identity.get("todo_id") == (TODO if todo_bound else None) + assert bool(identity.get("replan_obligation_id")) is not todo_bound + assert guard["interaction_contract"]["agent_channel"]["delivery_allowed"] is True + original_index = index.read_bytes() + base = ("native-child", "--goal-id", GOAL, "--agent-id", AGENT, + "--turn-instance-id", TURN) + decision = (*base, "record", "--operation-id", "op-independent-source", + "--stage", "decision", "--operation", "spawn", "--outcome", "started", + "--entrypoint-id", "generic_host", "--execute") + first = call(*decision) + assert first["appended"] is True + replay = call(*decision) + assert replay["appended"] is False + assert replay["receipt"]["event_id"] == first["receipt"]["event_id"] + call(*base, "record", "--operation-id", "op-independent-source", + "--stage", "result", "--outcome", "completed", "--execute") + reviewed = call(*base, "record", "--operation-id", "op-independent-source", + "--stage", "review", "--outcome", "accepted", "--evidence-ref", "evidence-source", + "--validation-ref", "validation-source", "--execute") + readback = call(*base, "read")["native_child_activity"] + assert readback == reviewed["native_child_activity"] + assert readback["observation"] == "coordinator_reported" + assert readback["host_attested"] is False + assert readback["parent_accepted_count"] == 1 + assert readback["quota_spend_slots"] == 0 + context = call("agent-context", "--goal-id", GOAL, "--agent-id", AGENT, + "--phase", "after_delegate_result", "--turn-instance-id", TURN) + assert context["native_child_activity"] == readback + assert context["host_receipts_observed"] is False + assert index.read_bytes() == original_index + events = load_rollout_events(rollout_event_log_path(runtime, GOAL)) + assert sum(row["event_kind"] == "native_child_decision" for row in events) == 1 + assert not any(row["event_kind"] in {"refresh_state", "quota_spend"} for row in events) diff --git a/tests/control_plane_ts/native_child_admission.test.ts b/tests/control_plane_ts/native_child_admission.test.ts new file mode 100644 index 0000000000..46a39c51b0 --- /dev/null +++ b/tests/control_plane_ts/native_child_admission.test.ts @@ -0,0 +1,146 @@ +import assert from "node:assert/strict"; +import {appendFile, mkdir, mkdtemp, rm, writeFile} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import test from "node:test"; + +import {nativeChildReportAdmission} from "../../loopx/control_plane/capabilities/native_child_admission.ts"; +import {settlementIdentity, settlementIdentityPayload} from "../../loopx/control_plane/effect_program.ts"; +import {QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA, readQuotaSettlement} from "../../loopx/control_plane/quota/settlement_readback.ts"; + +const proof = { + ok: true, should_run: true, must_attempt_work: true, + delivery_allowed: true, quiet_noop_allowed: false, +}; +const goalId = "native-report-fixture"; +const agentId = "generic-coordinator"; +const turnId = "turn-native-report"; +const todoId = "todo_source_validation"; + +test("committed work admission, not status spelling, admits a report", () => { + for (const status of ["normal_run", "turn_run_once", "autonomous_replan_required", "future_work_label"]) { + assert.deepEqual(nativeChildReportAdmission(proof, status, "effect-original", "open", false), { + schema_version: "native_child_report_admission_v0", + report_permission: "new_operation", reason_code: "turn_work_admitted", + settlement_effect_id: "effect-original", turn_phase: "open", + }); + } +}); + +test("late facts are allowed but no new operation once closeout has begun", () => { + for (const phase of ["settlement_pending", "settled"] as const) { + assert.equal(nativeChildReportAdmission(proof, "autonomous_replan_required", "original", phase, false) + .report_permission, "existing_operation_only"); + } + // A committed effect with a missing receipt is not permission for new work. + assert.equal(nativeChildReportAdmission(proof, "normal_run", "original", "open", true) + .report_permission, "existing_operation_only"); +}); + +test("legacy ordinary compatibility never overrides partial or negative facts", async (t) => { + for (const status of ["normal_run", "turn_run_once"]) { + assert.equal(nativeChildReportAdmission({}, status, "original", "open", false) + .report_permission, "new_operation"); + assert.equal(nativeChildReportAdmission({ok: true, should_run: true}, status, "original", "open", false) + .report_permission, "new_operation"); + } + for (const [label, facts] of [ + ["no work proof for replan", {}], + ["partial", {must_attempt_work: true}], + ["failure", {...proof, ok: false}], + ["wait", {...proof, should_run: false}], + ["no obligation", {...proof, must_attempt_work: false}], + ["no delivery", {...proof, delivery_allowed: false}], + ["quiet", {...proof, quiet_noop_allowed: true}], + ["truthy", {...proof, delivery_allowed: "true"}], + ["legacy failure", {ok: false}], + ["legacy wait", {should_run: false}], + ] as const) { + await t.test(label, () => { + const status = label === "no work proof for replan" ? "autonomous_replan_required" : "normal_run"; + assert.equal(nativeChildReportAdmission(facts, status, "original", "open", false) + .report_permission, "not_admitted"); + }); + } +}); + +async function fixture(replan = false) { + const root = await mkdtemp(join(tmpdir(), "loopx-native-report-admission-")); + const goalRoot = join(root, "goals", goalId); + await mkdir(join(goalRoot, "runs"), {recursive: true}); + const identity = settlementIdentity({goal_id: goalId, agent_id: agentId, + todo_id: replan ? null : todoId, turn_instance_id: turnId, + replan_obligation_id: replan ? "replan-0000000000000001" : null}); + const event = {schema_version: "loopx_rollout_event_v0", event_id: "event-original", + event_kind: "quota_should_run", goal_id: goalId, agent_id: agentId, run_id: turnId, + status: "autonomous_replan_required", details: {...proof, todo_id: identity.todo_id, + replan_obligation_id: identity.replan_obligation_id, settlement_effect_id: identity.effect_id}}; + const log = join(goalRoot, "rollout-event-log.jsonl"); + await writeFile(log, JSON.stringify(event) + "\n"); + await writeFile(join(goalRoot, "runs", "index.jsonl"), ""); + const request = {schema_version: QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA, + runtime_root: root, goal_id: goalId, agent_id: agentId, turn_instance_id: turnId, + todo_id: null, replan_obligation_id: null, infer_turn_instance_id: false, + allow_unbound_binding: false, resolve_original_binding: true}; + return {root, event, log, request, identity}; +} + +test("exact-Turn original binding is read without selecting or inferring new work", async () => { + for (const replan of [false, true]) { + const {root, request, identity} = await fixture(replan); + try { + const result = await readQuotaSettlement(request); + assert.equal(result.identity.result.failure, null); + assert.deepEqual(result.identity.result.value, settlementIdentityPayload(identity)); + assert.equal(result.native_child_admission.settlement_effect_id, identity.effect_id); + assert.equal(result.native_child_admission.report_permission, "new_operation"); + assert.equal(result.replay_phase, "open"); + } finally { + await rm(root, {recursive: true, force: true}); + } + } +}); + +test("exact binding resolution preserves explicit denials and conflicts", async (t) => { + for (const [label, mutate] of [ + ["other actor", {agent_id: "other-coordinator"}], + ["other Turn", {turn_instance_id: "other-turn"}], + ["other Todo", {todo_id: "todo_other_validation"}], + ["invalid Todo", {todo_id: "not-a-todo"}], + ["dual binding", {todo_id: todoId, replan_obligation_id: "replan-0000000000000002"}], + ] as const) { + await t.test(label, async () => { + const {root, request} = await fixture(); + try { + const result = await readQuotaSettlement({...request, ...mutate}); + assert.notEqual(result.identity.result.failure, null); + assert.equal(result.native_child_admission, null); + } finally {await rm(root, {recursive: true, force: true});} + }); + } + for (const [label, details] of [ + ["unbound", {...proof}], + ["effect mismatch", {...proof, todo_id: todoId, settlement_effect_id: "other-effect"}], + ["dual receipt", {...proof, todo_id: todoId, replan_obligation_id: "replan-0000000000000002"}], + ] as const) { + await t.test(label, async () => { + const {root, request, event, log} = await fixture(); + try { + await writeFile(log, JSON.stringify({...event, details}) + "\n"); + const result = await readQuotaSettlement(request); + assert.notEqual(result.identity.result.failure, null); + assert.equal(result.native_child_admission, null); + } finally {await rm(root, {recursive: true, force: true});} + }); + } + const {root, request, event, log} = await fixture(); + try { + await appendFile(log, JSON.stringify({...event, event_id: "event-conflict", + details: {...event.details, todo_id: "todo_other_validation", settlement_effect_id: "other-effect"}}) + "\n"); + assert.notEqual((await readQuotaSettlement(request)).identity.result.failure, null); + await assert.rejects(readQuotaSettlement({...request, resolve_original_binding: "yes"}), /must be a boolean/); + await assert.rejects(readQuotaSettlement({...request, infer_turn_instance_id: true}), /explicit Turn identity/); + // Existing callers still require an explicit binding by default. + assert.notEqual((await readQuotaSettlement({...request, resolve_original_binding: false})).identity.result.failure, null); + } finally {await rm(root, {recursive: true, force: true});} +});