From cbf8d1ff54ef991ca35d39fea1788e09b078c1a1 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Sun, 27 Sep 2026 10:17:41 +0800 Subject: [PATCH] fix(quota): retain auxiliary monitors after primary completion Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/control_plane/quota/cli_projection.py | 1 + loopx/control_plane/quota/monitor_poll.py | 12 + .../quota/monitor_poll_commit.ts | 64 ++++- loopx/quota.py | 172 +++++++------ .../test_monitor_observation_admission.py | 227 ++++++++++++++++++ .../auxiliary_monitor_settlement.test.ts | 134 +++++++++++ 6 files changed, 530 insertions(+), 80 deletions(-) create mode 100644 tests/control_plane_ts/auxiliary_monitor_settlement.test.ts diff --git a/loopx/control_plane/quota/cli_projection.py b/loopx/control_plane/quota/cli_projection.py index c977c9005c..10b8c4998c 100644 --- a/loopx/control_plane/quota/cli_projection.py +++ b/loopx/control_plane/quota/cli_projection.py @@ -134,6 +134,7 @@ "spend_after_validation", "spend_policy", "delivery_workspace_causality", + "settlement_resume_ref", ), } _RETAINED_MONITOR_POLL_RESPONSE_PLAN_FIELDS = ( diff --git a/loopx/control_plane/quota/monitor_poll.py b/loopx/control_plane/quota/monitor_poll.py index e98c87073c..20a14f04aa 100644 --- a/loopx/control_plane/quota/monitor_poll.py +++ b/loopx/control_plane/quota/monitor_poll.py @@ -698,6 +698,7 @@ def record_quota_monitor_poll_for_decision( task_lease_idempotency_key: str | None = None, task_lease_expected_version: int | None = None, use_current_task_lease: bool = False, + auxiliary_settlement_todo: Mapping[str, Any] | None = None, turn_instance_id: str | None = None, _index_lock_held: bool = False, status_reloader: Callable[[], dict[str, Any]] | None = None, @@ -750,6 +751,17 @@ def record_quota_monitor_poll_for_decision( registry_path=registry_path, runtime_root=runtime_root, ) + if auxiliary_settlement_todo is not None: + decision["auxiliary_settlement_todo"] = { + key: auxiliary_settlement_todo.get(key) + for key in ( + "todo_id", + "task_class", + "status", + "claimed_by", + "excluded_agents", + ) + } observation = _observation_packet( before=before, agent_id=agent_id, diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index 721a3c3292..8eb4e7010a 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -8,6 +8,7 @@ import { basename, dirname, join, resolve } from "node:path"; import { monitorSuccessorIntent, monitorSuccessorRoute } from "../scheduler/monitor_successor.ts"; import { normalizeTodoCapabilities } from "../todos/work_requirements.ts"; import { parseProjectionDelivery } from "../todos/projection_delivery.ts"; +import { readQuotaSettlement, QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA } from "./settlement_readback.ts"; import type { JsonObject } from "../effect_program.ts"; import { EffectRuntimeRequestError } from "../effect_runtime_errors.ts"; @@ -89,6 +90,7 @@ interface MonitorDecision extends JsonObject { vision_wait_state: JsonObject; due_monitor_candidates: JsonObject[]; registry_due_monitor: JsonObject; + auxiliary_settlement_todo: JsonObject | null; } interface MonitorObservation extends JsonObject { @@ -311,6 +313,8 @@ function decisionObject(value: unknown): MonitorDecision { decision.registry_due_monitor, "decision.registry_due_monitor", ), + auxiliary_settlement_todo: decision.auxiliary_settlement_todo == null ? null : + requiredObject(decision.auxiliary_settlement_todo, "decision.auxiliary_settlement_todo"), }; } @@ -537,10 +541,60 @@ interface Admission { external: boolean; } -function admission(request: MonitorRequest): Admission { +async function auxiliaryMonitorAllowed( + request: MonitorRequest, historicalAdmission: boolean, +): Promise { + const { decision, observation } = request; + const settlementTodo = observation.settlement_todo_id; + if (!settlementTodo || settlementTodo === observation.todo_id) return null; + const conflict = () => new EffectRuntimeRequestError( + "turn-scoped monitor-poll conflicts with the committed advancement settlement identity: " + + `settlement Todo ${settlementTodo}, observation Todo ${observation.todo_id}`, + "heartbeat_receipt_identity_conflict", + ); + const todo = decision.auxiliary_settlement_todo; + // A verified pending v1 receipt already admitted this exact observation. + // Earlier receipts predate the lifecycle-fact field; preserve their recovery + // basis, never reuse it for a new effect or a changed settlement binding. + if (historicalAdmission && todo === null) return null; + if (!request.runtime_root || !request.turn_instance_id || !decision.agent_id || + observation.actor_agent_id !== decision.agent_id || + todo === null || todo.todo_id !== settlementTodo || todo.task_class !== "advancement_task" || + !["open", "done"].includes(String(todo.status)) || + (todo.claimed_by != null && todo.claimed_by !== decision.agent_id) || + (todo.excluded_agents != null && (!Array.isArray(todo.excluded_agents) || + todo.excluded_agents.some(value => typeof value !== "string") || + todo.excluded_agents.includes(decision.agent_id)))) throw conflict(); + if (!historicalAdmission) { + const settlement = await readQuotaSettlement({ + schema_version: QUOTA_SETTLEMENT_READBACK_REQUEST_SCHEMA, + runtime_root: request.runtime_root, goal_id: request.goal_id, + agent_id: decision.agent_id, todo_id: settlementTodo, + turn_instance_id: request.turn_instance_id, replan_obligation_id: null, + infer_turn_instance_id: false, allow_unbound_binding: false, + }); + const progress = jsonObject(settlement.progress); + if (settlement.found !== true || jsonObject(jsonObject(settlement.identity)?.result)?.failure !== null || + progress?.schema_version !== "quota_settlement_progress_v0" || + progress.state === "identity_required" || progress.state === "settled") throw conflict(); + } + const monitor = decision.registry_due_monitor; + // Retain ordinary quota/due-work admission and capability/gate projections; + // lifecycle lookup must not become a second should-run bypass. + return !decision.requires_user_action && dueMonitorAllowed(decision, observation) && + monitor.due === true && candidateMatches(monitor, observation) && + (monitor.claimed_by == null || monitor.claimed_by === decision.agent_id); +} + +async function admission(request: MonitorRequest, historicalAdmission = false): Promise { const blocked = blockedSuccessorAllowed(request.decision); const external = externalMonitorAllowed(request.decision); - const due = dueMonitorAllowed(request.decision, request.observation); + const auxiliary = await auxiliaryMonitorAllowed(request, historicalAdmission); + const due = auxiliary ?? dueMonitorAllowed(request.decision, request.observation); + if (auxiliary === false) { + throw new EffectRuntimeRequestError("auxiliary monitor-poll requires its own due Monitor target", + "monitor_poll_admission_rejected"); + } if ( request.decision.effective_action !== EffectiveAction.MONITOR_QUIET_SKIP && !external && !due && !blocked @@ -1925,7 +1979,7 @@ export async function evaluateQuotaMonitorPollCommit( throw new EffectRuntimeRequestError("provider rejection recovery requires execute"); } if (request.phase === "event") { - const record = buildRecord(request, admission(request)); + const record = buildRecord(request, await admission(request)); return result( request, fingerprint, @@ -1946,7 +2000,7 @@ export async function evaluateQuotaMonitorPollCommit( if (!request.execute) { // A provider may mutate the Todo registry, so previews must pass admission // before returning a provider plan. - const allowed = admission(request); + const allowed = await admission(request); if (request.phase === "preflight") { const providerPlan = providerPlanFor(request); return result( @@ -2076,7 +2130,7 @@ export async function evaluateQuotaMonitorPollCommit( } let allowed: Admission; try { - allowed = admission(admittedRequest); + allowed = await admission(admittedRequest, existing?.schema_version === MONITOR_PENDING_ADMISSION_SCHEMA); } catch (error) { if (existing?.schema_version === QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA && error instanceof EffectRuntimeRequestError && error.code === "monitor_poll_admission_rejected") { diff --git a/loopx/quota.py b/loopx/quota.py index db1e8e02fb..e56f59edea 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -3,6 +3,7 @@ from collections.abc import Callable, Mapping from datetime import datetime, timedelta, timezone from pathlib import Path +import shlex from typing import Any from .control_plane import compact_control_plane_policy @@ -23,7 +24,6 @@ ReceiptBoundReplayPhase, ReceiptBoundTerminalPhase, ) -from .control_plane.quota.error_codes import HeartbeatReceiptIdentityConflictError from .control_plane.quota.goal_boundary import ( registry_goal_by_id as _registry_goal_by_id, ) @@ -34,7 +34,6 @@ from .control_plane.quota.monitor_poll import ( QUOTA_MONITOR_POLL_CLASSIFICATION as QUOTA_MONITOR_POLL_CLASSIFICATION, build_quota_monitor_poll_event as build_quota_monitor_poll_event, - find_quota_monitor_poll_turn, record_quota_monitor_poll_for_decision, resolve_due_monitor_candidate, ) @@ -98,8 +97,6 @@ CODEX_APP_SURFACE, ) from .control_plane.todos.contract import ( - TODO_TASK_CLASS_ADVANCEMENT, - TODO_TASK_CLASS_MONITOR, normalize_todo_claimed_by, normalize_todo_id, ) @@ -1138,87 +1135,31 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]: ) before = should_run(status_payload) + auxiliary_settlement_todo = None if ( normalized_receipt_todo_id and normalized_observation_todo_id and normalized_observation_todo_id != normalized_receipt_todo_id ): - selected = ( - before.get("selected_todo") - if isinstance(before.get("selected_todo"), Mapping) - else {} - ) - lane = ( - before.get("work_lane_contract") - if isinstance(before.get("work_lane_contract"), Mapping) - else {} - ) - summary = ( - before.get("agent_todo_summary") - if isinstance(before.get("agent_todo_summary"), Mapping) - else {} - ) - candidate_values = [ - *(lane.get("monitor_due_items") or []), - *(summary.get("monitor_due_items") or []), - ] - normalized_agent_id = normalize_todo_claimed_by(agent_id) - auxiliary_due_monitor = any( - isinstance(candidate, Mapping) - and normalize_todo_id(candidate.get("todo_id")) - == normalized_observation_todo_id - and candidate.get("task_class") == TODO_TASK_CLASS_MONITOR - and normalize_todo_claimed_by(candidate.get("claimed_by")) - in {None, normalized_agent_id} - for candidate in candidate_values - ) - auxiliary_registry_due = bool( - resolved_monitor - and normalize_todo_claimed_by(resolved_monitor.get("claimed_by")) - in {None, normalized_agent_id} - ) - existing_observation = ( - find_quota_monitor_poll_turn( - Path(str(raw_runtime_root)).expanduser(), + # Discovery/selection is an open-work projection, not the committed + # Turn's lifecycle authority. Read the exact bound record, including + # completed history; TS monitor admission owns its interpretation. + from .todos import list_goal_todos + + if registry_path is not None: + bound_records = list_goal_todos( + registry_path=registry_path, + runtime_root_arg=str(runtime_root) if runtime_root else None, goal_id=safe_goal_id, - agent_id=normalized_agent_id or "", - turn_instance_id=str(turn_instance_id or ""), - todo_id=normalized_observation_todo_id, - ) - if raw_runtime_root and agent_id and turn_instance_id - else None - ) - auxiliary_replay = bool( - isinstance(existing_observation, Mapping) - and normalize_todo_id(existing_observation.get("todo_id")) - == normalized_observation_todo_id - and normalize_todo_id( - existing_observation.get("settlement_todo_id") - ) - == normalized_receipt_todo_id - ) - auxiliary_observation_allowed = bool( - normalize_todo_id(selected.get("todo_id")) - == normalized_receipt_todo_id - and selected.get("task_class") == TODO_TASK_CLASS_ADVANCEMENT - and selected.get("selection_binding") == "heartbeat_receipt" - and ( - auxiliary_due_monitor - or auxiliary_registry_due - or auxiliary_replay - ) - ) - if not auxiliary_observation_allowed: - raise HeartbeatReceiptIdentityConflictError( - "turn-scoped monitor-poll Todo conflicts with the committed " - "heartbeat receipt: expected settlement Todo " - f"{normalized_receipt_todo_id}, requested observation Todo " - f"{normalized_observation_todo_id}" + todo_id=normalized_receipt_todo_id, + role="agent", ) + items = bound_records.get("todos") or [] + auxiliary_settlement_todo = items[0] if len(items) == 1 else None effective_todo_id = normalized_observation_todo_id or ( normalized_receipt_todo_id if not target_key else None ) - return record_quota_monitor_poll_for_decision( + result = record_quota_monitor_poll_for_decision( before, status_payload, goal_id=safe_goal_id, @@ -1230,6 +1171,7 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]: reason_summary=reason_summary, agent_id=agent_id, settlement_todo_id=normalized_receipt_todo_id, + auxiliary_settlement_todo=auxiliary_settlement_todo, todo_id=effective_todo_id, target_key=target_key, result_hash=result_hash, @@ -1251,6 +1193,86 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]: turn_instance_id=turn_instance_id, status_reloader=status_reloader, ) + continuation = result.get("turn_continuation") or {} + if ( + result.get("ok") is True + and continuation.get("settlement_binding_matches_observation") is False + ): + # The open-work readback may now select another Todo. Render the + # original receipt's closeout separately; never borrow that selection + # to construct this Turn's refresh/spend commands. + from .control_plane.agents.capability_gate import ( + runtime_capabilities_for_cli_projection, + ) + from .control_plane.quota.settlement import ( + attach_settlement_progress, + build_turn_scoped_cli_settlement_plan, + ) + + readback = read_heartbeat_settlement( + runtime_root, + goal_id=safe_goal_id, + agent_id=agent_id, + todo_id=normalized_receipt_todo_id, + turn_instance_id=turn_instance_id, + ) + if readback is None or readback.identity.value is None: + raise RuntimeError( + "auxiliary monitor receipt lost its original settlement readback" + ) + attach_settlement_progress( + result, readback, registry_path=registry_path, runtime_root=runtime_root, + ) + identity = readback.identity.value + prefix = "loopx" + if registry_path is not None: + prefix += f" --registry {shlex.quote(str(registry_path))}" + prefix += f" --runtime-root {shlex.quote(str(runtime_root))}" + scoped_args = "".join( + f" --available-capability {shlex.quote(capability)}" + for capability in runtime_capabilities_for_cli_projection( + available_capabilities + ) + ) + plan = build_turn_scoped_cli_settlement_plan( + goal_id=identity.goal_id, + agent_id=identity.agent_id, + todo_id=identity.todo_id, + turn_instance_id=identity.turn_instance_id, + replan_obligation_id=identity.replan_obligation_id, + command_prefix=prefix, + scoped_cli_args=scoped_args, + lifecycle_actor_args="", + quota_spend_source=readback.progress["quota_spend_source"], + ) + result["settlement_resume"] = { + "schema_version": "auxiliary_monitor_settlement_resume_v0", + "identity": identity.as_dict(), + "progress_ref": "$.settlement_progress", + "next_step": next( + ( + step for step in plan.as_dict()["ordered_steps"] + if step["kind"] == readback.progress["next_step"] + ), + None, + ), + "grants_new_delivery": False, + } + # A later discovery projection is not a second settlement plan. The + # top-level resume/readback above owns this original Turn's closeout. + after = result.get("after") or {} + cli = (after.get("interaction_contract") or {}).get("cli_channel") + if ( + isinstance(cli, dict) + and (after.get("selected_todo") or {}).get("todo_id") != identity.todo_id + ): + cli.pop("settlement_plan", None) + cli.pop("next_cli_actions", None) + cli.update( + spend_allowed_now=False, spend_after_validation=False, + settlement_resume_ref="$.settlement_resume", + ) + return result def build_quota_slot_void_preview( diff --git a/tests/control_plane/test_monitor_observation_admission.py b/tests/control_plane/test_monitor_observation_admission.py index e5aa1225ba..62787d0433 100644 --- a/tests/control_plane/test_monitor_observation_admission.py +++ b/tests/control_plane/test_monitor_observation_admission.py @@ -1,7 +1,10 @@ from __future__ import annotations +import json from pathlib import Path +import pytest + from tests.control_plane.test_quota_settlement_cli import ( AGENT_ID, DUE_MONITOR_TODO_ID, @@ -116,3 +119,227 @@ def test_receipt_bound_advancement_allows_one_auxiliary_due_monitor_receipt( assert final_guard["heartbeat_receipt"]["settlement_identity"]["todo_id"] == ( TODO_ID ) + + +@pytest.mark.parametrize( + ("provider", "writeback"), + [("legacy", False), ("legacy", True), ("file", True), ("sqlite", True)], +) +def test_completed_advancement_retains_auxiliary_monitor_admission( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, provider: str, writeback: bool +) -> None: + from canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, + ) + from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, + ) + from loopx.control_plane.coordination.local_authority import ( + read_canonical_todos_if_promoted, + ) + + if provider == "sqlite": + isolate_sqlite_runtime(tmp_path, monkeypatch) + project, runtime, registry_path = _write_fixture(tmp_path) + turn = "turn-completed-advancement-with-auxiliary-monitor" + binding = ( + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--turn-instance-id", + turn, + ) + capabilities = ( + "--available-capability", + "network", + "--available-capability", + "external_evidence_poll", + ) + guard = ( + "quota", + "should-run", + "--codex-app", + *binding, + *capabilities, + "--scan-path", + str(project), + ) + rc, first = _run_cli(registry_path, runtime, *guard) + assert rc == 0, first + _append_newly_due_monitor(project, watch_only=True) + rc, open_guard = _run_cli(registry_path, runtime, *guard) + assert rc == 0, open_guard + auxiliary = open_guard["interaction_contract"]["cli_channel"][ + "auxiliary_monitor_poll" + ] + poll_args = tuple( + "completed-primary-observation" + if token == "${LOOPX_MONITOR_RESULT_HASH:?}" + else token + for token in _projected_cli_args(auxiliary["command"], turn_instance_id=turn) + ) + ("--scan-path", str(project)) + rc, complete = _run_cli( + registry_path, + runtime, + "todo", + "complete", + "--goal-id", + GOAL_ID, + "--todo-id", + TODO_ID, + "--agent-id", + AGENT_ID, + "--claimed-by", + AGENT_ID, + "--evidence", + "read-only fixture delivery validated", + "--no-follow-up", + ) + assert rc == 0, complete + if provider != "legacy": + rc, listed = _run_cli( + registry_path, runtime, "todo", "list", "--goal-id", GOAL_ID + ) + assert rc == 0, listed + primary = next(todo for todo in listed["todos"] if todo["todo_id"] == TODO_ID) + assert primary["status"] == "done" + lease = { + "schema_version": "task_lease_v0", + "goal_id": GOAL_ID, + "todo_id": DUE_MONITOR_TODO_ID, + "owner": AGENT_ID, + "write_scopes": [], + "status": "active", + "idempotency_key": "own-monitor-execution", + "version": 3, + "lease_epoch": 2, + "acquired_at": "2026-01-01T00:00:00Z", + "expires_at": "2099-01-01T00:00:00Z", + } + projection = build_todo_runtime_shadow_projection( + goal_id=GOAL_ID, + todos=listed["todos"], + handoff_mode="hard_lease", + leases=[lease], + ) + initialize_canonical_authority( + runtime, + GOAL_ID, + projection, + state_path=project / f".codex/goals/{GOAL_ID}/ACTIVE_GOAL_STATE.md", + provider=provider, + ) + if writeback: + rc, refresh = _run_cli( + registry_path, + runtime, + "refresh-state", + *binding, + "--todo-id", + TODO_ID, + "--classification", + "fixture_delivery_validated", + "--delivery-batch-scale", + "single_surface", + "--delivery-outcome", + "outcome_progress", + "--delivery-workspace-path", + str(project), + "--no-global-sync", + "--suppress-external-sinks", + ) + assert rc == 0, refresh + rc, poll = _run_cli(registry_path, runtime, *poll_args) + assert rc == 0, (poll.get("error_code"), poll.get("reason"), poll.get("error")) + assert poll["todo_id"] == DUE_MONITOR_TODO_ID + assert poll["settlement_todo_id"] == TODO_ID + assert poll["turn_continuation"]["current_turn_settled"] is False + if provider != "legacy": + canonical = read_canonical_todos_if_promoted( + runtime_root=runtime, goal_id=GOAL_ID, include_leases=True + ) + assert canonical["source_authority"] == f"{provider}_v0" + observed = next( + todo + for todo in canonical["todos"] + if todo["todo_id"] == DUE_MONITOR_TODO_ID + ) + assert observed["result_hash"] == "completed-primary-observation" + assert poll["todo_writeback"]["lease_proof"] == { + "idempotency_key": "own-monitor-execution", + "expected_version": 3, + } + resume = poll["settlement_resume"] + assert resume["identity"]["todo_id"] == TODO_ID + assert resume["identity"]["turn_instance_id"] == turn + assert resume["grants_new_delivery"] is False + assert resume["progress_ref"] == "$.settlement_progress" + assert len(json.dumps(resume, ensure_ascii=False, indent=2)) < 2_000 + after_cli = poll["after"]["interaction_contract"]["cli_channel"] + assert after_cli["settlement_resume_ref"] == "$.settlement_resume" + assert after_cli["spend_after_validation"] is False + assert poll["settlement_progress"]["state"] == ( + "spend_required" if writeback else "writeback_required" + ) + rc, replay = _run_cli(registry_path, runtime, *poll_args) + assert rc == 0, replay + assert replay["replayed"] is True + assert _classification_count(runtime, "quota_monitor_poll") == 1 + assert _spend_run_count(runtime) == 0 + if provider != "legacy": + rc, released = _run_cli( + registry_path, + runtime, + "task-lease", + "release", + "--goal-id", + GOAL_ID, + "--todo-id", + DUE_MONITOR_TODO_ID, + "--owner", + AGENT_ID, + "--idempotency-key", + "own-monitor-execution", + "--expected-version", + "3", + ) + assert rc == 0, released + if not writeback: + assert resume["next_step"]["kind"] == "durable_writeback" + writeback_command = resume["next_step"]["command_template"] + args = tuple( + { + "": "fixture_delivery_validated", + "": "single_surface", + "": "outcome_progress", + }.get(token, token) + for token in _projected_cli_args(writeback_command, turn_instance_id=turn) + ) + rc, refresh = _run_cli( + registry_path, + runtime, + *args, + "--delivery-workspace-path", + str(project), + "--no-global-sync", + "--suppress-external-sinks", + ) + assert rc == 0, refresh + spend_command = refresh["settlement_owed"]["command"] + else: + assert resume["next_step"]["kind"] == "quota_spend" + spend_command = poll["settlement_owed"]["command"] + spend_args = _projected_cli_args(spend_command, turn_instance_id=turn) + for _ in range(2): + rc, spend = _run_cli(registry_path, runtime, *spend_args) + assert rc == 0, spend + assert spend["settlement_identity"]["todo_id"] == TODO_ID + assert _spend_run_count(runtime) == 1 + rc, replay = _run_cli(registry_path, runtime, *poll_args) + assert rc == 0, replay + assert replay["replayed"] is True + assert replay["settlement_progress"]["state"] == "settled" + assert replay["settlement_resume"]["next_step"] is None + assert _classification_count(runtime, "quota_monitor_poll") == 1 diff --git a/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts new file mode 100644 index 0000000000..fb9bdc83e6 --- /dev/null +++ b/tests/control_plane_ts/auxiliary_monitor_settlement.test.ts @@ -0,0 +1,134 @@ +import assert from "node:assert/strict"; +import { createHash } from "node:crypto"; +import { appendFile, mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { settlementIdentity, type JsonObject } from "../../loopx/control_plane/effect_program.ts"; +import { EffectRuntimeRequestError } from "../../loopx/control_plane/effect_runtime_errors.ts"; +import { + evaluateQuotaMonitorPollCommit, QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, +} from "../../loopx/control_plane/quota/monitor_poll_commit.ts"; + +const goal = "auxiliary-monitor-goal"; +const agent = "agent-monitor-fixture"; +const turn = "turn-completed-advancement"; +const primary = "todo_primary"; +const monitor = "todo_monitor"; +const identity = settlementIdentity({ goal_id: goal, agent_id: agent, turn_instance_id: turn, todo_id: primary }); + +function event(kind: string): JsonObject { + return { schema_version: "loopx_rollout_event_v0", event_id: `event-${kind}`, + event_kind: kind, goal_id: goal, agent_id: agent, run_id: turn, + details: { todo_id: primary, settlement_effect_id: identity.effect_id } }; +} + +async function fixture(t: test.TestContext): Promise<{ runtime: string; params: JsonObject }> { + const runtime = await mkdtemp(join(tmpdir(), "loopx-auxiliary-settlement-")); + t.after(() => rm(runtime, { recursive: true, force: true })); + await mkdir(join(runtime, "goals", goal, "runs"), { recursive: true }); + await writeFile(join(runtime, "goals", goal, "runs", "index.jsonl"), ""); + await writeFile(join(runtime, "goals", goal, "rollout-event-log.jsonl"), `${JSON.stringify(event("quota_should_run"))}\n`); + const candidate = { todo_id: monitor, task_class: "continuous_monitor", claimed_by: agent, + target_key: "public-watch", due: true }; + return { runtime, params: { + schema_version: QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + phase: "preflight", execute: true, effect_id: "auxiliary-observation", runtime_root: runtime, + goal_id: goal, source: "heartbeat", generated_at: "2026-09-01T00:00:00Z", + expected_index_digest: `sha256:${createHash("sha256").update("").digest("hex")}`, turn_instance_id: turn, + decision: { goal_id: goal, agent_id: agent, should_run: true, + normal_delivery_allowed: true, recovery_delivery_allowed: false, effective_action: "normal_run", + self_repair_allowed: false, capability_repair_allowed: false, workspace_repair_allowed: false, + state: "active", safe_bypass_allowed: false, safe_bypass_kind: null, blocked_action_scope: null, + compute: 1, window_hours: 24, slot_minutes: 1, spent_slots: 0, allowed_slots: 10, + recommended_action: "Continue the original settlement", reason: "Due observation", + requires_user_action: false, heartbeat_recommendation: {}, external_evidence_observation: null, + vision_wait_state: null, work_lane_contract: { must_attempt_work: true, obligation: "attempt_due_monitor" }, + due_monitor_candidates: [candidate], registry_due_monitor: candidate, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "done", + claimed_by: agent, excluded_agents: [] } }, + observation: { actor_agent_id: agent, settlement_todo_id: primary, todo_id: monitor, + target_key: "public-watch", result_hash: "unchanged", material_change: false, reason_summary: null, + cadence: "1h", next_due_at: "2026-09-01T01:00:00Z", next_agent_todo: null, + next_action_kind: null, next_task_repository: null, next_required_capabilities: [], + next_continuation_policy: null, next_target_key: null, next_user_todo: null, + next_user_task_class: null, next_claimed_by: null }, provider_receipt: null, status_reload_warning: null, + } }; +} + +async function settlePrimary(runtime: string): Promise { + await appendFile(join(runtime, "goals", goal, "rollout-event-log.jsonl"), + [event("refresh_state"), event("quota_spend")].map(row => JSON.stringify(row)).join("\n") + "\n"); + await appendFile(join(runtime, "goals", goal, "runs", "index.jsonl"), [ + { classification: "state_refreshed", delivery_outcome: "outcome_progress" }, + { classification: "quota_slot_spent" }, + ].map(row => JSON.stringify({ ...row, goal_id: goal, agent_id: agent, todo_id: primary, + turn_instance_id: turn, settlement_identity: identity })).join("\n") + "\n"); +} + +test("fresh auxiliary admission requires the exact unsettled advancement and ordinary due-work gates", async t => { + const cases: Array<[string, (params: JsonObject) => JsonObject, string]> = [ + ["other Turn", p => ({ ...p, turn_instance_id: "other-turn" }), "heartbeat_receipt_identity_conflict"], + ["other binding", p => ({ ...p, observation: { ...p.observation as JsonObject, settlement_todo_id: "todo_other" } }), "heartbeat_receipt_identity_conflict"], + ["missing lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, auxiliary_settlement_todo: null } }), "heartbeat_receipt_identity_conflict"], + ["blocked lifecycle", p => ({ ...p, decision: { ...p.decision as JsonObject, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "blocked" } } }), "heartbeat_receipt_identity_conflict"], + ["foreign primary", p => ({ ...p, decision: { ...p.decision as JsonObject, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "done", claimed_by: "peer" } } }), "heartbeat_receipt_identity_conflict"], + ["excluded actor", p => ({ ...p, decision: { ...p.decision as JsonObject, + auxiliary_settlement_todo: { todo_id: primary, task_class: "advancement_task", status: "done", excluded_agents: [agent] } } }), "heartbeat_receipt_identity_conflict"], + ["foreign observer", p => ({ ...p, observation: { ...p.observation as JsonObject, actor_agent_id: "peer" } }), "heartbeat_receipt_identity_conflict"], + ["no due target", p => ({ ...p, decision: { ...p.decision as JsonObject, registry_due_monitor: {} } }), "monitor_poll_admission_rejected"], + ["foreign monitor", p => ({ ...p, decision: { ...p.decision as JsonObject, + registry_due_monitor: { ...(p.decision as JsonObject).registry_due_monitor as JsonObject, claimed_by: "peer" } } }), "monitor_poll_admission_rejected"], + ["ordinary gate closed", p => ({ ...p, decision: { ...p.decision as JsonObject, + work_lane_contract: { must_attempt_work: false }, should_run: false } }), "monitor_poll_admission_rejected"], + ["user action", p => ({ ...p, decision: { ...p.decision as JsonObject, requires_user_action: true } }), "monitor_poll_admission_rejected"], + ]; + for (const [name, change, code] of cases) { + await t.test(name, async st => { + const { runtime, params } = await fixture(st); + const changed = change(params); + await assert.rejects(evaluateQuotaMonitorPollCommit(changed), error => { + assert.ok(error instanceof EffectRuntimeRequestError); + assert.equal(error.code, code); + if (code === "heartbeat_receipt_identity_conflict") { + const observation = changed.observation as JsonObject; + assert.ok(error.message.includes(String(observation.settlement_todo_id))); + assert.ok(error.message.includes(String(observation.todo_id))); + } + return true; + }); + await assert.rejects(readFile(join(runtime, "goals", goal, "runs", ".transactions", "quota-monitor-poll")), + { code: "ENOENT" }); + assert.equal(await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8"), ""); + }); + } +}); + +test("completed primary preserves an admitted pending effect across settlement, but cannot admit a new one", async t => { + const { runtime, params } = await fixture(t); + const preflight = await evaluateQuotaMonitorPollCommit(params); + assert.equal(preflight.status, "provider_required", JSON.stringify(preflight)); + await settlePrimary(runtime); + const receipt = { schema_version: "monitor_poll_todo_writeback_v0", monitor_effect_id: params.effect_id, + goal_id: goal, todo_id: monitor, target_key: "public-watch", result_hash: "unchanged", + dry_run: false, material_change: false, material_change_generation: 0, consecutive_no_change: 1, + last_checked_at: params.generated_at, next_due_at: "2026-09-01T01:00:00Z", cadence: "1h", + todo_update: { ok: true }, next_todos: [], successor_receipts: [] }; + const postBusiness = { ...params, phase: "commit", provider_receipt: receipt, + decision: { ...params.decision as JsonObject, registry_due_monitor: {}, auxiliary_settlement_todo: null, + due_monitor_candidates: [], work_lane_contract: { must_attempt_work: false }, should_run: false } }; + assert.equal((await evaluateQuotaMonitorPollCommit(postBusiness)).status, "written"); + assert.equal((await evaluateQuotaMonitorPollCommit(postBusiness)).status, "replayed"); + await assert.rejects(evaluateQuotaMonitorPollCommit({ ...params, effect_id: "new-after-settlement" }), error => { + assert.ok(error instanceof EffectRuntimeRequestError); + assert.equal(error.code, "heartbeat_receipt_identity_conflict"); + return true; + }); + const rows = (await readFile(join(runtime, "goals", goal, "runs", "index.jsonl"), "utf8")) + .trim().split("\n").map(line => JSON.parse(line)); + assert.equal(rows.filter(row => row.classification === "quota_monitor_poll").length, 1); + assert.equal(rows.filter(row => row.classification === "quota_slot_spent").length, 1); +});