From d068c2599c8d65798b7774b27fafeed60f4891c4 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 18 Sep 2026 01:27:56 +0800 Subject: [PATCH 1/4] feat(monitor): bind leased observations to recoverable quota settlement Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- loopx/cli_commands/quota_monitor_poll.py | 2 + loopx/cli_commands/quota_request.py | 13 ++ .../coordination/local_authority_runtime.ts | 14 +- .../coordination/task_lease_proof.ts | 60 +++++++++ .../coordination/todo_monitor_poll.ts | 40 +++++- .../control_plane/coordination/todo_update.ts | 26 +--- loopx/control_plane/quota/monitor_poll.py | 20 ++- .../quota/monitor_poll_commit.ts | 126 ++++++++++++++---- .../scheduler/monitor_poll_writeback.py | 10 ++ .../scheduler/provider_monitor_poll.py | 5 +- .../control_plane/todos/active_state_todos.py | 5 +- loopx/quota.py | 4 + 12 files changed, 265 insertions(+), 60 deletions(-) create mode 100644 loopx/control_plane/coordination/task_lease_proof.ts diff --git a/loopx/cli_commands/quota_monitor_poll.py b/loopx/cli_commands/quota_monitor_poll.py index 185e479155..2e46f8a371 100644 --- a/loopx/cli_commands/quota_monitor_poll.py +++ b/loopx/cli_commands/quota_monitor_poll.py @@ -83,6 +83,8 @@ def record_quota_monitor_poll_for_cli( next_user_todo=args.next_user_todo, next_user_task_class=args.next_user_task_class, next_claimed_by=args.next_claimed_by, + task_lease_idempotency_key=getattr(args, "task_lease_idempotency_key", None), + task_lease_expected_version=getattr(args, "task_lease_expected_version", None), turn_instance_id=turn_instance_id, receipt_bound_todo_id=_receipt_bound_monitor_todo_id( args, diff --git a/loopx/cli_commands/quota_request.py b/loopx/cli_commands/quota_request.py index 01902a9e89..99fe35f307 100644 --- a/loopx/cli_commands/quota_request.py +++ b/loopx/cli_commands/quota_request.py @@ -90,11 +90,24 @@ def register_quota_monitor_poll_request_arguments( "that must not block the bound agent lane." ), ) + quota_parser.add_argument("--task-lease-idempotency-key", + help="Current Monitor execution key for canonical quota monitor-poll; requires --task-lease-expected-version.") + quota_parser.add_argument("--task-lease-expected-version", type=int, + help="Current Monitor lease version, checked atomically with observation and successors; never renews the lease.") quota_parser.add_argument("--next-claimed-by", help="Registered agent id to claim the `--next-agent-todo` follow-up.") def validate_quota_command_request(args: argparse.Namespace) -> None: command = args.quota_command + lease_key = getattr(args, "task_lease_idempotency_key", None) + lease_version = getattr(args, "task_lease_expected_version", None) + if lease_key is not None or lease_version is not None: + if command != "monitor-poll": + raise QuotaCommandValidationError("task lease proof is only valid with quota monitor-poll") + if not lease_key or lease_version is None or lease_version < 1 or lease_version > 9007199254740991: + raise QuotaCommandValidationError("Monitor lease proof requires an execution key and positive safe-integer version") + if not (args.todo_id or args.target_key): + raise QuotaCommandValidationError("Monitor lease proof requires --todo-id or --target-key") begin_turn = bool(getattr(args, "begin_turn", False)) if command not in {"status", "plan"} and not args.goal_id: raise QuotaCommandValidationError( diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index e66c140496..6914145692 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -1,3 +1,4 @@ +import {decodeTaskLeaseProof} from "./task_lease_proof.ts"; import {COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA} from "./todo_archive.ts"; import {readCoordinationOwnership} from "./ownership_observation.ts"; import {executeTodoContinuation} from "./todo_continuation.ts"; @@ -6,7 +7,7 @@ import { ShadowManagementError } from "./shadow_management.ts"; import { isAbsolute, join } from "node:path"; import type { JsonObject } from "../effect_program.ts"; -import {executeCoordinationMonitorPoll, COORDINATION_MONITOR_POLL_REQUEST_SCHEMA, +import {executeCoordinationMonitorPoll, COORDINATION_MONITOR_POLL_REQUEST_SCHEMA, COORDINATION_LEASED_MONITOR_POLL_REQUEST_SCHEMA, COORDINATION_MONITOR_POLL_RESULT_SCHEMA} from "./todo_monitor_poll.ts"; import { requireJsonObject } from "../runtime_decode.ts"; import { @@ -112,7 +113,15 @@ export async function pollLocalCoordinationMonitor(value: unknown, const evidence = {source_authority: "file_v0", decision_read_from_provider: true, legacy_fallback_used: false}; try { const input = requireJsonObject(value, "local Monitor poll request"); - if (input.schema_version !== COORDINATION_MONITOR_POLL_REQUEST_SCHEMA) throw new TypeError("Monitor poll schema mismatch"); + if (input.schema_version !== COORDINATION_MONITOR_POLL_REQUEST_SCHEMA && + input.schema_version !== COORDINATION_LEASED_MONITOR_POLL_REQUEST_SCHEMA) throw new TypeError("Monitor poll schema mismatch"); + if (input.lease_proof != null && input.schema_version !== COORDINATION_LEASED_MONITOR_POLL_REQUEST_SCHEMA) { + throw new TypeError("lease-backed Monitor poll requires request v1"); + } + const proof = decodeTaskLeaseProof(input.lease_proof); + if (input.schema_version === COORDINATION_LEASED_MONITOR_POLL_REQUEST_SCHEMA && !proof) { + throw new TypeError("Monitor poll request v1 requires lease_proof"); + } const root = runtimeRoot(input.runtime_root); const goalId = requireAuthorityStoreId(input.goal_id, "goal id"); if (!Array.isArray(input.registered_agents)) throw new TypeError("registered_agents must be an array"); @@ -126,6 +135,7 @@ export async function pollLocalCoordinationMonitor(value: unknown, registered_agents: registered, dry_run: input.dry_run as boolean, observation: requireJsonObject(input.observation, "Monitor observation"), intent: requireJsonObject(input.intent, "Monitor successor intent"), + lease_proof: proof, now: new Date(), }), ...evidence}; }); } catch (error) { diff --git a/loopx/control_plane/coordination/task_lease_proof.ts b/loopx/control_plane/coordination/task_lease_proof.ts new file mode 100644 index 0000000000..f6fb157cf1 --- /dev/null +++ b/loopx/control_plane/coordination/task_lease_proof.ts @@ -0,0 +1,60 @@ +/** Current nonterminal mutation proof. Admission stays in the lifecycle owner; + * this boundary never acquires, renews, releases or transfers a lease. */ +import type {JsonObject} from "../effect_program.ts"; +import {requireJsonObject} from "../runtime_decode.ts"; +import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; +import {parseIsoTimestamp} from "../runtime_timestamp.ts"; +import {leaseEpoch} from "../work_items/task_lease_acquire.ts"; +import {evaluateCoordinationTerminalFence, COORDINATION_TERMINAL_FENCE_REQUEST_SCHEMA} from "./todo_lifecycle_decision.ts"; + +export interface TaskLeaseProof { + idempotency_key: string; + expected_version: number; +} + +export function decodeTaskLeaseProof(value: unknown): TaskLeaseProof | null { + if (value == null) return null; + const proof = requireJsonObject(value, "lease_proof"); + if (Object.keys(proof).some(key => key !== "idempotency_key" && key !== "expected_version")) { + throw new EffectRuntimeRequestError("lease_proof accepts only idempotency_key and expected_version"); + } + if (typeof proof.idempotency_key !== "string" || !proof.idempotency_key.trim() || + proof.idempotency_key !== proof.idempotency_key.trim() || + typeof proof.expected_version !== "number" || !Number.isSafeInteger(proof.expected_version) || proof.expected_version < 1) { + throw new EffectRuntimeRequestError("lease_proof requires an unpadded execution key and positive safe-integer version"); + } + return {idempotency_key: proof.idempotency_key, expected_version: proof.expected_version}; +} + +export function evaluateCanonicalTaskLeaseProof(input: { + todo: JsonObject; lease: JsonObject | undefined; handoff_mode: string; + actor_agent_id: string | null; registered_agents: readonly string[]; + lease_idempotency_key: string | null; lease_expected_version: number | null; now: Date; +}) { + const {lease} = input; + if (!(input.now instanceof Date) || !Number.isFinite(input.now.valueOf())) { + throw new EffectRuntimeRequestError("lease proof requires a valid transaction clock"); + } + const expires = lease === undefined ? null : + typeof lease.expires_at === "string" ? parseIsoTimestamp(lease.expires_at) : null; + if (lease?.status === "active" && expires === null) { + throw new EffectRuntimeRequestError("active lease expiry is invalid"); + } + const decision = evaluateCoordinationTerminalFence({ + schema_version: COORDINATION_TERMINAL_FENCE_REQUEST_SCHEMA, + todo: input.todo, registered_agents: input.registered_agents, + actor_agent_id: input.actor_agent_id, + // Retained execution lineage must not become an unfenced metadata edit. + handoff_mode: lease !== undefined ? "hard_lease" : input.handoff_mode, + lease: lease === undefined ? null : {...lease, present: true, + active: lease.status === "active" && expires !== null && expires > input.now, + lease_epoch: leaseEpoch(lease)}, + lease_idempotency_key: input.lease_idempotency_key, + lease_expected_version: input.lease_expected_version, + allow_user_gate_auto_acquire: false, delegated_authority: false, + require_active_when_fence_supplied: true, + }); + // Strip the terminal owner's release proposal: callers receive only a + // decision, so an observation cannot accidentally apply terminal effects. + return {outcome: decision.outcome, code: decision.code}; +} diff --git a/loopx/control_plane/coordination/todo_monitor_poll.ts b/loopx/control_plane/coordination/todo_monitor_poll.ts index cc9027e143..ddb968cd22 100644 --- a/loopx/control_plane/coordination/todo_monitor_poll.ts +++ b/loopx/control_plane/coordination/todo_monitor_poll.ts @@ -12,8 +12,12 @@ import {planMonitorSuccessor, selectMonitorTodo, MONITOR_SUCCESSOR_REQUEST_SCHEM import {optionalNonEmptyString, requireBoolean} from "../runtime_decode.ts"; import {planTodoAuthoringScope, TODO_AUTHORING_SCOPE_REQUEST_SCHEMA} from "../todos/authoring_scope.ts"; import {CoordinationCommandReceipt} from "./command_receipt.ts"; +import {decodeTaskLeaseProof, evaluateCanonicalTaskLeaseProof, type TaskLeaseProof} from "./task_lease_proof.ts"; +import {HANDOFF_MODES} from "./handoff_mode_policy.ts"; +import {requireStringLiteral} from "../runtime_decode.ts"; export const COORDINATION_MONITOR_POLL_REQUEST_SCHEMA = "loopx_coordination_monitor_poll_request_v0"; +export const COORDINATION_LEASED_MONITOR_POLL_REQUEST_SCHEMA = "loopx_coordination_monitor_poll_request_v1"; export const COORDINATION_MONITOR_POLL_RESULT_SCHEMA = "loopx_coordination_monitor_poll_result_v0"; const RECEIPT_SCHEMA = "loopx_coordination_monitor_poll_receipt_v0"; @@ -25,6 +29,9 @@ export interface CoordinationMonitorPollInput { dry_run: boolean; observation: JsonObject; intent: JsonObject; + lease_proof?: TaskLeaseProof | null; + /** Authority clock supplied by the runtime, never observation.generated_at. */ + now?: Date; } function failure(reason_code: string, reason: string): JsonObject & {schema_version: typeof COORDINATION_MONITOR_POLL_RESULT_SCHEMA} { @@ -42,18 +49,27 @@ function monitorReceipt(input: CoordinationMonitorPollInput, hash: string) { !Array.isArray(writeback.next_todos)) { throw new AuthorityStoreProtocolError("Monitor receipt writeback identity or successors invalid"); } + let proof: TaskLeaseProof | null; + try { proof = decodeTaskLeaseProof(writeback.lease_proof); } + catch { throw new AuthorityStoreProtocolError("Monitor receipt lease proof is malformed"); } + if (canonicalAuthoritySha256(proof ?? {}) !== canonicalAuthoritySha256(input.lease_proof ?? {})) { + throw new AuthorityStoreProtocolError("Monitor receipt belongs to a different lease proof"); + } return {fields: {writeback: {...writeback, provider_replayed: phase === "replayed"}}, changed: true}; }}); } -function normalize(raw: CoordinationMonitorPollInput): CoordinationMonitorPollInput { +type NormalizedMonitorPollInput = CoordinationMonitorPollInput & {now: Date; lease_proof: TaskLeaseProof | null}; + +function normalize(raw: CoordinationMonitorPollInput): NormalizedMonitorPollInput { const input = {...raw, goal_id: requireAuthorityStoreId(raw.goal_id, "goal id"), operation_id: requireAuthorityStoreId(raw.operation_id, "operation id"), actor_agent_id: raw.actor_agent_id == null ? null : normalizeTodoAgent(raw.actor_agent_id, "actor_agent_id"), registered_agents: normalizeRegisteredTodoAgents(raw.registered_agents), dry_run: requireBoolean(raw.dry_run, "dry_run"), observation: canonicalAuthorityObject(raw.observation, "Monitor observation"), - intent: canonicalAuthorityObject(raw.intent, "Monitor successor intent")}; + intent: canonicalAuthorityObject(raw.intent, "Monitor successor intent"), + lease_proof: decodeTaskLeaseProof(raw.lease_proof), now: raw.now ?? new Date()}; const allowed = new Set(["todo_id", "target_key", "generated_at", "result_hash", "material_change", "cadence", "next_due_at", "reason_summary"]); for (const key of Object.keys(input.observation)) if (!allowed.has(key)) throw new Error(`unsupported Monitor observation field: ${key}`); const intentFields = new Set(["next_agent_todo", "next_action_kind", "next_task_repository", "next_required_capabilities", @@ -62,7 +78,7 @@ function normalize(raw: CoordinationMonitorPollInput): CoordinationMonitorPollIn return input; } -function planWriteback(input: CoordinationMonitorPollInput, head: JsonObject) { +function planWriteback(input: NormalizedMonitorPollInput, head: JsonObject) { const indexed = indexCoordinationProjection(head, input.goal_id); validateCoordinationTodoReadModel(head, input.goal_id); const observation = input.observation; @@ -75,7 +91,17 @@ function planWriteback(input: CoordinationMonitorPollInput, head: JsonObject) { if (Array.isArray(monitor.excluded_agents) && monitor.excluded_agents.includes(actor)) throw new Error("Monitor actor is excluded"); if (monitor.bound_agent && monitor.bound_agent !== actor) throw new Error("Monitor bound agent mismatch"); if (monitor.claimed_by && monitor.claimed_by !== actor) throw new Error("Monitor claim owner mismatch"); - if (indexed.leases.has(String(monitor.todo_id))) throw new Error("Monitor observation with a lease is not supported; no state was written"); + const lease = indexed.leases.get(String(monitor.todo_id)); + const mode = requireStringLiteral(head.handoff_mode ?? "legacy", HANDOFF_MODES, "canonical handoff_mode"); + if (lease !== undefined || mode === "hard_lease" || input.lease_proof != null) { + if (mode === "soft_claim") throw new Error("soft_claim forbids lease-backed Monitor observation"); + const fence = evaluateCanonicalTaskLeaseProof({todo: monitor, lease, handoff_mode: mode, + actor_agent_id: actor, registered_agents: input.registered_agents, + lease_idempotency_key: input.lease_proof?.idempotency_key ?? null, + lease_expected_version: input.lease_proof?.expected_version ?? null, now: input.now}); + if (fence.outcome !== "apply") throw new Error(`Monitor requires current lease proof: ${fence.code}`); + if (lease !== undefined && monitor.claimed_by !== actor) throw new Error("Leased Monitor observation requires the current claim owner"); + } const successorPlan = planMonitorSuccessor({schema_version: MONITOR_SUCCESSOR_REQUEST_SCHEMA, todo_id: monitor.todo_id, result_hash: observation.result_hash, source_task_repository: monitor.task_repository ?? null, intent: {...input.intent, material_change: observation.material_change}}); @@ -143,6 +169,7 @@ function planWriteback(input: CoordinationMonitorPollInput, head: JsonObject) { consecutive_no_change: transition.consecutive_no_change, last_checked_at: observation.generated_at, next_due_at: transition.next_due_at ?? null, cadence: transition.cadence || null, todo_update: {ok: true, todo_id: monitor.todo_id, monitor_poll_transition: transition}, + ...(input.lease_proof ? {lease_proof: {...input.lease_proof}} : {}), next_todos: nextTodos, successor_receipts: nextTodos.map(todo => Object.fromEntries( receiptFields.filter(key => todo[key] != null).map(key => [key, todo[key]]))), provider_replayed: false}; return {mutations, writeback}; @@ -150,12 +177,13 @@ function planWriteback(input: CoordinationMonitorPollInput, head: JsonObject) { export async function executeCoordinationMonitorPoll(store: AuthorityStore, raw: CoordinationMonitorPollInput): Promise { - let input: CoordinationMonitorPollInput; + let input: NormalizedMonitorPollInput; try { input = normalize(raw); } catch (error) { return failure("invalid_monitor_poll_request", String(error)); } // Original wire identity, before any normalization/default route inference. const hash = canonicalAuthoritySha256({goal_id: input.goal_id, observation: input.observation, - intent: input.intent, actor_agent_id: input.actor_agent_id, dry_run: input.dry_run}); + intent: input.intent, actor_agent_id: input.actor_agent_id, dry_run: input.dry_run, + ...(input.lease_proof ? {lease_proof: input.lease_proof} : {})}); const receipt = monitorReceipt(input, hash); const previous = await receipt.read(store); if (previous) return previous; diff --git a/loopx/control_plane/coordination/todo_update.ts b/loopx/control_plane/coordination/todo_update.ts index 535972e7d2..eb92b64ab5 100644 --- a/loopx/control_plane/coordination/todo_update.ts +++ b/loopx/control_plane/coordination/todo_update.ts @@ -17,12 +17,10 @@ import { } from "./coordination_projection.ts"; import { normalizeRegisteredTodoAgents, normalizeTodoAgent } from "./todo_agents.ts"; -import { evaluateCoordinationTerminalFence, COORDINATION_TERMINAL_FENCE_REQUEST_SCHEMA, - evaluateCoordinationTodoMutationDecision, +import { evaluateCoordinationTodoMutationDecision, COORDINATION_TODO_MUTATION_DECISION_REQUEST_SCHEMA } from "./todo_lifecycle_decision.ts"; -import { leaseEpoch } from "../work_items/task_lease_acquire.ts"; -import { parseIsoTimestamp } from "../runtime_timestamp.ts"; +import {evaluateCanonicalTaskLeaseProof} from "./task_lease_proof.ts"; import { normalizeNativePlanningIntent, planNativeTodoUpdate } from "../todos/native_update_plan.ts"; import { CoordinationCommandReceipt } from "./command_receipt.ts"; @@ -226,24 +224,10 @@ function targetRejection( if (lease !== undefined || mode === "hard_lease" || input.lease_idempotency_key != null || input.lease_expected_version != null) { try { - const expires = lease === undefined ? null : - typeof lease.expires_at === "string" ? parseIsoTimestamp(lease.expires_at) : null; - if (lease?.status === "active" && expires === null) { - return failure("invalid_coordination_projection", "active lease expiry is invalid"); - } - const fence = evaluateCoordinationTerminalFence({ - schema_version: COORDINATION_TERMINAL_FENCE_REQUEST_SCHEMA, - todo, registered_agents: input.registered_agents, actor_agent_id: input.actor_agent_id, - // A historical lease never licenses an unfenced edit. No acquisition or override. - handoff_mode: lease !== undefined ? "hard_lease" : mode, - lease: lease === undefined ? null : {...lease, present: true, - active: lease.status === "active" && expires !== null && expires > input.now, - lease_epoch: leaseEpoch(lease)}, + const fence = evaluateCanonicalTaskLeaseProof({todo, lease, handoff_mode: mode, + registered_agents: input.registered_agents, actor_agent_id: input.actor_agent_id, lease_idempotency_key: input.lease_idempotency_key ?? null, - lease_expected_version: input.lease_expected_version ?? null, - allow_user_gate_auto_acquire: false, delegated_authority: false, - require_active_when_fence_supplied: true, - }); + lease_expected_version: input.lease_expected_version ?? null, now: input.now}); if (fence.outcome !== "apply") { return failure(String(fence.code), "Todo update requires the current active lease execution proof"); } diff --git a/loopx/control_plane/quota/monitor_poll.py b/loopx/control_plane/quota/monitor_poll.py index 7d047f61c5..fa136205ff 100644 --- a/loopx/control_plane/quota/monitor_poll.py +++ b/loopx/control_plane/quota/monitor_poll.py @@ -212,8 +212,14 @@ def _observation_packet( next_user_todo: str | None, next_user_task_class: str | None, next_claimed_by: str | None, + task_lease_idempotency_key: str | None = None, + task_lease_expected_version: int | None = None, ) -> dict[str, Any]: + proof = ({"idempotency_key": task_lease_idempotency_key, + "expected_version": task_lease_expected_version} + if task_lease_idempotency_key is not None or task_lease_expected_version is not None else None) return { + **({"lease_proof": proof} if proof is not None else {}), "actor_agent_id": normalize_todo_claimed_by(agent_id) or quota_decision_agent_id(before), "settlement_todo_id": settlement_todo_id, @@ -278,7 +284,8 @@ def _request( status_reload_warning: Mapping[str, Any] | None = None, ) -> dict[str, Any]: return { - "schema_version": QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + "schema_version": ("loopx_quota_monitor_poll_commit_request_v1" if observation.get("lease_proof") is not None + else QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA), "phase": phase, "effect_id": effect_id, "runtime_root": str(runtime_root) if runtime_root is not None else None, @@ -652,6 +659,8 @@ def _provider_writeback( next_user_task_class=plan.get("next_user_task_class"), next_claimed_by=plan.get("next_claimed_by"), agent_id=plan.get("agent_id"), + task_lease_idempotency_key=(plan.get("lease_proof") or {}).get("idempotency_key"), + task_lease_expected_version=(plan.get("lease_proof") or {}).get("expected_version"), ) if not isinstance(result, dict): raise TypeError("monitor Todo provider returned no writeback receipt") @@ -686,6 +695,8 @@ def record_quota_monitor_poll_for_decision( next_user_todo: str | None = None, next_user_task_class: str | None = None, next_claimed_by: str | None = None, + task_lease_idempotency_key: str | None = None, + task_lease_expected_version: int | None = None, turn_instance_id: str | None = None, _index_lock_held: bool = False, status_reloader: Callable[[], dict[str, Any]] | None = None, @@ -746,6 +757,8 @@ def record_quota_monitor_poll_for_decision( next_user_todo=next_user_todo, next_user_task_class=next_user_task_class, next_claimed_by=next_claimed_by, + task_lease_idempotency_key=task_lease_idempotency_key, + task_lease_expected_version=task_lease_expected_version, ) generated_at = _now_local() @@ -835,13 +848,16 @@ def transact() -> tuple[dict[str, Any], dict[str, Any]]: else: native, after_status = transact() except ValueError as exc: - return failure( + payload = failure( str(exc), include_capability_retry=( isinstance(exc, _NativeMonitorPollRejected) and exc.diagnostic_code == "monitor_poll_admission_rejected" ), ) + if isinstance(exc, _NativeMonitorPollRejected): + payload["error_code"] = exc.diagnostic_code + return payload if native.get("status") == "conflict": payload = failure( diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index 19ea7e1750..986388e99c 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -1,3 +1,4 @@ +import {decodeTaskLeaseProof, type TaskLeaseProof} from "../coordination/task_lease_proof.ts"; import { EffectiveAction, type QuotaEffectiveActionValue } from "./effective_action.generated.ts"; import { AgentScopeFrontierAction } from "../agents/agent_scope_frontier.generated.ts"; import { createHash } from "node:crypto"; @@ -28,14 +29,17 @@ import { export const QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA = "loopx_quota_monitor_poll_commit_request_v0"; +export const QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA = "loopx_quota_monitor_poll_commit_request_v1"; export const QUOTA_MONITOR_POLL_COMMIT_RESULT_SCHEMA = "loopx_quota_monitor_poll_commit_result_v0"; export const QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA = "quota_monitor_poll_commit_receipt_v0"; +const MONITOR_PENDING_ADMISSION_SCHEMA = "quota_monitor_poll_pending_admission_v1"; export const QUOTA_MONITOR_POLL_CLASSIFICATION = "quota_monitor_poll"; const MONITOR_TARGET_SCHEMA = "quota_monitor_target_v0"; const MONITOR_TODO_PROVIDER_PLAN_SCHEMA = "monitor_poll_todo_provider_plan_v0"; +const LEASED_MONITOR_TODO_PROVIDER_PLAN_SCHEMA = "monitor_poll_todo_provider_plan_v1"; const MONITOR_TODO_WRITEBACK_SCHEMA = "monitor_poll_todo_writeback_v0"; const MONITOR_PHASES = ["event", "preflight", "commit"] as const; const MONITOR_SOURCES = ["heartbeat", "controller", "adapter", "visible-goal"] as const; @@ -86,6 +90,7 @@ interface MonitorDecision extends JsonObject { } interface MonitorObservation extends JsonObject { + lease_proof?: TaskLeaseProof; actor_agent_id: string | null; settlement_todo_id: string | null; reason_summary: string | null; @@ -107,7 +112,7 @@ interface MonitorObservation extends JsonObject { } interface MonitorRequest { - schema_version: typeof QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA; + schema_version: typeof QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA | typeof QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA; phase: MonitorPhase; effect_id: string; runtime_root: string | null; @@ -124,7 +129,8 @@ interface MonitorRequest { } interface MonitorProviderPlan extends JsonObject { - schema_version: typeof MONITOR_TODO_PROVIDER_PLAN_SCHEMA; + schema_version: typeof MONITOR_TODO_PROVIDER_PLAN_SCHEMA | typeof LEASED_MONITOR_TODO_PROVIDER_PLAN_SCHEMA; + lease_proof?: TaskLeaseProof; monitor_effect_id: string; goal_id: string; generated_at: string; @@ -148,8 +154,7 @@ interface MonitorProviderPlan extends JsonObject { agent_id: string | null; } -interface PendingMonitorReceipt extends JsonObject { - schema_version: typeof QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA; +interface PendingMonitorReceiptFields extends JsonObject { effect_id: string; request_digest: string; status: "provider_pending"; @@ -159,6 +164,11 @@ interface PendingMonitorReceipt extends JsonObject { provider_plan: JsonObject; } +type PendingMonitorReceipt = PendingMonitorReceiptFields & ( + | {schema_version: typeof QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA} + | {schema_version: typeof MONITOR_PENDING_ADMISSION_SCHEMA; admitted_decision: MonitorDecision} +); + interface DurableMonitorReceipt extends JsonObject { schema_version: typeof QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA; effect_id: string; @@ -304,11 +314,12 @@ function decisionObject(value: unknown): MonitorDecision { function observationObject(value: unknown): MonitorObservation { const observation = requiredObject(value, "observation"); + const proof = decodeTaskLeaseProof(observation.lease_proof); const materialChange = requireBoolean( observation.material_change, "observation.material_change", ); - const result = { + const result: MonitorObservation = { ...observation, actor_agent_id: optionalString( observation.actor_agent_id, @@ -367,12 +378,18 @@ function observationObject(value: unknown): MonitorObservation { observation.next_claimed_by, "observation.next_claimed_by", )?.trim() ?? null, - } satisfies MonitorObservation; + }; if (materialChange && !result.todo_id && !result.target_key) { throw new EffectRuntimeRequestError( "`quota monitor-poll --material-change` requires --todo-id or --target-key", ); } + if (proof) { + if (!result.todo_id && !result.target_key) throw new EffectRuntimeRequestError("lease proof requires a Monitor target"); + result.lease_proof = proof; + } else { + delete result.lease_proof; + } // Validate the route without rewriting the persisted observation fingerprint. // Pending receipts from earlier versions must remain replayable. monitorSuccessorIntent(result); @@ -381,7 +398,8 @@ function observationObject(value: unknown): MonitorObservation { function requestObject(value: unknown): MonitorRequest { const request = requiredObject(value, "quota.monitor_poll.commit params"); - if (request.schema_version !== QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA) { + if (request.schema_version !== QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA && + request.schema_version !== QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA) { throw new EffectRuntimeRequestError("Quota monitor-poll commit request schema mismatch"); } const phase = requireStringLiteral(request.phase, MONITOR_PHASES, "phase"); @@ -409,8 +427,15 @@ function requestObject(value: unknown): MonitorRequest { "turn-scoped monitor-poll requires a registered --agent-id", ); } + const observation = observationObject(request.observation); + if (observation.lease_proof && request.schema_version !== QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA) { + throw new EffectRuntimeRequestError("lease-backed monitor-poll requires request v1"); + } + if (request.schema_version === QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA && !observation.lease_proof) { + throw new EffectRuntimeRequestError("monitor-poll request v1 requires lease_proof"); + } return { - schema_version: QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + schema_version: request.schema_version, phase, effect_id: requiredString(request.effect_id, "effect_id").trim(), runtime_root: runtimeRoot, @@ -429,7 +454,7 @@ function requestObject(value: unknown): MonitorRequest { ), turn_instance_id: turnId, decision, - observation: observationObject(request.observation), + observation, provider_receipt: jsonObject(request.provider_receipt), status_reload_warning: jsonObject(request.status_reload_warning), }; @@ -645,6 +670,7 @@ function compactProviderWriteback(receipt: JsonObject): JsonObject { ]) { compact[field] = receipt[field] ?? null; } + if (receipt.lease_proof != null) compact.lease_proof = receipt.lease_proof; Object.assign(compact, monitorProjectionDelivery(receipt)); return compact; } @@ -669,8 +695,7 @@ function monitorProjectionDelivery(receipt: JsonObject): JsonObject { return {projection_delivery: status, projection_outbox: diagnostic}; } -function buildRecord(request: MonitorRequest): JsonObject { - const allowed = admission(request); +function buildRecord(request: MonitorRequest, allowed: Admission): JsonObject { const material = request.observation.material_change; let kind = "monitor"; let prefix = "monitor"; @@ -789,7 +814,8 @@ function providerPlanFor(request: MonitorRequest): MonitorProviderPlan { throw new EffectRuntimeRequestError("monitor todo writeback requires --result-hash"); } return { - schema_version: MONITOR_TODO_PROVIDER_PLAN_SCHEMA, + schema_version: request.observation.lease_proof ? LEASED_MONITOR_TODO_PROVIDER_PLAN_SCHEMA : MONITOR_TODO_PROVIDER_PLAN_SCHEMA, + ...(request.observation.lease_proof ? {lease_proof: request.observation.lease_proof} : {}), monitor_effect_id: request.effect_id, goal_id: request.goal_id, generated_at: request.generated_at, @@ -816,11 +842,17 @@ function providerPlanFor(request: MonitorRequest): MonitorProviderPlan { function providerPlanObject(value: unknown): MonitorProviderPlan { const plan = requiredObject(value, "receipt.provider_plan"); - if (plan.schema_version !== MONITOR_TODO_PROVIDER_PLAN_SCHEMA) { + if (plan.schema_version !== MONITOR_TODO_PROVIDER_PLAN_SCHEMA && + plan.schema_version !== LEASED_MONITOR_TODO_PROVIDER_PLAN_SCHEMA) { throw new EffectRuntimeRequestError("Monitor Todo provider plan schema mismatch"); } + const proof = decodeTaskLeaseProof(plan.lease_proof); + if ((plan.schema_version === LEASED_MONITOR_TODO_PROVIDER_PLAN_SCHEMA) !== (proof !== null)) { + throw new EffectRuntimeRequestError("Monitor provider plan lease proof/schema mismatch"); + } return { - schema_version: MONITOR_TODO_PROVIDER_PLAN_SCHEMA, + schema_version: plan.schema_version, + ...(proof ? {lease_proof: proof} : {}), monitor_effect_id: requiredString( plan.monitor_effect_id, "provider_plan.monitor_effect_id", @@ -1135,8 +1167,14 @@ function validatedProviderReceipt( "provider_receipt.successor_receipts", ); validateSuccessorReceipts(successors, nextTodos, plan, todoId); + const proof = decodeTaskLeaseProof(receipt.lease_proof); + requireProviderMatch(proof?.idempotency_key ?? null, plan.lease_proof?.idempotency_key ?? null, + "provider_receipt.lease_proof.idempotency_key"); + requireProviderMatch(proof?.expected_version ?? null, plan.lease_proof?.expected_version ?? null, + "provider_receipt.lease_proof.expected_version"); return { schema_version: MONITOR_TODO_WRITEBACK_SCHEMA, + ...(proof ? {lease_proof: proof} : {}), dry_run: dryRun, goal_id: plan.goal_id, todo_id: todoId, @@ -1434,7 +1472,8 @@ function payloadFor( function receiptObject(value: unknown): MonitorReceipt { const receipt = requiredObject(value, "quota monitor-poll transaction receipt"); - if (receipt.schema_version !== QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA) { + if (receipt.schema_version !== QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA && + receipt.schema_version !== MONITOR_PENDING_ADMISSION_SCHEMA) { throw new EffectRuntimeRequestError( "Quota monitor-poll transaction receipt schema mismatch", ); @@ -1463,12 +1502,21 @@ function receiptObject(value: unknown): MonitorReceipt { "receipt.status", ); if (status === "provider_pending") { + if (receipt.schema_version === MONITOR_PENDING_ADMISSION_SCHEMA) { + return {...common, schema_version: MONITOR_PENDING_ADMISSION_SCHEMA, status, + provider_plan: providerPlanObject(receipt.provider_plan), + admitted_decision: decisionObject(receipt.admitted_decision)}; + } return { ...common, status, provider_plan: providerPlanObject(receipt.provider_plan), }; } + if (receipt.schema_version !== QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA) { + throw new EffectRuntimeRequestError("admitted pending receipt cannot represent a completed settlement", + "malformed_transaction_receipt"); + } return { ...common, status, @@ -1757,6 +1805,11 @@ function conflictFields( if (requested.material_change !== (recorded.material_change === true)) { conflicts.push("material_change"); } + const recordedProof = decodeTaskLeaseProof(recorded.lease_proof ?? jsonObject(recorded.todo_writeback)?.lease_proof); + if ((requested.lease_proof?.idempotency_key ?? null) !== (recordedProof?.idempotency_key ?? null) || + (requested.lease_proof?.expected_version ?? null) !== (recordedProof?.expected_version ?? null)) { + conflicts.push("lease_proof"); + } return conflicts; } @@ -1789,7 +1842,7 @@ export async function evaluateQuotaMonitorPollCommit( const request = requestObject(value); const fingerprint = requestDigest(request); if (request.phase === "event") { - const record = buildRecord(request); + const record = buildRecord(request, admission(request)); return result( request, fingerprint, @@ -1810,7 +1863,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. - admission(request); + const allowed = admission(request); if (request.phase === "preflight") { const providerPlan = providerPlanFor(request); return result( @@ -1842,7 +1895,7 @@ export async function evaluateQuotaMonitorPollCommit( validatedProviderReceipt(request.provider_receipt, plan), ); } - const record = buildRecord(effectiveRequest); + const record = buildRecord(effectiveRequest, allowed); const runsDir = request.runtime_root ? join(request.runtime_root, "goals", request.goal_id, "runs") : null; @@ -1905,10 +1958,30 @@ export async function evaluateQuotaMonitorPollCommit( } } - // A completed or provider-pending receipt already proves that the original - // request passed admission. Revalidate only new effects so a changed Todo - // projection cannot block exact-effect replay or crash recovery. - if (!existing) admission(request); + // A pending v1 receipt preserves the decision that admitted this effect. + // Validate that historical basis, never the post-business-commit projection. + // It only authorizes settlement; the provider still fences any new mutation. + let admittedRequest = request; + if (existing?.schema_version === MONITOR_PENDING_ADMISSION_SCHEMA) { + const decision = existing.admitted_decision; + if (decision.goal_id !== request.goal_id || decision.agent_id !== request.decision.agent_id) { + throw new EffectRuntimeRequestError("pending Monitor admission goal/agent does not match request", + "malformed_transaction_receipt"); + } + admittedRequest = {...request, decision}; + } + let allowed: Admission; + try { + allowed = admission(admittedRequest); + } catch (error) { + if (existing?.schema_version === QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA && + error instanceof EffectRuntimeRequestError && error.code === "monitor_poll_admission_rejected") { + throw new EffectRuntimeRequestError( + "legacy pending Monitor receipt has no frozen admission; current admission is unavailable; preserve the receipt for reconciliation", + "legacy_monitor_admission_unavailable"); + } + throw error; + } const indexBytes = await readOptionalBytes(indexPath); const indexContent = indexBytes?.toString("utf8") ?? null; @@ -1977,7 +2050,7 @@ export async function evaluateQuotaMonitorPollCommit( } if (!existing) { const pending = { - schema_version: QUOTA_MONITOR_POLL_COMMIT_RECEIPT_SCHEMA, + schema_version: MONITOR_PENDING_ADMISSION_SCHEMA, effect_id: request.effect_id, request_digest: fingerprint, status: "provider_pending", @@ -1985,6 +2058,7 @@ export async function evaluateQuotaMonitorPollCommit( expected_index_digest: currentDigest, expected_index_bytes: indexBytes?.length ?? 0, provider_plan: plan, + admitted_decision: request.decision, } satisfies PendingMonitorReceipt; await atomicWriteJson(receiptPath, pending); } @@ -2000,7 +2074,7 @@ export async function evaluateQuotaMonitorPollCommit( ); } - let effectiveRequest = request; + let effectiveRequest = admittedRequest; let expectedDigest = currentDigest; let expectedBytes = indexBytes?.length ?? 0; if (providerNeeded) { @@ -2043,7 +2117,7 @@ export async function evaluateQuotaMonitorPollCommit( ); } effectiveRequest = requestWithProvider( - request, + admittedRequest, plan, validatedProviderReceipt(request.provider_receipt, plan), ); @@ -2053,7 +2127,7 @@ export async function evaluateQuotaMonitorPollCommit( return effectConflict(request, fingerprint, currentDigest, existing); } - const record = buildRecord(effectiveRequest); + const record = buildRecord(effectiveRequest, allowed); const { jsonPath, markdownPath } = await nextArtifactPaths( runsDir, effectiveRequest.generated_at, diff --git a/loopx/control_plane/scheduler/monitor_poll_writeback.py b/loopx/control_plane/scheduler/monitor_poll_writeback.py index 2eae209b89..20aaa0c1f9 100644 --- a/loopx/control_plane/scheduler/monitor_poll_writeback.py +++ b/loopx/control_plane/scheduler/monitor_poll_writeback.py @@ -70,6 +70,8 @@ def write_monitor_poll_todo_state( next_user_task_class: str | None = None, next_claimed_by: str | None = None, agent_id: str | None = None, + task_lease_idempotency_key: str | None = None, + task_lease_expected_version: int | None = None, ) -> dict[str, Any] | None: """Apply one monitor poll observation as a complete Todo writeback. @@ -86,13 +88,19 @@ def write_monitor_poll_todo_state( require_legacy_coordination_write_allowed, ) + lease_proof = ({"idempotency_key": task_lease_idempotency_key, + "expected_version": task_lease_expected_version} + if task_lease_idempotency_key is not None or task_lease_expected_version is not None else None) if not todo_id and not target_key: + if lease_proof is not None: + raise ValueError("lease proof requires a Monitor target") return None from .provider_monitor_poll import poll_canonical_monitor_if_promoted canonical = poll_canonical_monitor_if_promoted( registry_path=registry_path, runtime_root=runtime_root, goal_id=goal_id, execute=execute, monitor_effect_id=monitor_effect_id, agent_id=agent_id, + lease_proof=lease_proof, observation={"todo_id": todo_id, "target_key": target_key, "result_hash": result_hash, "material_change": material_change, "generated_at": generated_at, "cadence": cadence, @@ -106,6 +114,8 @@ def write_monitor_poll_todo_state( ) if canonical is not None: return canonical + if lease_proof is not None: + raise ValueError("Monitor lease proof requires promoted canonical authority; no legacy write attempted") if execute: require_legacy_coordination_write_allowed( runtime_root=runtime_root, diff --git a/loopx/control_plane/scheduler/provider_monitor_poll.py b/loopx/control_plane/scheduler/provider_monitor_poll.py index 7eab4ca4f8..ee3187b06a 100644 --- a/loopx/control_plane/scheduler/provider_monitor_poll.py +++ b/loopx/control_plane/scheduler/provider_monitor_poll.py @@ -25,11 +25,14 @@ def poll_canonical_monitor_if_promoted( *, registry_path: Path, runtime_root: Path, goal_id: str, execute: bool, monitor_effect_id: str | None, agent_id: str | None, observation: dict[str, Any], intent: dict[str, Any], + lease_proof: dict[str, Any] | None = None, ) -> dict[str, Any] | None: if not local_authority_is_promoted(runtime_root=runtime_root, goal_id=goal_id): return None result = effect_runtime_result("coordination.local_authority.monitor_poll", { - "schema_version": "loopx_coordination_monitor_poll_request_v0", + "schema_version": ("loopx_coordination_monitor_poll_request_v1" if lease_proof is not None + else "loopx_coordination_monitor_poll_request_v0"), + **({"lease_proof": lease_proof} if lease_proof is not None else {}), "runtime_root": str(runtime_root.expanduser().resolve()), "goal_id": goal_id, "operation_id": monitor_effect_id or f"monitor-poll:{goal_id}:{uuid4().hex}", "actor_agent_id": agent_id, diff --git a/loopx/control_plane/todos/active_state_todos.py b/loopx/control_plane/todos/active_state_todos.py index 72472e0db4..51f470bd4e 100644 --- a/loopx/control_plane/todos/active_state_todos.py +++ b/loopx/control_plane/todos/active_state_todos.py @@ -133,8 +133,9 @@ def active_state_todo_fields( ) if canonical is not None: fields = canonical_todo_summary_fields(canonical["todos"], rollout_events=rollout_events) - # Reading canonical Todos does not qualify legacy monitor writeback. - monitor_writeback_contract_writer(fields, supported=False, source="file_authority") + # Canonical observation/successor transactions now support current + # lease proof. Scheduling exposes due work; mutation admission still + # validates the caller's proof and never falls back to the old writer. elif event_fields.get("user_todos") or event_fields.get("agent_todos"): fields = event_fields markdown_fields = parse_active_state_todos( diff --git a/loopx/quota.py b/loopx/quota.py index d365f184bc..5998aa034e 100644 --- a/loopx/quota.py +++ b/loopx/quota.py @@ -1046,6 +1046,8 @@ def record_quota_monitor_poll( next_user_todo: str | None = None, next_user_task_class: str | None = None, next_claimed_by: str | None = None, + task_lease_idempotency_key: str | None = None, + task_lease_expected_version: int | None = None, turn_instance_id: str | None = None, receipt_bound_todo_id: str | None = None, scheduler_execution_context: Mapping[str, Any] @@ -1213,6 +1215,8 @@ def should_run(current_status: dict[str, Any]) -> dict[str, Any]: next_user_todo=next_user_todo, next_user_task_class=next_user_task_class, next_claimed_by=next_claimed_by, + task_lease_idempotency_key=task_lease_idempotency_key, + task_lease_expected_version=task_lease_expected_version, turn_instance_id=turn_instance_id, status_reloader=status_reloader, ) From 342f6b9ac3e6549fa187814adbcd67c0a36f7f5c Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 18 Sep 2026 01:28:04 +0800 Subject: [PATCH 2/4] test(monitor): qualify lease fences and crash recovery on real providers Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../authority-monitor-poll-rehearsal.py | 201 ++++++++++++++++++ .../test_canonical_status_todos.py | 4 +- .../control_plane/test_leased_monitor_poll.py | 134 ++++++++++++ .../control_plane/test_native_monitor_poll.py | 12 +- .../authority_store_conformance.ts | 7 +- .../coordination_receipt_conformance.ts | 3 + .../monitor_poll_lease_conformance.ts | 121 +++++++++++ .../production_scale_coordination_fixture.ts | 24 +++ .../quota_monitor_poll_commit.test.ts | 86 ++++++++ .../todo_monitor_poll.test.ts | 55 +++++ 10 files changed, 642 insertions(+), 5 deletions(-) create mode 100644 examples/control_plane/authority-monitor-poll-rehearsal.py create mode 100644 tests/control_plane/test_leased_monitor_poll.py create mode 100644 tests/control_plane_ts/monitor_poll_lease_conformance.ts diff --git a/examples/control_plane/authority-monitor-poll-rehearsal.py b/examples/control_plane/authority-monitor-poll-rehearsal.py new file mode 100644 index 0000000000..197482654f --- /dev/null +++ b/examples/control_plane/authority-monitor-poll-rehearsal.py @@ -0,0 +1,201 @@ +#!/usr/bin/env python3 +"""Compare a frozen Monitor baseline and real providers on a read-only Goal snapshot. + +Only disposable copies receive synthetic Monitor/lease records. The report is +bounded and excludes source text, identifiers, paths and connection strings. +""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import subprocess +import sys +from pathlib import Path + +REPOSITORY = Path(__file__).resolve().parents[2] +if str(REPOSITORY) not in sys.path: + sys.path.insert(0, str(REPOSITORY)) + +from loopx.control_plane.coordination.runtime_shadow import build_runtime_shadow_source_snapshot # noqa: E402 +from loopx.history import load_registry # noqa: E402 +from loopx.paths import resolve_runtime_root # noqa: E402 +from loopx.state_refresh import resolve_goal_state # noqa: E402 + +NODE_REHEARSAL = r""" +import assert from 'node:assert/strict'; +import {createHash, randomUUID} from 'node:crypto'; +import {mkdtemp, mkdir, rm, writeFile} from 'node:fs/promises'; +import {tmpdir} from 'node:os'; +import {join} from 'node:path'; +import {pathToFileURL} from 'node:url'; +import {Pool} from 'pg'; +let raw=''; for await (const chunk of process.stdin) raw+=chunk; +const input=JSON.parse(raw); +const moduleAt=(root,path)=>import(pathToFileURL(join(root,'loopx/control_plane',path)).href); +const {pollLocalCoordinationMonitor: current}=await moduleAt(input.repo,'coordination/local_authority_runtime.ts'); +const {pollLocalCoordinationMonitor: baseline}=await moduleAt(input.baseline_repo,'coordination/local_authority_runtime.ts'); +const {FileAuthorityStore}=await moduleAt(input.repo,'coordination/file_authority_store.ts'); +const {SqliteAuthorityStore}=await moduleAt(input.repo,'coordination/sqlite_authority_store.ts'); +const {PostgreSqlAuthorityStore,installPostgreSqlAuthorityStoreSchema}=await moduleAt(input.repo,'coordination/postgresql_authority_store.ts'); +const {PostgreSqlAuthorityService}=await moduleAt(input.repo,'coordination/postgresql_authority_service.ts'); +const {selectLocalSqliteAuthority}=await moduleAt(input.repo,'coordination/local_authority_provider.ts'); +const {coordinationTodoReadModel}=await moduleAt(input.repo,'coordination/coordination_projection.ts'); +const {canonicalAuthoritySha256: digest}=await moduleAt(input.repo,'coordination/authority_store_codec.ts'); +const {engageLegacyCoordinationWriterFence}=await moduleAt(input.repo,'coordination/legacy_writer_fence.ts'); +const goal=input.goal_id, target='todo_monitor_rehearsal'; +const initial=structuredClone(input.projection); +assert(!initial.todos.some(t=>t.todo_id===target)); +initial.todos.push({schema_version:'todo_item_v0',todo_id:target,role:'agent',status:'open',done:false, + text:'Observe isolated public changes',archive_state:'active',source_section:'Agent Todo', + claimed_by:'agent-a',excluded_agents:[],task_class:'continuous_monitor',target_key:'isolated-watch', + cadence:'1h',watch_only:'true',next_due_at:'2026-09-01T00:00:00Z'}); +initial.todos.sort((a,b)=>a.todo_idb.todo_id?1:0); +initial.todo_read_model=coordinationTodoReadModel(initial.todos,initial.todo_read_model.schema_version); +initial.handoff_mode='legacy'; +const proof={idempotency_key:'monitor-rehearsal-execution',expected_version:3}; +const lease={schema_version:'task_lease_v0',goal_id:goal,todo_id:target,owner:'agent-a', + idempotency_key:proof.idempotency_key,version:3,lease_epoch:2,status:'active',write_scopes:[], + acquired_at:'2026-01-01T00:00:00Z',updated_at:'2026-01-01T00:00:00Z',expires_at:'2099-01-01T00:00:00Z'}; +const root=await mkdtemp(join(tmpdir(),'loopx-monitor-rehearsal-')); +const pool=new Pool({connectionString:process.env.LOOPX_TEST_POSTGRES_URL,max:4}); +const database={connect:async()=>{const c=await pool.connect();return {query:async(t,v)=>c.query(t,v),release:e=>c.release(e)};}}; +const tenant=`monitor-rehearsal-${randomUUID()}`; +const heads={}, compatible={}, report={}; +try { + await installPostgreSqlAuthorityStoreSchema(database,`postgresql:${'b'.repeat(32)}`); + for (const arm of ['baseline','file','sqlite','postgresql']) { + const runtime=join(root,arm); await mkdir(runtime,{recursive:true}); + const display=join(runtime,'state.md'); await writeFile(display,'# Disposable display\n'); + let store=new FileAuthorityStore(join(runtime,'authority/file-v0'),goal), dependencies={}; + if(arm==='sqlite') { + assert.equal((await selectLocalSqliteAuthority(runtime,goal,true)).ok,true); + store=new SqliteAuthorityStore(join(runtime,'authority/sqlite-v0'),goal); + } + if(arm==='postgresql') { + store=new PostgreSqlAuthorityStore(database,{tenant_id:tenant,goal_id:goal}); + const identity=await store.storeIdentity(); assert.equal(identity.status,'available'); + const selection={schema_version:'loopx_local_authority_provider_v0',provider:'postgresql',goal_id:goal, + tenant_id:tenant,store_identity:identity.store_identity}; + await mkdir(join(runtime,'authority'),{recursive:true}); + await writeFile(join(runtime,'authority',`provider-${createHash('sha256').update(goal).digest('hex')}.json`),JSON.stringify(selection)); + const service=new PostgreSqlAuthorityService({database, + authenticatePrincipal:()=>({status:'authenticated',principal:{principal_id:'isolated-rehearsal'}}), + authorizeTenant:(principal,selected)=>principal==='isolated-rehearsal'&&selected===tenant?{status:'allowed'}: + {status:'denied',reason_code:'wrong_tenant',reason:'outside disposable tenant'}}); + dependencies={openPostgresqlStore:async selected=>{ + const opened=await service.openStore({credential:null,...selected}); + assert.equal(opened.status,'opened'); return opened.store; + }}; + } + const seed=await store.commitAuthority({expected_provider_revision:null,operation_id:'seed',events:[],receipts:[],next_projection:initial}); + assert.equal(seed.status,'applied'); + assert.equal((await engageLegacyCoordinationWriterFence({schema_version:'loopx_legacy_coordination_writer_fence_engage_request_v0', + runtime_root:runtime,goal_id:goal,state_path:display,fence:{schema_version:'loopx_legacy_coordination_writer_fence_v0', + state:'engaged',goal_id:goal,fence_id:'rehearsal',source_version:'state:1',source_projection_sha256:digest(initial), + expected_shadow_provider_revision:seed.provider_revision}})).status,'applied'); + await rm(display); + const poll=arm==='baseline'?baseline:current; + const request={schema_version:'loopx_coordination_monitor_poll_request_v0',runtime_root:runtime,goal_id:goal, + operation_id:'unchanged',actor_agent_id:'agent-a',registered_agents:['agent-a','agent-b'],dry_run:false, + observation:{todo_id:target,generated_at:'2026-09-01T00:00:00Z',result_hash:'first',material_change:false},intent:{}}; + const firstPoll=await poll(request,dependencies); + assert.equal(firstPoll.status,'applied',JSON.stringify(firstPoll)); + const first=await store.loadAuthority(); assert.equal(first.status,'loaded'); compatible[arm]=first.head; + assert.equal((await poll(request,dependencies)).status,'replayed'); + const leased={...first.head,handoff_mode:'hard_lease',leases:[...first.head.leases,lease].sort((a,b)=>a.todo_idb.todo_id?1:0)}; + assert.equal((await store.commitAuthority({operation_id:'seed-execution',expected_provider_revision:first.provider_revision, + events:[],receipts:[],next_projection:leased})).status,'applied'); + const before=await store.loadAuthority(); + const changed={...request,schema_version:'loopx_coordination_monitor_poll_request_v1',operation_id:'changed',lease_proof:proof, + observation:{...request.observation,generated_at:'2026-09-01T01:00:00Z',result_hash:'second',material_change:true}, + intent:{next_agent_todo:'Validate isolated change',next_action_kind:'validate'}}; + if(arm==='baseline') { + assert.equal((await poll({...changed,schema_version:request.schema_version},dependencies)).status,'failed'); + assert.deepEqual(await store.loadAuthority(),before); + report[arm]={compatible_no_lease_replay:true,leased_poll:'unsupported_no_write'}; continue; + } + for(const invalid of [null,{...proof,expected_version:2},{...proof,idempotency_key:'wrong'}]) { + assert.equal((await poll({...changed,lease_proof:invalid},dependencies)).status,'failed'); + assert.deepEqual(await store.loadAuthority(),before); + } + assert.equal((await poll({...changed,dry_run:true},dependencies)).status,'planned'); + assert.deepEqual(await store.loadAuthority(),before); + const applied=await poll(changed,dependencies); assert.equal(applied.status,'applied'); + const after=await store.loadAuthority(); assert.equal(after.status,'loaded'); + assert.deepEqual(after.head.leases,leased.leases); + assert.equal(after.head.todos.length,initial.todos.length+1); + assert.equal(after.head.todos.find(t=>t.todo_id===target).material_change_generation,1); + assert.deepEqual(after.head.todos.filter(t=>input.projection.todos.some(old=>old.todo_id===t.todo_id)),input.projection.todos); + assert.equal((await poll(changed,dependencies)).status,'replayed'); + assert.deepEqual(await store.loadAuthority(),after); + const expired={...after.head,leases:after.head.leases.map(l=>l.todo_id===target?{...l,expires_at:'2020-01-01T00:00:00Z'}:l)}; + await store.commitAuthority({operation_id:'expire-fixture',expected_provider_revision:after.provider_revision,events:[],receipts:[],next_projection:expired}); + const retired=await store.loadAuthority(); + assert.equal((await poll({...changed,operation_id:'expired-new-observation',intent:{}, + observation:{...changed.observation,generated_at:'2026-09-01T02:00:00Z',result_hash:'third',material_change:false}},dependencies)).status,'failed', + 'expired execution cannot commit another observation'); + assert.equal((await poll(changed,dependencies)).status,'replayed'); + assert.deepEqual(await store.loadAuthority(),retired); + heads[arm]=retired.head; + report[arm]={leased_poll:'applied',replay:'historical',expired_new_poll:'rejected',non_target_unchanged:true}; + } + for(const arm of ['file','sqlite','postgresql']) assert.deepEqual(compatible[arm],compatible.baseline); + assert.deepEqual(heads.file,heads.sqlite); assert.deepEqual(heads.file,heads.postgresql); + process.stdout.write(JSON.stringify({schema_version:'authority_monitor_rehearsal_v0', + source_todos:input.projection.todos.length,source_leases:input.projection.leases.length, + compatible_head_sha256:digest(compatible.baseline),final_head_sha256:digest(heads.file), + provider_heads_equal:true,arms:report})); +} finally { + await pool.query('DELETE FROM loopx_control_plane.authority_heads WHERE tenant_id=$1',[tenant]); + await pool.end(); await rm(root,{recursive:true,force:true}); +} +""" + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--registry", type=Path, required=True) + parser.add_argument("--goal-id", required=True) + parser.add_argument("--baseline-repo", type=Path, required=True) + parser.add_argument("--execute-isolated-postgresql", action="store_true") + parser.add_argument("--private-diagnostics", type=Path) + args = parser.parse_args() + if not args.execute_isolated_postgresql or not os.environ.get("LOOPX_TEST_POSTGRES_URL"): + raise SystemExit("an explicitly isolated PostgreSQL server is required") + baseline = args.baseline_repo.resolve() + def git(*command): + return subprocess.check_output(["git", "-C", str(baseline), *command], text=True).strip() + if git("status", "--porcelain", "--untracked-files=no"): + raise SystemExit("baseline must have no tracked changes") + revision = git("rev-parse", "HEAD") + registry_path = args.registry.resolve() + registry_bytes = registry_path.read_bytes() + registry = load_registry(registry_path) + goal = next(g for g in registry["goals"] if g["id"] == args.goal_id) + runtime = resolve_runtime_root(registry, None, registry_path=registry_path) + _, _, state = resolve_goal_state(registry=registry, goal_id=args.goal_id, project_override=None, state_file_override=None) + projection, snapshot = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=runtime, state_path=state, registry_path=registry_path) + source_digest = hashlib.sha256(json.dumps(snapshot, sort_keys=True).encode()).hexdigest() + child = subprocess.run(["node", "--no-warnings", "--experimental-sqlite", "--experimental-strip-types", "--input-type=module", "-e", NODE_REHEARSAL], + input=json.dumps({"repo": str(REPOSITORY), "baseline_repo": str(baseline), "goal_id": args.goal_id, "projection": projection}), + cwd=REPOSITORY, capture_output=True, text=True, timeout=180, check=False) + if child.returncode: + if args.private_diagnostics: + descriptor = os.open(args.private_diagnostics, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + os.fchmod(descriptor, 0o600) + with os.fdopen(descriptor, "w") as stream: + stream.write(child.stderr) + raise SystemExit(f"isolated Monitor rehearsal failed (exit {child.returncode}); no qualification; diagnostics stay private") + after_projection, after_snapshot = build_runtime_shadow_source_snapshot(goal=goal, runtime_root=runtime, state_path=state, registry_path=registry_path) + if projection != after_projection or snapshot != after_snapshot or registry_path.read_bytes() != registry_bytes: + raise SystemExit("live source changed during rehearsal; rerun from a stable snapshot") + result = json.loads(child.stdout) + result.update(baseline_revision=revision, source_snapshot_sha256=source_digest, source_unchanged=True) + print(json.dumps(result, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/control_plane/test_canonical_status_todos.py b/tests/control_plane/test_canonical_status_todos.py index 3a62425235..97d2f5414f 100644 --- a/tests/control_plane/test_canonical_status_todos.py +++ b/tests/control_plane/test_canonical_status_todos.py @@ -71,7 +71,7 @@ def promoted_goal(tmp_path: Path, request): @pytest.mark.parametrize("display", ["stale", "missing", "invalid_utf8"]) -def test_status_uses_provider_even_when_markdown_is_unusable(promoted_goal, display, monkeypatch): +def test_status_uses_provider_and_allows_native_monitor_writeback_without_display(promoted_goal, display, monkeypatch): goal, runtime, state = promoted_goal if display == "invalid_utf8": state.write_bytes(b"\xff") @@ -89,7 +89,7 @@ def test_status_uses_provider_even_when_markdown_is_unusable(promoted_goal, disp assert items[0]["claimed_by"] == "agent-a" assert fields["standing_decision_authority"]["active_count"] == 0 assert fields["standing_decision_authority"]["entries"][0]["outcome"] == "reject" - assert fields["agent_todos"]["monitor_writeback"]["supported"] is False + assert "monitor_writeback" not in fields["agent_todos"] # Native observation owns writeback. if display == "missing": assert not state.exists() # A read must not recreate the display. else: diff --git a/tests/control_plane/test_leased_monitor_poll.py b/tests/control_plane/test_leased_monitor_poll.py new file mode 100644 index 0000000000..057727caed --- /dev/null +++ b/tests/control_plane/test_leased_monitor_poll.py @@ -0,0 +1,134 @@ +"""Public CLI proof transport and recovery across business/quota authority.""" +from __future__ import annotations + +import json +import subprocess +import sys + +import pytest + +from canonical_authority_fixture import isolate_sqlite_runtime +from test_native_monitor_poll import _canonical +from test_monitor_followthrough_contract import GOAL_ID, AGENT_ID, _write_fixture, _add_monitor +from loopx.control_plane.coordination.local_authority import read_canonical_todos_if_promoted +from loopx.control_plane.scheduler.monitor_poll_writeback import write_monitor_poll_todo_state +from loopx.control_plane.testing.canary_harness import run_json_cli, run_json_cli_result + + +LEASE = {"status": "active", "idempotency_key": "monitor-execution", "version": 3, + "lease_epoch": 2, "acquired_at": "2026-01-01T00:00:00Z", "expires_at": "2099-01-01T00:00:00Z"} +PROOF = {"idempotency_key": "monitor-execution", "expected_version": 3} + + +def arguments(monitor): + return ["quota", "monitor-poll", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--runtime-profile", "generic_cli", "--todo-id", monitor["todo_id"], + "--result-hash", "observed-a", "--material-change", "--next-agent-todo", "Validate changed evidence", + "--next-action-kind", "validate", "--task-lease-idempotency-key", PROOF["idempotency_key"], + "--task-lease-expected-version", str(PROOF["expected_version"]), "--execute"] + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("native", [False, True]) +def test_leased_monitor_public_cli_settles_once_without_display(tmp_path, monkeypatch, provider, native): + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, state, monitor = _canonical(tmp_path, provider=provider, native=native, lease=LEASE) + before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) + state.unlink() + result = run_json_cli(*arguments(monitor), registry_path=registry, runtime_root=runtime) + assert result["ok"] is True + assert result["todo_writeback"]["lease_proof"] == PROOF + assert result["todo_writeback"]["material_change_generation"] == 1 + assert len(result["todo_writeback"]["next_todos"]) == 1 + after = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) + assert after["leases"] == before["leases"] + assert len(after["todos"]) == len(before["todos"]) + 1 + assert state.exists() + records = [json.loads(line) for line in (runtime / "goals" / GOAL_ID / "runs" / "index.jsonl").read_text().splitlines()] + assert sum(row.get("classification") == "quota_monitor_poll" for row in records) == 1 + assert all(row.get("classification") != "quota_slot_spend" for row in records) + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_old_observation_time_cannot_revive_expired_proof(tmp_path, monkeypatch, provider): + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, _state, monitor = _canonical(tmp_path, provider=provider, + lease={**LEASE, "expires_at": "2026-01-01T01:00:00Z"}) + before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) + with pytest.raises(RuntimeError, match="current lease proof"): + write_monitor_poll_todo_state(registry_path=registry, runtime_root=runtime, goal_id=GOAL_ID, + execute=True, todo_id=monitor["todo_id"], agent_id=AGENT_ID, + generated_at="2026-01-01T00:30:00Z", result_hash="old-but-once-valid", material_change=False, + task_lease_idempotency_key=PROOF["idempotency_key"], task_lease_expected_version=3) + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) == before + + +def test_explicit_proof_never_falls_back_to_legacy_writer(tmp_path): + registry, runtime, state = _write_fixture(tmp_path) + monitor = _add_monitor(registry, text="Observe public progress", target_key="watch", next_due_at="2000-01-01T00:00:00Z") + before = state.read_bytes() + with pytest.raises(ValueError, match="requires promoted canonical authority"): + write_monitor_poll_todo_state(registry_path=registry, runtime_root=runtime, goal_id=GOAL_ID, + execute=True, todo_id=monitor["todo_id"], agent_id=AGENT_ID, + generated_at="2026-09-01T00:00:00Z", result_hash="a", material_change=False, + task_lease_idempotency_key=PROOF["idempotency_key"], task_lease_expected_version=3) + assert state.read_bytes() == before + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +def test_process_death_after_business_commit_recovers_after_lease_release(tmp_path, monkeypatch, provider): + isolate_sqlite_runtime(tmp_path, monkeypatch) + registry, runtime, _state, monitor = _canonical(tmp_path, provider=provider, lease=LEASE) + turn_id = "leased-monitor-crash" + guard = run_json_cli("quota", "should-run", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--runtime-profile", "generic_cli", "--turn-instance-id", turn_id, + "--available-capability", "network", "--available-capability", "external_evidence_poll", + registry_path=registry, runtime_root=runtime) + assert guard["selected_todo"]["todo_id"] == monitor["todo_id"] + args = [*arguments(monitor), "--turn-instance-id", turn_id, + "--available-capability", "network", "--available-capability", "external_evidence_poll"] + # Kill the actual CLI process after the authority transaction returned its + # receipt, before the separate quota transaction can settle it. + script = """ +import os, sys +from loopx.cli import main +import loopx.control_plane.quota.monitor_poll as monitor +original = monitor._native_result +def crash(request): + if request.get('phase') == 'commit' and request.get('provider_receipt'): + os._exit(73) + return original(request) +monitor._native_result = crash +main(sys.argv[1:]) +""" + process = subprocess.run([sys.executable, "-c", script, "--registry", str(registry), + "--runtime-root", str(runtime), "--format", "json", *args], capture_output=True, text=True, timeout=60) + assert process.returncode == 73, process.stderr + process.stdout + committed = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) + assert len(committed["todos"]) == 2 + run_json_cli("task-lease", "release", "--goal-id", GOAL_ID, "--todo-id", monitor["todo_id"], + "--owner", AGENT_ID, "--idempotency-key", PROOF["idempotency_key"], "--expected-version", "3", + registry_path=registry, runtime_root=runtime) + retired = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) + assert retired["leases"][0]["status"] == "released" + # An old pending receipt cannot invent the missing admission snapshot. + pending_path, = (runtime / "goals" / GOAL_ID / "runs" / ".transactions" / "quota-monitor-poll").glob("*.json") + pending_bytes = pending_path.read_bytes() + legacy = json.loads(pending_bytes) + legacy["schema_version"] = "quota_monitor_poll_commit_receipt_v0" + del legacy["admitted_decision"] + pending_path.write_text(json.dumps(legacy)) + code, failure = run_json_cli_result(*args, registry_path=registry, runtime_root=runtime) + assert code != 0 + assert failure["error_code"] == "legacy_monitor_admission_unavailable" + assert json.loads(pending_path.read_text()) == legacy + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) == retired + pending_path.write_bytes(pending_bytes) + recovered = run_json_cli(*args, registry_path=registry, runtime_root=runtime) + assert recovered["ok"] is True + assert recovered["todo_writeback"]["lease_proof"] == PROOF + again = run_json_cli(*args, registry_path=registry, runtime_root=runtime) + assert again["replayed"] is True + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID, include_leases=True) == retired + records = [json.loads(line) for line in (runtime / "goals" / GOAL_ID / "runs" / "index.jsonl").read_text().splitlines()] + assert sum(row.get("classification") == "quota_monitor_poll" for row in records) == 1 diff --git a/tests/control_plane/test_native_monitor_poll.py b/tests/control_plane/test_native_monitor_poll.py index c96f0473de..8f8fd64e2d 100644 --- a/tests/control_plane/test_native_monitor_poll.py +++ b/tests/control_plane/test_native_monitor_poll.py @@ -14,12 +14,20 @@ from loopx.todos import list_goal_todos -def _canonical(tmp_path, native=False, provider="file"): +def _canonical(tmp_path, native=False, provider="file", lease=None): registry, runtime, state = _write_fixture(tmp_path) monitor = _add_monitor(registry, text="Observe a public target", target_key="public-watch", next_due_at="2000-01-01T00:00:00Z") items = list_goal_todos(registry_path=registry, goal_id=GOAL_ID, role="agent")["todos"] - projection = build_todo_runtime_shadow_projection(goal_id=GOAL_ID, todos=items, leases=[], handoff_mode="soft_claim") + leases = [] + if lease is not None: + for item in items: + if item["todo_id"] == monitor["todo_id"]: + item["claimed_by"] = AGENT_ID + leases = [{"schema_version": "task_lease_v0", "goal_id": GOAL_ID, + "todo_id": monitor["todo_id"], "owner": AGENT_ID, "write_scopes": [], **lease}] + projection = build_todo_runtime_shadow_projection(goal_id=GOAL_ID, todos=items, leases=leases, + handoff_mode="hard_lease" if lease is not None else "soft_claim") if native: from loopx.control_plane.coordination.coordination_state_contract import TODO_DOMAIN_RECORD_FIELDS from loopx.control_plane.coordination.local_authority_shadow_projection import canonical_bytes diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index ef1cb020b6..85f19c63ae 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -1,4 +1,5 @@ import {registerLeaseAcquisitionConformance} from "./lease_acquisition_conformance.ts"; +import {registerLeasedMonitorConformance} from "./monitor_poll_lease_conformance.ts"; import {registerLeaseLifecycleConformance} from "./lease_lifecycle_conformance.ts"; import {registerMonitorConfigurationConformance} from "./monitor_configuration_conformance.ts"; import {registerAuthorityScanConformance} from "./authority_scan_conformance.ts"; @@ -217,6 +218,7 @@ export function registerAuthorityStoreConformance( registerOwnershipObservationConformance(providerName, factory); registerNativePlanningUpdateConformance(providerName, factory); registerMonitorConfigurationConformance(providerName, factory); + registerLeasedMonitorConformance(providerName, factory); registerCoordinationReceiptConformance(providerName, factory); registerHandoffModeConformance(providerName, factory); for (const native of [false, true]) test(`${providerName} conformance: standing revocation survives canonical ordering and archive (${native ? "native" : "legacy"})`, async (t) => { @@ -256,11 +258,14 @@ export function registerAuthorityStoreConformance( assert.equal((await executeCoordinationTodoArchiveCompleted(store, request)).status, "replayed"); }); - for (const native of [false, true]) test(`${providerName} conformance: atomic Monitor observation and successor (${native ? "native" : "legacy"})`, async (t) => { + for (const native of [false, true]) test(`${providerName} conformance: lease-free legacy-mode Monitor observation and successor (${native ? "native" : "legacy"})`, async (t) => { const {store, contender} = await factory(t); const goal = "goal-monitor"; const fixture = productionScaleCoordinationFixture(goal, native ? "native" : "legacy"); const projection = structuredClone(fixture.projection); + // A lease-free observation is legal in legacy mode. Hard mode now requires + // current execution proof; its positive/negative cases have their own fixture. + projection.handoff_mode = "legacy"; const records = projection.todos as Record[]; const monitor = records.find(todo => todo.task_class === "continuous_monitor" && todo.status !== "done" && !(projection.leases as Record[]).some(lease => lease.todo_id === todo.todo_id)); diff --git a/tests/control_plane_ts/coordination_receipt_conformance.ts b/tests/control_plane_ts/coordination_receipt_conformance.ts index bd87d98d7f..175e02ad0d 100644 --- a/tests/control_plane_ts/coordination_receipt_conformance.ts +++ b/tests/control_plane_ts/coordination_receipt_conformance.ts @@ -24,6 +24,9 @@ async function commandFixture(store: AuthorityStore, command: Command) { const goal_id = "goal-a"; const fixture = productionScaleCoordinationFixture(goal_id); const projection = fixture.projection; + // This recovery arm deliberately has no Monitor execution lease. The leased + // arm separately proves current proof and historical replay in hard mode. + if (command === "monitor") projection.handoff_mode = "legacy"; const todos = projection.todos as JsonObject[]; const leased = new Set((projection.leases as JsonObject[]).map(row => row.todo_id)); const claimTodo = todos.find(row => row.role === "agent" && row.status === "open" && diff --git a/tests/control_plane_ts/monitor_poll_lease_conformance.ts b/tests/control_plane_ts/monitor_poll_lease_conformance.ts new file mode 100644 index 0000000000..dddee8381d --- /dev/null +++ b/tests/control_plane_ts/monitor_poll_lease_conformance.ts @@ -0,0 +1,121 @@ +/** Real providers run the same current-proof, atomicity and recovery contract. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; +import {coordinationTodoReadModel} from "../../loopx/control_plane/coordination/coordination_projection.ts"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; +import {productionScaleLeasedMonitorFixture} from "./production_scale_coordination_fixture.ts"; + +export function registerLeasedMonitorConformance(provider: string, factory: AuthorityStoreConformanceFactory) { + for (const schema of ["legacy", "native"] as const) test(`${provider}: leased Monitor ${schema} keeps execution and commits one generation with successors`, async t => { + const {store, contender} = await factory(t); + const goal = "leased-monitor"; + const fixture = productionScaleLeasedMonitorFixture(goal, schema); + assert.equal((await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + events: [], receipts: [], next_projection: fixture.projection})).status, "applied"); + const request = {goal_id: goal, operation_id: "poll", actor_agent_id: fixture.actor, + registered_agents: fixture.registered_agents, now: fixture.now, lease_proof: fixture.proof, + dry_run: false, observation: {todo_id: fixture.target, generated_at: fixture.now.toISOString(), + result_hash: "changed-evidence", material_change: true}, + intent: {next_agent_todo: "Validate changed evidence", next_action_kind: "validate", + next_user_todo: "Review changed evidence", next_user_task_class: "user_action"}}; + const before = await store.loadAuthority(); + for (const extra of [{lease_proof: null}, {lease_proof: {...fixture.proof, expected_version: 999}}, + {lease_proof: {...fixture.proof, idempotency_key: "wrong"}}, {actor_agent_id: "agent-b"}, + {registered_agents: ["agent-b"]}, {now: new Date("2099-01-01T00:00:00Z")}]) { + const result = await executeCoordinationMonitorPoll(store, {...request, ...extra}); + assert.equal(result.status, "failed", JSON.stringify(result)); + assert.deepEqual(await store.loadAuthority(), before); + } + assert.equal((await executeCoordinationMonitorPoll(store, {...request, dry_run: true})).status, "planned"); + assert.deepEqual(await store.loadAuthority(), before); + const first = await executeCoordinationMonitorPoll(store, request); + assert.equal(first.status, "applied", JSON.stringify(first)); + const after = await store.loadAuthority(); + assert.equal(after.status, "loaded"); + if (after.status !== "loaded") return; + const rows = after.head.todos as JsonObject[]; + assert.deepEqual(after.head.leases, fixture.projection.leases); + const original = fixture.projection.todos as JsonObject[]; + assert.deepEqual(rows.filter(todo => original.some(old => old.todo_id === todo.todo_id) && todo.todo_id !== fixture.target), + original.filter(todo => todo.todo_id !== fixture.target)); + assert.equal(rows.length, original.length + 2); + assert.equal(rows.find(todo => todo.todo_id === fixture.target)?.material_change_generation, 5); + const unchanged = await executeCoordinationMonitorPoll(store, {...request, operation_id: "unchanged", + observation: {...request.observation, generated_at: "2026-09-01T02:00:00Z", material_change: false}, intent: {}}); + assert.equal(unchanged.status, "applied"); + assert.equal((unchanged.writeback as JsonObject).material_change_generation, 5); + assert.deepEqual((unchanged.writeback as JsonObject).next_todos, []); + const current = await store.loadAuthority(); + assert.equal(current.status, "loaded"); + if (current.status !== "loaded") return; + // Execution retirement and Todo archival cannot make historical settlement + // re-run a poll or require a new execution's proof. + const retired = (current.head.todos as JsonObject[]).map(todo => todo.todo_id === fixture.target + ? {...todo, archive_state: "archive", status: "done", done: true} : todo); + await store.commitAuthority({operation_id: "retire", expected_provider_revision: current.provider_revision, + events: [], receipts: [], next_projection: {...current.head, todos: retired, + todo_read_model: coordinationTodoReadModel(retired, (current.head.todo_read_model as JsonObject).schema_version), + leases: (current.head.leases as JsonObject[]).map(lease => lease.todo_id === fixture.target + ? {...lease, status: "released"} : lease)}}); + const retiredHead = await store.loadAuthority(); + const replay = await executeCoordinationMonitorPoll(contender, {...request, now: new Date("2099-01-01T00:00:00Z")}); + assert.equal(replay.status, "replayed"); + assert.deepEqual((replay.writeback as JsonObject).next_todos, (first.writeback as JsonObject).next_todos); + assert.deepEqual(await store.loadAuthority(), retiredHead); + assert.equal((await executeCoordinationMonitorPoll(store, {...request, + lease_proof: {...fixture.proof, expected_version: fixture.proof.expected_version + 1}})).reason_code, + "coordination_operation_identity_mismatch"); + }); + + test(`${provider}: leased Monitor recovers a lost commit and rejects a racing renewal`, async t => { + const {store, contender} = await factory(t); + const fixture = productionScaleLeasedMonitorFixture("monitor-recovery"); + await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + events: [], receipts: [], next_projection: fixture.projection}); + const request = {goal_id: "monitor-recovery", operation_id: "lost-response", actor_agent_id: fixture.actor, + registered_agents: fixture.registered_agents, now: fixture.now, lease_proof: fixture.proof, dry_run: false, + observation: {todo_id: fixture.target, generated_at: fixture.now.toISOString(), result_hash: "first", material_change: true}, + intent: {next_agent_todo: "Validate once", next_action_kind: "validate"}}; + let responseLost = false; + const lost = new Proxy(store, {get(target, property) { + if (property === "commitAuthority") return async (commit: AuthorityStoreCommit) => { + await target.commitAuthority(commit); responseLost = true; throw new Error("response lost"); + }; + if (property === "readReceipt") return async (operation: string) => { + if (responseLost) throw new Error("readback unavailable"); + return target.readReceipt(operation); + }; + const value = Reflect.get(target, property); return typeof value === "function" ? value.bind(target) : value; + }}) as AuthorityStore; + assert.equal((await executeCoordinationMonitorPoll(lost, request)).status, "ambiguous"); + const committed = await store.loadAuthority(); + assert.equal((await executeCoordinationMonitorPoll(contender, request)).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), committed); + let race = true; + const renewed = new Proxy(store, {get(target, property) { + if (property === "commitAuthority") return async (commit: AuthorityStoreCommit) => { + if (race) { + race = false; + const loaded = await contender.loadAuthority(); + assert.equal(loaded.status, "loaded"); + if (loaded.status !== "loaded") throw new Error("missing race source"); + await contender.commitAuthority({operation_id: "renew", expected_provider_revision: loaded.provider_revision, + events: [], receipts: [], next_projection: {...loaded.head, + leases: (loaded.head.leases as JsonObject[]).map(lease => lease.todo_id === fixture.target + ? {...lease, version: fixture.proof.expected_version + 1} : lease)}}); + } + return target.commitAuthority(commit); + }; + const value = Reflect.get(target, property); return typeof value === "function" ? value.bind(target) : value; + }}) as AuthorityStore; + const next = {...request, operation_id: "racing-poll", + observation: {...request.observation, generated_at: "2026-09-01T02:00:00Z", result_hash: "second"}, intent: {}}; + const raced = await executeCoordinationMonitorPoll(renewed, next); + assert.equal(raced.status, "conflict", JSON.stringify(raced)); + assert.equal((await executeCoordinationMonitorPoll(store, next)).status, "failed"); + assert.equal((await store.readReceipt(next.operation_id)).status, "missing"); + }); +} diff --git a/tests/control_plane_ts/production_scale_coordination_fixture.ts b/tests/control_plane_ts/production_scale_coordination_fixture.ts index 27eb179517..c02b9da3d2 100644 --- a/tests/control_plane_ts/production_scale_coordination_fixture.ts +++ b/tests/control_plane_ts/production_scale_coordination_fixture.ts @@ -507,3 +507,27 @@ export function productionScaleLeaseAcquisitionFixture(goalId: string, {source_authority: "synthetic_production_scale_fixture", handoff_mode: "hard_lease"}); return {...fixture, projection, acquisition: scenario}; } + +/** A leased Monitor in the complete mixed work graph, with prior observations + * and a waiting dependent. Only these target rows differ from the base fixture. */ +export function productionScaleLeasedMonitorFixture(goalId: string, + schema: AuthorityProjectionSchema = "native") { + const fixture = productionScaleCoordinationFixture(goalId, schema); + const target = fixture.completion_todo_id; + const rows = fixture.projection.todos as Record[]; + const dependent = [...rows].reverse().find(todo => todo.role === "agent" && todo.status === "open" && todo.todo_id !== target)!; + const todos = rows.map(todo => todo.todo_id === target ? {...todo, + task_class: "continuous_monitor", target_key: "leased-public-watch", cadence: "1h", + material_change_generation: 4, consecutive_no_change: "2", result_hash: "previous-evidence", + last_checked_at: "2026-09-01T00:00:00Z", next_due_at: "2026-09-01T01:00:00Z"} : + todo.todo_id === dependent.todo_id ? {...todo, resume_when: `monitor_changed:${target}`, + resume_monitor_generation: 4} : todo); + const projection = authorityProjectionFixture(goalId, todos, + fixture.projection.leases as Record[], schema, + {source_authority: "synthetic_production_scale_fixture", handoff_mode: "hard_lease"}); + return {projection, target, dependent: String(dependent.todo_id), + actor: "agent-a", registered_agents: fixture.registered_agents, + now: new Date("2026-09-01T01:00:00Z"), + proof: {idempotency_key: fixture.completion_lease_idempotency_key, + expected_version: fixture.completion_lease_expected_version}}; +} diff --git a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts index 8a87da1ed0..9949c0eda2 100644 --- a/tests/control_plane_ts/quota_monitor_poll_commit.test.ts +++ b/tests/control_plane_ts/quota_monitor_poll_commit.test.ts @@ -9,6 +9,7 @@ import test from "node:test"; import { evaluateQuotaMonitorPollCommit, QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, } from "../../loopx/control_plane/quota/monitor_poll_commit.ts"; import { EffectRuntimeRequestError } from "../../loopx/control_plane/effect_runtime_errors.ts"; @@ -871,3 +872,88 @@ test("durable replay rejects receipt path escape and lost pre-existing index his /artifact paths do not match the transaction scope/, ); }); + +test("leased Monitor pending settlement binds the original proof and rejects missing or substituted receipts", async t => { + const runtime = await tempRuntime(t); + const proof = {idempotency_key: "monitor-execution", expected_version: 3}; + const params = request({schema_version: QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + phase: "preflight", runtime_root: runtime, execute: true, effect_id: "leased-poll", + observation: observation({todo_id: "todo_public_monitor", result_hash: "observed-a", lease_proof: proof})}); + await assert.rejects(evaluateQuotaMonitorPollCommit({...params, + schema_version: QUOTA_MONITOR_POLL_COMMIT_REQUEST_SCHEMA}), /requires request v1/); + const preflight = await evaluateQuotaMonitorPollCommit(params); + assert.equal(preflight.status, "provider_required"); + assert.deepEqual(preflight.provider_plan?.lease_proof, proof); + assert.equal(preflight.provider_plan?.schema_version, "monitor_poll_todo_provider_plan_v1"); + const receipt = {schema_version: "monitor_poll_todo_writeback_v0", monitor_effect_id: "leased-poll", + goal_id: goalId, todo_id: "todo_public_monitor", dry_run: false, target_key: null, + result_hash: "observed-a", material_change: false, material_change_generation: 0, + consecutive_no_change: 1, last_checked_at: params.generated_at, next_due_at: null, cadence: null, + todo_update: {ok: true}, next_todos: [], successor_receipts: [], lease_proof: proof}; + for (const invalid of [undefined, {...proof, expected_version: 4}, {...proof, idempotency_key: "another"}]) { + await assert.rejects(evaluateQuotaMonitorPollCommit({...params, phase: "commit", + provider_receipt: {...receipt, lease_proof: invalid}}), /lease_proof/); + } + const changedDecision = decision({effective_action: "normal_delivery", heartbeat_recommendation: {}, + recommended_action: "Work on the new successor", reason: "Monitor is no longer due"}); + const retry = {...params, decision: changedDecision, generated_at: "2099-01-01T00:00:00Z"}; + const replay = await evaluateQuotaMonitorPollCommit(retry); + assert.deepEqual(replay.provider_plan, preflight.provider_plan); + const written = await evaluateQuotaMonitorPollCommit({...retry, phase: "commit", provider_receipt: receipt}); + assert.equal(written.status, "written"); + assert.equal(written.record?.recommended_action, "Watch the public release queue."); + assert.equal(((written.record?.monitor_event as Record).before as Record).effective_action, + "monitor_quiet_skip"); + assert.deepEqual((written.payload.todo_writeback as Record).lease_proof, proof); + assert.equal((await evaluateQuotaMonitorPollCommit({...params, phase: "commit", provider_receipt: receipt})).status, "replayed"); + const changed = await evaluateQuotaMonitorPollCommit({...params, + observation: observation({todo_id: "todo_public_monitor", result_hash: "observed-a", lease_proof: {...proof, expected_version: 4}})}); + assert.equal(changed.status, "conflict"); +}); + +test("pending admission is scoped, mandatory in v1, and preserves bounded v0 recovery", async t => { + const runtime = await tempRuntime(t); + const effect = "pending-admission-compatibility"; + const params = request({phase: "preflight", runtime_root: runtime, execute: true, effect_id: effect, + observation: observation({todo_id: "todo_public_monitor", result_hash: "observed-a"})}); + const preflight = await evaluateQuotaMonitorPollCommit(params); + const path = join(runtime, "goals", goalId, "runs", ".transactions", "quota-monitor-poll", + `${createHash("sha256").update(effect).digest("hex").slice(0, 24)}.json`); + const pending = JSON.parse(await readFile(path, "utf8")); + for (const admitted of [null, {...pending.admitted_decision, goal_id: "another-goal"}, + {...pending.admitted_decision, agent_id: "another-agent"}, + {...pending.admitted_decision, effective_action: "normal_delivery", heartbeat_recommendation: {}}]) { + const bytes = JSON.stringify({...pending, admitted_decision: admitted}); + await writeFile(path, bytes); + await assert.rejects(evaluateQuotaMonitorPollCommit(params), EffectRuntimeRequestError); + assert.equal(await readFile(path, "utf8"), bytes); + } + const legacy = {...pending, schema_version: "quota_monitor_poll_commit_receipt_v0"}; + delete legacy.admitted_decision; + await writeFile(path, JSON.stringify(legacy)); + assert.deepEqual((await evaluateQuotaMonitorPollCommit(params)).provider_plan, preflight.provider_plan); + await assert.rejects(evaluateQuotaMonitorPollCommit({...params, + decision: decision({effective_action: "normal_delivery", heartbeat_recommendation: {}})}), + (error: unknown) => error instanceof EffectRuntimeRequestError && error.code === "legacy_monitor_admission_unavailable"); + assert.deepEqual(JSON.parse(await readFile(path, "utf8")), legacy); + const receipt = {schema_version: "monitor_poll_todo_writeback_v0", monitor_effect_id: effect, + goal_id: goalId, todo_id: "todo_public_monitor", dry_run: false, target_key: null, + result_hash: "observed-a", material_change: false, material_change_generation: 0, + consecutive_no_change: 1, last_checked_at: params.generated_at, next_due_at: null, cadence: null, + todo_update: {}, next_todos: [], successor_receipts: []}; + assert.equal((await evaluateQuotaMonitorPollCommit({...params, phase: "commit", provider_receipt: receipt})).status, "written"); +}); + +test("lease proof and schema must agree before any provider intent is journaled", async t => { + const runtime = await tempRuntime(t); + const params = request({schema_version: QUOTA_LEASED_MONITOR_POLL_COMMIT_REQUEST_SCHEMA, + phase: "preflight", runtime_root: runtime, execute: true, effect_id: "invalid-proof"}); + for (const proof of [undefined, {}, {idempotency_key: " x", expected_version: 1}, + {idempotency_key: "x", expected_version: 0}, {idempotency_key: "x", expected_version: 1.5}, + {idempotency_key: "x", expected_version: Number.MAX_SAFE_INTEGER + 1}, + {idempotency_key: "x", expected_version: 1, owner: "injected"}]) { + await assert.rejects(evaluateQuotaMonitorPollCommit({...params, + observation: observation({todo_id: "todo_public_monitor", result_hash: "a", lease_proof: proof})})); + } + await assert.rejects(readFile(join(runtime, "goals", goalId, "runs", "index.jsonl")), {code: "ENOENT"}); +}); diff --git a/tests/control_plane_ts/todo_monitor_poll.test.ts b/tests/control_plane_ts/todo_monitor_poll.test.ts index 93a2e18bcf..23f66991b4 100644 --- a/tests/control_plane_ts/todo_monitor_poll.test.ts +++ b/tests/control_plane_ts/todo_monitor_poll.test.ts @@ -128,6 +128,25 @@ test("retained lease and stale observation cannot mutate either half", async () assert.deepEqual(await store.loadAuthority(), leased); }); +test("hard mode cannot bypass a missing lease; soft mode cannot accept supplied execution proof", async () => { + for (const mode of ["hard_lease", "soft_claim"] as const) { + const {store, request} = await seeded({claimed_by: "agent-a"}); + const loaded = await store.loadAuthority(); + assert.equal(loaded.status, "loaded"); + if (loaded.status !== "loaded") return; + await store.commitAuthority({operation_id: "mode", expected_provider_revision: loaded.provider_revision, + events: [], receipts: [], next_projection: {...loaded.head, handoff_mode: mode}}); + const before = await store.loadAuthority(); + for (const proof of [null, {idempotency_key: "missing-execution", expected_version: 1}]) { + if (mode === "soft_claim" && proof === null) continue; + const result = await executeCoordinationMonitorPoll(store, {...request, lease_proof: proof}); + assert.equal(result.status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal((await store.readReceipt(request.operation_id)).status, "missing"); + } + } +}); + test("target-key selection ignores completed history but never guesses between live monitors", () => { const active = {todo_id: "todo_current", role: "agent", task_class: "continuous_monitor", target_key: "watch", status: "open", archive_state: "active"}; @@ -136,3 +155,39 @@ test("target-key selection ignores completed history but never guesses between l assert.throws(() => selectMonitorTodo([history, active], history.todo_id, "watch"), /unfinished/); assert.throws(() => selectMonitorTodo([active, {...active, todo_id: "todo_other"}], null, "watch"), /multiple/); }); + +test("current leased Monitor atomically observes and creates work without renewing or releasing its lease", async () => { + const {store, request} = await seeded({claimed_by: "agent-a"}); + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); + if (head.status !== "loaded") return; + const lease = {schema_version: "task_lease_v0", goal_id: request.goal_id, + todo_id: "todo_monitor", owner: "agent-a", status: "active", + idempotency_key: "execution-a", version: 4, lease_epoch: 2, + acquired_at: "2026-09-01T00:00:00Z", expires_at: "2026-09-01T01:00:00Z"}; + await store.commitAuthority({operation_id: "seed-lease", expected_provider_revision: head.provider_revision, + events: [], receipts: [], next_projection: {...head.head, handoff_mode: "hard_lease", leases: [lease]}}); + const fenced = {...request, now: new Date("2026-09-01T00:05:00Z"), + lease_proof: {idempotency_key: "execution-a", expected_version: 4}}; + const before = await store.loadAuthority(); + for (const proof of [null, {idempotency_key: "wrong", expected_version: 4}, + {idempotency_key: "execution-a", expected_version: 3}]) { + assert.equal((await executeCoordinationMonitorPoll(store, {...fenced, lease_proof: proof})).status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + } + const applied = await executeCoordinationMonitorPoll(store, fenced); + assert.equal(applied.status, "applied", JSON.stringify(applied)); + assert.deepEqual((applied.writeback as JsonObject).lease_proof, fenced.lease_proof); + const after = await store.loadAuthority(); + assert.equal(after.status, "loaded"); + if (after.status !== "loaded") return; + assert.deepEqual(after.head.leases, [lease]); + assert.equal((after.head.todos as JsonObject[]).length, 2); + // A receipt settles past work even after the old execution expires; it + // must not execute the observation again or grant another lease. + assert.equal((await executeCoordinationMonitorPoll(store, {...fenced, + now: new Date("2026-09-02T00:00:00Z")})).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), after); + assert.equal((await executeCoordinationMonitorPoll(store, {...fenced, operation_id: "fresh-expired", + now: new Date("2026-09-02T00:00:00Z")})).status, "failed"); +}); From 5b1cca0708b7717ade07dcaa44f69c6de14b7ab7 Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 18 Sep 2026 01:28:04 +0800 Subject: [PATCH 3/4] docs(monitor): explain leased polling and historical settlement boundaries Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- ...shared-goal-authority-state-provider-v0.md | 7 +- ...-goal-authority-state-provider-v0.zh-CN.md | 5 +- .../typescript-control-plane-migration-v0.md | 31 +++++-- ...script-control-plane-migration-v0.zh-CN.md | 23 +++-- .../quota-monitor-observation-receipt-v0.md | 91 +++++++++++++++++++ 5 files changed, 137 insertions(+), 20 deletions(-) diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md index 91db2e1ee2..d13c6a2e01 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2949,8 +2949,9 @@ outbox renderer: display failure is pending, not rollback or successor recreatio Quota consumes the same v0 business receipt for its separate settlement. Validation covers the real CLI with missing display, operation replay after a renderer failure, and complex-data concurrency/lost-acknowledgement recovery on -File, NoKV and isolated real PostgreSQL. Retained Monitor leases and cross-owner -successor claims remain explicitly unsupported. This slice changes neither the +File, NoKV and isolated real PostgreSQL. The L4 extension below closes retained +Monitor lease proof and crash-safe settlement; cross-owner successor claims +remain explicitly unsupported. This slice changes neither the provider default, writer fence nor promotion approval, and does not replace the independent legacy three-arm comparison or D2 soak. @@ -3040,7 +3041,7 @@ or moving a helper is not by itself a package exit. | A / L1: Monitor configuration (this slice) | Existing `todo update` config enters the TS planner/CAS/receipt; delete Python's duplicate intent field catalog. Separate authoring from observed hashes, times and generations. | Ordinary CLI/API, clear/omission, active lease proof, no-op/replay, failed display delivery, complete fixture and real providers. This does not complete delegated Chat or leased polling. | | A / L2: Complete public mutation admission | Inventory actual CLI/Turn/Chat callers; close remaining effect-owned user decisions, delegated owner actions and Monitor lifecycle transitions with validated actor/grant facts. | Build on merged T1 owners, not a generic raw patch. Prove permission rejection and exact caller response; remove replaced Python admission and name every remaining unsupported command. | | A / L3: Canonical lease lifecycle | Standalone acquire/takeover, atomic claim lease admission and maintenance reuse TS facts/decision/materialization and one provider opening fence. Acquire success verifies current execution proof; canonical completion can recover missing display. | Full-head scope conflict, archived/ineffective holders, exact create-CAS retry, stale execution, process loss and real CLI/four-arm rehearsal are covered. [Operation and remaining callers](../../reference/canonical-lease-renew.md). Executor-held external-effect fences remain explicit work; D1–D3/default holds remain. | -| B / L4: Leased Monitor poll and settlement | Compose observation, generation and independent successors with the current lease fence. Reuse the existing quota settlement protocol and exact business receipt. | L2/L3; real polling failure, duplicate/no-change observations, crash between business and quota settlement, and competing writers. Do not pretend separate authorities share a database transaction. | +| B / L4: Leased Monitor poll and settlement | Current execution proof now binds CLI intent, observation/generation/independent-successor CAS and historical business receipt. Quota pending admission is frozen before the business write; recovery preserves that decision after lease retirement. | Existing L3 lease lifecycle, real File/SQLite/PostgreSQL, mixed fixtures, process death between business/quota commits, competing renewal and unchanged polling. [Operation and snapshot rehearsal](../../reference/protocols/quota-monitor-observation-receipt-v0.md). No lease lifecycle effects or quota spend; separate authorities stay separate. Event callers, wider L2 admission and D1–D3/default remain open. | | B / L5: Consumer and display closure | Reconcile #4316, audit Turn/quota/Dashboard/Chat source reads, and finish D1 freshness/recovery through the existing projection outbox. | CLI, Lark/Chat and packaged frontend read back their affected interactions; absent/stale display, empty canonical state, pending projection and data beyond UI limits. Delete post-promotion legacy fallbacks with each consumer. | | A–C / L6: Local durability qualification | Continue contributor-owned #4224/#4328 on the selected SQLite profile; reuse File/NoKV references and complete 7.2's ledger. | Capacity, real process/crash/restore/upgrade, retained receipts/scans, consumer lag, supported runtimes/OS and the separately authorized >=10-day synthetic soak. Missing measurements remain holds. | | A–C / L7: Capture continuity | Resolve #4315 with source-correlated archive retirement and identical lease membership at bootstrap and later writers; activate its row/mutant and complete the mixed-writer/event-source matrix. | Real CLI/File capture, history retained, partial drain unqualified, crash/replay and a new lease after archive/rebootstrap. Keep the legacy migration window provable; T4 cannot be used to skip this row. | diff --git a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md index d86fc08a7c..aacd5c24a4 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.zh-CN.md @@ -2340,7 +2340,8 @@ route planner 本身仍不授予权限。CLI 将已提交回执交给既有 jour 展示失败标为 pending,不回滚提交、不重新生成后继;quota 继续消费同一 v0 业务回执 完成独立记账。验证覆盖真实 CLI 的缺失 display、renderer 失败后的 operation 重放, 及 File/NoKV/真实隔离 PostgreSQL 的复杂数据、并发和丢回执恢复。 -带 lease Monitor、跨 owner claim 等未闭合能力仍明确拒绝;此切片不改变 provider +下述 L4 扩展闭合了 Monitor lease proof 与崩溃后结算;跨 owner claim 仍明确拒绝。 +此切片不改变 provider 默认、writer fence 或 promotion 审批,也不替代三臂 legacy 对照及 D2 soak。 - 从 `loopx/control_plane/todos/provider_projection.py`、既有 Todo-section renderer、 @@ -2410,7 +2411,7 @@ canonical renew 候选,#4328 是 SQLite D2 首批测量/恢复候选;它 | A/L1:Monitor 配置(本切片) | 现有 `todo update` 配置进入 TS planner/CAS/receipt,删除 Python 重复 intent 字段表;区分配置与观察 hash、时间、代数。 | 普通 CLI/API、清除/省略、active lease proof、no-op/replay、展示失败恢复、完整 fixture 和真实 provider。不宣称完成委托 Chat 或 leased polling。 | | A/L2:公共 mutation admission 闭合 | 盘点 CLI/Turn/Chat 实际 caller;以可信 actor/grant 事实闭合剩余 effect-owned 用户决策、委托 owner 动作和 Monitor lifecycle。 | 复用已合并 T1 owner,不开通通用 raw patch;验证权限拒绝和 caller 响应,删除替代的 Python admission,列全未支持命令。 | | A/L3:canonical lease 生命周期 | 独立 acquire/接管、原子 claim 的 lease 准入及维护复用 TS facts/decision/materializer 与同一 provider opening fence。Acquire 成功必须校验当前执行 proof;canonical 完成可恢复缺失展示。 | 已覆盖完整 head scope 冲突、归档/失效 holder、创建 CAS 原样重试、旧执行、进程中断、真实 CLI 与四臂演练。[操作及剩余 caller](../../reference/canonical-lease-renew.md)。跨外部 effect 的 executor 持锁 fence 仍为明确工作;保留 D1–D3/default hold。 | -| B/L4:leased Monitor poll 与 settlement | 组合观察、变化代数、独立 successor 和现有 lease fence;复用 quota settlement 与精确业务回执。 | L2/L3;真实 polling 失败、重复/无变化、业务提交到 quota settlement 间崩溃和并发。不能假装不同 authority 共享一个数据库事务。 | +| B/L4:leased Monitor poll 与 settlement | 当前 execution proof 贯穿 CLI intent、观察/generation/独立 successor CAS 和历史业务回执;业务写入前冻结 quota 准入,租约结束后仍按原决策恢复结算。 | 既有 L3 lease lifecycle、真实 File/SQLite/PostgreSQL、混合 fixture、业务与 quota 间真实进程退出、并发 renewal 和 unchanged poll;见[操作与快照演练](../../reference/protocols/quota-monitor-observation-receipt-v0.md)。不操作 lease lifecycle、不消耗 quota,不把两个 authority 假装成同一事务。Event caller、更广 L2 准入及 D1–D3/default 仍开放。 | | B/L5:consumer 与展示闭合 | 核对 #4316,审计 Turn/quota/Dashboard/Chat 的来源,复用 projection outbox 完成 D1 新鲜度和恢复。 | 验证 CLI、Lark/Chat、打包 frontend 的受影响交互;缺失/陈旧展示、权威空状态、pending 投影及超过 UI 上限的数据。逐个删除晋升后的 legacy fallback。 | | A–C/L6:本地持久化资格 | 延续 contributor 认领的 #4224/#4328,在选定 SQLite profile 上补齐第 7.2 节 ledger,复用 File/NoKV 对照。 | capacity、真实进程/crash/restore/upgrade、历史 receipt/scan、consumer lag、支持的 runtime/OS,以及另行授权的 >=10 天合成 soak。缺项继续 hold。 | | A–C/L7:capture 连续性 | 修复 #4315:归档的源事务明确退休 lease 引用,bootstrap 与后续 writer 使用一致成员范围;执行 row/mutant 和 mixed-writer/event-source 矩阵。 | 真实 CLI/File capture、保留历史、半完成 drain 不合格、crash/replay,以及归档/rebootstrap 后再申请 lease。不能借 T4 跳过迁移窗口证明。 | diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 6fc16bca00..cc8ceac6cb 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -675,14 +675,29 @@ Retrying the original operation recovers the original successors instead of creating new work. A fresh observation with no successor remains valid. User gates use the existing actor-bound scope, never an inferred global gate. -Boundaries still open: any retained Monitor lease fails closed in this native -operation; cross-owner successor claims are not implicitly authorized. Unpromoted -Goals retain their legacy writer. Quota accounting stays in its existing -preflight/writeback/settlement protocol and reuses the v0 receipt shape and raw -observation identity. Canonical commit success is independent of pending Markdown -delivery. This does not finish all T2 commands or authorize whole-Goal promotion. - -- Finish the retained lease and event callers of `monitor_poll_writeback.py`. +The leased Monitor path now reuses the current nonterminal lease fence with +Todo metadata updates. Public `quota monitor-poll` carries the execution key and +version through its pending plan, canonical transaction and business receipt. +Observation, generation and independent successors share one CAS; the lease is +unchanged. Canonical due monitors are selectable again, but scheduling does not +grant mutation authority. Runtime time, not the supplied observation timestamp, +determines whether a fresh write's lease is active. + +Quota preflight now freezes its admitted decision in a versioned pending receipt. +Recovery settles the original business receipt even after the Monitor becomes +not due or its lease is released; it never substitutes a new lease or re-runs +business effects. Proof-less v0 request identities and completed receipts remain +compatible. Old pending receipts without an admission basis retain current-state +admission and explicitly report when historical recovery cannot be proven. +See [Monitor observation and recovery](../../reference/protocols/quota-monitor-observation-receipt-v0.md). + +Boundaries still open: cross-owner successor claims are not implicitly authorized; +unpromoted Goals retain their legacy writer and reject explicit lease proof. +Quota and business authority remain separate recoverable transactions. Canonical +commit success is independent of pending Markdown delivery. This does not finish +all T2 commands or authorize whole-Goal promotion. + +- Finish the retained event callers of `monitor_poll_writeback.py`. Reuse existing monitor generation, independent-successor and settlement owners. Compose one transaction rather than adding a second monitor engine. - Preserve unchanged polling/reschedule behavior, generation fences, diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md index c6295b7dcf..deb484f62c 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.zh-CN.md @@ -529,13 +529,22 @@ generation,不能对相同证据重复声明 `material_change=true` 就继续 原 operation 重试恢复原后继,不创建新工作;不附带后继的新 observation 仍可接受。 User gate 复用既有 actor-bound scope,不推导全局 gate。 -尚未闭合:原生操作遇到任何保留的 Monitor lease 仍 fail closed,不隐式授权跨 owner -的 successor claim;未 promotion Goal 仍走旧 writer。Quota 记账继续使用现有 -preflight/writeback/settlement 协议,沿用 v0 回执形状及原始 observation identity。 -Canonical 提交成功独立于 Markdown delivery pending。这不代表全部 T2 命令或整 Goal -promotion 已完成。 - -- 继续闭合 `monitor_poll_writeback.py` 保留的 lease 与 event caller,复用 monitor +带 lease Monitor 现与 Todo metadata update 共用当前非终结 lease fence。公开 +`quota monitor-poll` 将 execution key/version 贯穿 pending plan、canonical transaction +和业务回执;观察、generation 与独立后继在同一 CAS 提交,lease 保持不变。 +Canonical 到期 Monitor 恢复可选,但调度不授予写权限;租约是否有效取 runtime +当前时间,不取调用方提交的观察时间。 + +Quota preflight 将原始准入决策冻结到版本化 pending receipt;即使 Monitor 已不再 +到期或 lease 已释放,恢复仍可凭原业务回执结算,不替换租约、不重做业务。 +无 proof 的 v0 request identity 和已完成回执保持兼容;无原准入依据的旧 pending +沿用当前准入,无法证明历史恢复时明确报错。见[观察与恢复协议](../../reference/protocols/quota-monitor-observation-receipt-v0.md)。 + +尚未闭合:不隐式授权跨 owner successor claim;未晋升 Goal 保留旧 writer,并拒绝 +显式 lease proof。业务与 quota 仍是分别可恢复的事务,canonical 成功独立于 Markdown +delivery pending;这不代表全部 T2 命令或整 Goal promotion 已完成。 + +- 继续闭合 `monitor_poll_writeback.py` 保留的 event caller,复用 monitor generation、独立 successor 和 settlement owner,组成一笔事务,不建第二套引擎。 - 保持 unchanged poll/reschedule、generation fence、material-change successor 去重和可归属 settlement。Monitor 不是 delivery 执行任务;独立 advancement Todo diff --git a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md index 56756b508b..9456393c6e 100644 --- a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md +++ b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md @@ -42,6 +42,66 @@ and replay idempotently in one settlement Turn without spending quota, while a guard replay continues to select the original advancement Todo. Existing wrong-Todo tests for receipt-bound monitor Turns must remain passing. +### Canonical leased observations and recovery + +For a promoted Goal, an existing Monitor execution can supply its current lease +key and version. Read the existing lease first; this command neither acquires +nor renews/releases it: + +```bash +loopx task-lease inspect --goal-id "$GOAL" --todo-id "$MONITOR" +loopx quota monitor-poll --goal-id "$GOAL" --agent-id "$AGENT" \ + --todo-id "$MONITOR" --result-hash "$OBSERVED_HASH" \ + --task-lease-idempotency-key "$LEASE_KEY" --task-lease-expected-version "$VERSION" \ + --turn-instance-id "$TURN" --execute +loopx todo list --goal-id "$GOAL" +``` + +The caller must already have an admitted Monitor Turn and any required host +capabilities. External observation remains the caller's effect; `--result-hash` +does not perform a network poll. Add `--material-change --next-agent-todo ... +--next-action-kind ...` only for changed evidence that needs independent work. + +- The registered actor, claim, binding, exclusions and current lease are checked + against the same canonical revision as observation/generation/successor writes. + A stale key/version, expired lease, soft-claim lease request or missing hard-lease + proof rejects the complete mutation. Observation timestamps cannot revive a + lease. The existing lease remains byte-for-byte unchanged. + This closes a previous hard-mode hole: a Monitor with no lease could write + merely because the old implementation only rejected retained leases. +- Canonical due monitors no longer carry the blanket unsupported-writeback + marker. The scheduler can select them; selection itself is not a lease grant. +- New quota pending receipts use `quota_monitor_poll_pending_admission_v1` to + freeze the original admitted decision and provider plan before business writes. + After process loss, retry the **same Turn, observation, intent and original + proof**. Its durable business receipt can settle after lease release/expiry, + with the original decision as the event's `before` state. A new observation + still needs current authority. Quota settlement never spends a delivery slot. +- Lease-bearing quota/provider requests and provider plans use v1; proof-less + requests retain v0, including their original digests. Completed v0 settlement + receipts remain readable. Old v0 *pending* receipts have no frozen admission: + they recover if current admission still holds, otherwise report + `legacy_monitor_admission_unavailable` and preserve evidence for reconciliation. + Do not fabricate a historical decision, delete the receipt or repeat a known + committed effect under a fresh identity to bypass that condition. + +File, SQLite and service-opened PostgreSQL share these transaction semantics. +This is not provider activation, new-Goal default selection, cross-host service +qualification or whole-Goal promotion. Unpromoted legacy Goals reject explicit +lease proof. Before downgrading to a version without this protocol, finish v1 +pending settlements; older binaries cannot interpret them. Never disable a +writer fence to recover by writing an older Markdown projection. + +The bounded real-source rehearsal clones read-only source state and adds only +synthetic records in disposable runtimes. It requires an isolated PostgreSQL URL: + +```bash +LOOPX_TEST_POSTGRES_URL="$DISPOSABLE_POSTGRES_URL" \ +uv run --extra test python examples/control_plane/authority-monitor-poll-rehearsal.py \ + --registry "$REGISTRY" --goal-id "$GOAL" --baseline-repo "$FROZEN_BASELINE" \ + --execute-isolated-postgresql +``` + ## 中文 ### 问题 @@ -74,3 +134,34 @@ wrong-Todo tests for receipt-bound monitor Turns must remain passing. CLI 端到端测试必须证明:多个到期 monitor 能在同一结算 Turn 中分别更新周期并 幂等重放、全程不消耗配额;随后重放 guard 仍选择原 advancement Todo。同时, receipt-bound monitor Turn 的错误 Todo 替换测试必须继续通过。 + +### Canonical 带租约观察与恢复 + +已晋升 Goal 的 Monitor 若已有执行租约,可用上述命令读取租约,并传入当前 key 和 +version。调用方须已有合法 Monitor Turn 和所需 host capabilities;外部观察由调用方 +执行,`--result-hash` 不会自行发起网络轮询。只有新证据需要独立工作时才增加 +`--material-change --next-agent-todo ... --next-action-kind ...`。 + +- 注册 actor、claim、binding、exclusion 和当前租约在同一 canonical revision 上校验; + 观察、generation 与 successor 原子提交。过期、错 key/version、soft-claim 下携带租约 + 或 hard-lease 下缺少凭据均整笔拒绝。旧观察时间不能复活租约;租约内容保持不变。 + 同时修复旧 hard-mode 缺口:过去仅检查“是否保留 lease”,反而让没有租约的 Monitor + 在 hard mode 下写入;现在该分支明确拒绝。 +- Canonical 到期 Monitor 不再被统一标为“不支持 writeback”,可进入调度选择;被选中 + 本身不授予执行租约。此默认变化只影响 canonical Monitor 的调度可见性。 +- 新 pending receipt 使用 `quota_monitor_poll_pending_admission_v1`,在业务写入前 + 冻结原准入决策与 provider plan。进程丢失后,用**相同 Turn、观察、intent 和原凭据** + 重试;已提交业务的回执可在租约释放/过期后补结算,event 的 `before` 仍是原决策。 + 新观察依旧需要当前权限;结算不消耗 delivery slot。 +- 携带租约的 quota/provider request 和 provider plan 使用 v1;无 proof 的请求保留 + v0 及原 digest,已完成的 v0 结算继续可读。旧 v0 pending 没有冻结准入:当前准入仍 + 成立时可恢复,否则报告 `legacy_monitor_admission_unavailable`,保留证据供核对。 + 不得伪造历史准入、删除回执,或换新 identity 重做已知提交的业务来绕过它。 + +File、SQLite 与经 service 打开的 PostgreSQL 共用此事务语义;这不等于 provider +启用、新 Goal 默认切换、跨 host 资格或整 Goal 晋升。未晋升 Goal 拒绝显式租约凭据。 +降级到不支持此协议的版本前,应先完成 v1 pending settlement;旧程序无法解释它们。 +不得通过关闭 writer fence、恢复旧 Markdown 写入来绕过恢复要求。 + +上述 rehearsal 命令只读源数据,在临时 runtime 中增加合成记录,比较冻结基线及三个 +真实后端;要求隔离 PostgreSQL,输出限于计数和摘要,不写回原 Goal。 From 11487a643128ebcca562814df1a521ca1d75697d Mon Sep 17 00:00:00 2001 From: huangruiteng <14976749+huangruiteng@users.noreply.github.com> Date: Fri, 18 Sep 2026 01:31:21 +0800 Subject: [PATCH 4/4] test(monitor): verify dependent wakeup and state settlement CAS limits Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com> --- .../protocols/quota-monitor-observation-receipt-v0.md | 3 +++ tests/control_plane_ts/monitor_poll_lease_conformance.ts | 9 +++++++++ 2 files changed, 12 insertions(+) diff --git a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md index 9456393c6e..cd5236d596 100644 --- a/docs/reference/protocols/quota-monitor-observation-receipt-v0.md +++ b/docs/reference/protocols/quota-monitor-observation-receipt-v0.md @@ -77,6 +77,8 @@ does not perform a network poll. Add `--material-change --next-agent-todo ... proof**. Its durable business receipt can settle after lease release/expiry, with the original decision as the event's `before` state. A new observation still needs current authority. Quota settlement never spends a delivery slot. + The existing quota-index CAS remains enforced: an intervening index write + causes an explicit conflict, not an unconditional settlement append. - Lease-bearing quota/provider requests and provider plans use v1; proof-less requests retain v0, including their original digests. Completed v0 settlement receipts remain readable. Old v0 *pending* receipts have no frozen admission: @@ -153,6 +155,7 @@ version。调用方须已有合法 Monitor Turn 和所需 host capabilities; 冻结原准入决策与 provider plan。进程丢失后,用**相同 Turn、观察、intent 和原凭据** 重试;已提交业务的回执可在租约释放/过期后补结算,event 的 `before` 仍是原决策。 新观察依旧需要当前权限;结算不消耗 delivery slot。 + 既有 quota-index CAS 继续生效;若期间有其他 index 写入,则明确冲突,不能无条件追加。 - 携带租约的 quota/provider request 和 provider plan 使用 v1;无 proof 的请求保留 v0 及原 digest,已完成的 v0 结算继续可读。旧 v0 pending 没有冻结准入:当前准入仍 成立时可恢复,否则报告 `legacy_monitor_admission_unavailable`,保留证据供核对。 diff --git a/tests/control_plane_ts/monitor_poll_lease_conformance.ts b/tests/control_plane_ts/monitor_poll_lease_conformance.ts index dddee8381d..0475c6286d 100644 --- a/tests/control_plane_ts/monitor_poll_lease_conformance.ts +++ b/tests/control_plane_ts/monitor_poll_lease_conformance.ts @@ -5,6 +5,7 @@ import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; import type {AuthorityStore, AuthorityStoreCommit} from "../../loopx/control_plane/coordination/authority_store.ts"; import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; import {coordinationTodoReadModel} from "../../loopx/control_plane/coordination/coordination_projection.ts"; +import {evaluateTodoResumeConditions} from "../../loopx/control_plane/todos/resume_condition.ts"; import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; import {productionScaleLeasedMonitorFixture} from "./production_scale_coordination_fixture.ts"; @@ -43,6 +44,14 @@ export function registerLeasedMonitorConformance(provider: string, factory: Auth original.filter(todo => todo.todo_id !== fixture.target)); assert.equal(rows.length, original.length + 2); assert.equal(rows.find(todo => todo.todo_id === fixture.target)?.material_change_generation, 5); + const ready = [original, rows].map(source => { + const result = evaluateTodoResumeConditions({schema_version: "todo_resume_evaluation_request_v0", + items: source.filter(todo => todo.todo_id === fixture.dependent), source_items: source, + kinds: ["monitor_changed"]}); + const condition = (result.conditions as JsonObject[])[0]!.condition as JsonObject; + return condition.satisfied; + }); + assert.deepEqual(ready, [false, true], "the committed generation releases the existing dependent's wait"); const unchanged = await executeCoordinationMonitorPoll(store, {...request, operation_id: "unchanged", observation: {...request.observation, generated_at: "2026-09-01T02:00:00Z", material_change: false}, intent: {}}); assert.equal(unchanged.status, "applied");