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 2d9ab64641..aed207e6ad 100644 --- a/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md +++ b/docs/architecture/rfcs/shared-goal-authority-state-provider-v0.md @@ -2642,10 +2642,17 @@ source paths, authorize monitor writeback, or change provider/promotion holds. **D1 — qualify permanent projection delivery; may overlap T1/T2.** -The T2 monitor successor route owner is now shared across preflight, the legacy -effect adapter and receipt checks. Its result proves only normalized intent, -not actor authority, provider commit or atomic monitor-plus-successor durability. -Keep the monitor writer fence and promotion hold until that transaction closes. +T2 now commits a lease-free native Monitor observation and its independent +successors in one canonical CAS/receipt; the route planner alone still grants +no authority. The CLI delivers committed state through the existing journal/ +outbox renderer: display failure is pending, not rollback or successor recreation. +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 +provider default, writer fence nor promotion approval, and does not replace the +independent legacy three-arm comparison or D2 soak. - Start from `loopx/control_plane/todos/provider_projection.py`, the existing Todo-section renderer and canonical journal/outbox. #4097 already recovers 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 031191057c..9016d73a44 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 @@ -2095,10 +2095,13 @@ summary,之前消费 legacy summary;真实 CLI 覆盖容量变化和 promote **D1 — 资格化永久投影交付,可与 T1/T2 重叠推进。** -T2 monitor successor 的路由 owner 已由 preflight、legacy effect adapter 和回执校验 -共享。其结果仅证明规范化 intent,不证明 actor authority、provider commit 或 -monitor-plus-successor 原子持久化;事务闭合前继续保留 monitor writer fence 与 -promotion hold。 +T2 的无 lease 原生 Monitor 观察与独立后继现由同一 canonical CAS/receipt 提交; +route planner 本身仍不授予权限。CLI 将已提交回执交给既有 journal/outbox renderer, +展示失败标为 pending,不回滚提交、不重新生成后继;quota 继续消费同一 v0 业务回执 +完成独立记账。验证覆盖真实 CLI 的缺失 display、renderer 失败后的 operation 重放, +及 File/NoKV/真实隔离 PostgreSQL 的复杂数据、并发和丢回执恢复。 +带 lease Monitor、跨 owner claim 等未闭合能力仍明确拒绝;此切片不改变 provider +默认、writer fence 或 promotion 审批,也不替代三臂 legacy 对照及 D2 soak。 - 从 `loopx/control_plane/todos/provider_projection.py`、既有 Todo-section renderer、 canonical journal/outbox 入手。复用 #4097 已有的缺失 Todo section 恢复及 diff --git a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md index 52904ac60a..4522820456 100644 --- a/docs/architecture/rfcs/typescript-control-plane-migration-v0.md +++ b/docs/architecture/rfcs/typescript-control-plane-migration-v0.md @@ -449,11 +449,30 @@ are compared as the same route at readback. The original wire observation still owns the v0 replay digest; normalization must not silently invalidate pending receipts. The node-independent repository/bootstrap codec remains separately characterized, not replaced by a runtime dependency. -This is **not** the T2 atomic transaction: monitor mutation and successor writes -still use existing fenced effects. Cross-effect crash recovery, native writer -closure and whole-Goal promotion remain held; do not infer them from a route plan. - -- Inventory `monitor_poll_writeback.py` and its event/Todo/lease callers. +The native `coordination.local_authority.monitor_poll` transaction now commits a +lease-free Monitor observation and its requested independent successors against +one canonical revision, with one CAS and durable operation receipt. It composes +the existing generation, successor-route, User authoring-scope and Todo-create +planners. Public create and Monitor batches share create admission/duplicate +planning; target selection is shared by legacy preflight and native commit. +Python only routes provider intent and drains the existing projection outbox. + +Explicit semantic corrections: completed/archived Monitor targets are rejected; +target-key selection ignores finished history but never guesses between live matches; +successor authoring requires an actually advanced material-change generation, +not merely a repeated `material_change=true` assertion for the same evidence. +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`. 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 7a6e27c4c3..1cfab1cedb 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 @@ -346,10 +346,25 @@ guard/resolver 及 TS 回执端独立的默认值/capability 解释。非法 c action/claim/capability 别名和 Git transport 在回执核对时指向同一路由。v0 replay digest 仍绑定原始 wire observation,不能因规范化而悄悄使 pending receipt 失效。 无需 Node 的 repository/bootstrap codec 暂留并做跨运行时对照,不引入启动依赖。 -这**不是** T2 原子事务:monitor mutation 和 successor 写入仍通过既有 fenced effect -执行;跨 effect crash 恢复、native writer 闭合及整 Goal promotion 仍未放行。 - -- 盘点 `monitor_poll_writeback.py` 及 event/Todo/lease caller,复用 monitor +原生 `coordination.local_authority.monitor_poll` 现将无 lease Monitor 的观察及请求的 +独立后继,绑定同一个 canonical revision,以一次 CAS 和持久 operation receipt +提交。它组合已有 generation、successor route、User authoring scope 和 Todo create +planner;单项 create 与 Monitor 批次共用创建准入/语义去重,legacy preflight 与 +native commit 共用目标选择。Python 只路由意图并交付既有 projection outbox。 + +明确的语义修正:拒绝已完成/归档的 Monitor;target-key 选择排除结束的历史项, +但多个活跃匹配仍要求显式 id;创建后继必须实际推进 material-change +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 generation、独立 successor 和 settlement owner,组成一笔事务,不建第二套引擎。 - 保持 unchanged poll/reschedule、generation fence、material-change successor 去重和可归属 settlement。Monitor 不是 delivery 执行任务;独立 advancement Todo diff --git a/loopx/cli_commands/quota.py b/loopx/cli_commands/quota.py index 8c168f2931..7630daec0e 100644 --- a/loopx/cli_commands/quota.py +++ b/loopx/cli_commands/quota.py @@ -45,6 +45,7 @@ from ..control_plane.coordination.legacy_writer_fence import ( LegacyCoordinationWriterFenced, ) +from ..control_plane.coordination.local_authority import LocalCoordinationAuthorityUnavailable from ..control_plane.effect_runtime import EffectRuntimeRejected from ..control_plane.scheduler.execution_context import ( GUIDED_START_TURN_RUNTIME_PROFILES, @@ -241,7 +242,7 @@ def _quota_failure_payload( ) if error.agent_id is not None: payload["agent_id"] = error.agent_id - elif isinstance(error, LegacyCoordinationWriterFenced): + elif isinstance(error, (LegacyCoordinationWriterFenced, LocalCoordinationAuthorityUnavailable)): payload.update( { "error_code": error.code, diff --git a/loopx/cli_commands/quota_context.py b/loopx/cli_commands/quota_context.py index d32f5fb594..e3c6004e30 100644 --- a/loopx/cli_commands/quota_context.py +++ b/loopx/cli_commands/quota_context.py @@ -5,8 +5,8 @@ from dataclasses import dataclass from pathlib import Path -from ..control_plane.coordination.legacy_writer_fence import ( - require_legacy_coordination_write_allowed, +from ..control_plane.scheduler.provider_monitor_poll import ( + require_monitor_poll_source_available, ) from ..control_plane.quota.error_codes import QuotaCommandValidationError from ..control_plane.runtime.status_projection_cache import ( @@ -225,10 +225,9 @@ def prepare_quota_command_context( runtime_root_override=runtime_root_arg, ) if command == "monitor-poll" and args.execute and (args.todo_id or args.target_key): - # This command still uses the legacy Todo writer. Preserve its typed - # rejection before collecting a promoted read model (which may itself - # be unavailable). The writer repeats the check under its mutation lock. - require_legacy_coordination_write_allowed( + # Canonical availability precedes unrelated status/quota preparation. + # The eventual transaction repeats its fence check under the writer lock. + require_monitor_poll_source_available( runtime_root=runtime_root, goal_id=args.goal_id, ) status_goal_id = args.goal_id if command not in {"status", "plan"} else None diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index a1f2bb130b..0c4449527a 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -3,6 +3,8 @@ import { ShadowManagementError, requireShadowPrimaryWriteAllowed, shadowMaintena import { isAbsolute, join } from "node:path"; import type { JsonObject } from "../effect_program.ts"; +import {executeCoordinationMonitorPoll, COORDINATION_MONITOR_POLL_REQUEST_SCHEMA, + COORDINATION_MONITOR_POLL_RESULT_SCHEMA} from "./todo_monitor_poll.ts"; import { requireJsonObject } from "../runtime_decode.ts"; import { LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA, @@ -101,6 +103,34 @@ async function withCanonicalWriter(root: string, goalId: string, dryRun: bool }); } +/** Monitor observation and successors share the existing writer/fence lifetime. */ +export async function pollLocalCoordinationMonitor(value: unknown, + dependencies: LocalAuthorityRuntimeDependencies = {}): Promise { + 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"); + 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"); + const registered = input.registered_agents.map(agent => claimAgentValue(agent, "registered agent")); + return await withCanonicalWriter(root, goalId, input.dry_run === true, async () => ({ + ...await executeCoordinationMonitorPoll(dependencies.createStore?.(authorityDirectory(root), goalId) ?? + new FileAuthorityStore(authorityDirectory(root), goalId), { + goal_id: goalId, operation_id: requireAuthorityStoreId(input.operation_id, "operation id"), + actor_agent_id: input.actor_agent_id == null ? null : claimAgentValue(input.actor_agent_id, "actor_agent_id"), + registered_agents: registered, dry_run: input.dry_run as boolean, + observation: requireJsonObject(input.observation, "Monitor observation"), + intent: requireJsonObject(input.intent, "Monitor successor intent"), + }), ...evidence, + })); + } catch (error) { + return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status: "failed", changed: false, + reason_code: error instanceof ShadowManagementError ? error.reason_code : "invalid_local_monitor_poll_request", + reason: error instanceof Error ? error.message : String(error), ...evidence}; + } +} + interface LocalAuthorityRuntimeDependencies { createStore?: (directory: string, goalId: string) => AuthorityStore; createShadowStore?: (directory: string, goalId: string) => AuthorityStore; diff --git a/loopx/control_plane/coordination/todo_create.ts b/loopx/control_plane/coordination/todo_create.ts index 30c74c08cf..eafb21d76e 100644 --- a/loopx/control_plane/coordination/todo_create.ts +++ b/loopx/control_plane/coordination/todo_create.ts @@ -166,6 +166,26 @@ function createCandidate( }; } +/** In-process create planning for a caller-owned canonical transaction. Never + * commits or authorizes the enclosing operation; single create and Monitor + * batches share record admission, attribution and semantic duplicate rules. */ +export function planCoordinationTodoCreate( + rawInput: CoordinationTodoCreateInput, + todos: ReadonlyMap, + readModelSchema: unknown, +): CoordinationTodoCreateResult { + const input = normalizeCreateInput(rawInput); + const duplicate = [...todos.values()].find((todo) => + todo.role === input.todo.role && todo.archive_state === "active" && + todo.status !== "done" && todo.status !== "deferred" && todo.text === input.todo.text); + if (duplicate !== undefined) return semanticDuplicateResult(input.todo, duplicate, null, null); + if (todos.has(String(input.todo.todo_id))) { + return failure("todo_already_exists", "Todo id already exists in canonical authority", {todo_id: input.todo.todo_id}); + } + return {schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, status: "planned", + changed: true, todo_id: input.todo.todo_id, todo: createCandidate(input, readModelSchema)}; +} + async function commitCreate( store: AuthorityStore, input: CoordinationTodoCreateInput, @@ -241,18 +261,10 @@ export async function executeCoordinationTodoCreate( ); } const todoId = requireAuthorityStoreId(input.todo.todo_id, "todo id"); - const duplicate = [...projection.todos.values()].find((todo) => - todo.role === input.todo.role && todo.archive_state === "active" && - todo.status !== "done" && todo.status !== "deferred" && todo.text === input.todo.text - ); - if (duplicate !== undefined) { - return semanticDuplicateResult(input.todo, duplicate, head.provider_revision, head.cursor); - } - if (projection.todos.has(todoId)) { - return failure("todo_already_exists", "Todo id already exists in canonical authority", {todo_id: todoId}); - } const readModel = canonicalAuthorityObject(head.head.todo_read_model, "Todo read model"); - const created = createCandidate(input, readModel.schema_version); + const plan = planCoordinationTodoCreate(input, projection.todos, readModel.schema_version); + if (plan.status !== "planned") return {...plan, provider_revision: head.provider_revision, cursor: head.cursor}; + const created = canonicalAuthorityObject(plan.todo, "created Todo"); if (input.dry_run) { return { schema_version: COORDINATION_TODO_CREATE_RESULT_SCHEMA, diff --git a/loopx/control_plane/coordination/todo_monitor_poll.ts b/loopx/control_plane/coordination/todo_monitor_poll.ts new file mode 100644 index 0000000000..264ee05390 --- /dev/null +++ b/loopx/control_plane/coordination/todo_monitor_poll.ts @@ -0,0 +1,177 @@ +/** One canonical transaction for an observation and its independent successors. + * Network polling, quota settlement and display delivery are separate effects. */ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore, AuthorityStoreReceiptResult} from "./authority_store.ts"; +import {canonicalAuthorityObject, canonicalAuthoritySha256, requireAuthorityStoreId} from "./authority_store_codec.ts"; +import {normalizeRegisteredTodoAgents, normalizeTodoAgent} from "./todo_agents.ts"; +import {indexCoordinationProjection, prepareCoordinationProjectionCommit, validateCoordinationTodoReadModel} from "./coordination_projection.ts"; +import {TODO_DOMAIN_ITEM_SCHEMA} from "./coordination_state_contract.ts"; +import {planCoordinationTodoCreate} from "./todo_create.ts"; +import {planMonitorMetadata, TODO_MONITOR_METADATA_REQUEST_SCHEMA} from "../todos/monitor_metadata.ts"; +import {planMonitorSuccessor, selectMonitorTodo, MONITOR_SUCCESSOR_REQUEST_SCHEMA} from "../scheduler/monitor_successor.ts"; +import {optionalNonEmptyString, requireBoolean} from "../runtime_decode.ts"; +import {planTodoAuthoringScope, TODO_AUTHORING_SCOPE_REQUEST_SCHEMA} from "../todos/authoring_scope.ts"; + +export const COORDINATION_MONITOR_POLL_REQUEST_SCHEMA = "loopx_coordination_monitor_poll_request_v0"; +export const COORDINATION_MONITOR_POLL_RESULT_SCHEMA = "loopx_coordination_monitor_poll_result_v0"; +const RECEIPT_SCHEMA = "loopx_coordination_monitor_poll_receipt_v0"; + +export interface CoordinationMonitorPollInput { + goal_id: string; + operation_id: string; + actor_agent_id: string | null; + registered_agents: readonly string[]; + dry_run: boolean; + observation: JsonObject; + intent: JsonObject; +} + +function failure(reason_code: string, reason: string): JsonObject { + return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status: "failed", changed: false, reason_code, reason}; +} + +function replay(receipt: AuthorityStoreReceiptResult, input: CoordinationMonitorPollInput, + hash: string, status: "applied" | "recovered" | "replayed"): JsonObject | null { + if (receipt.status === "missing") return null; + if (receipt.status !== "found") return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...receipt}; + const original = receipt.receipts[0]; + if (receipt.receipts.length !== 1 || original?.schema_version !== RECEIPT_SCHEMA || + original.goal_id !== input.goal_id || original.operation_id !== input.operation_id || original.request_sha256 !== hash) { + return failure("coordination_operation_identity_mismatch", "operation id already names a different Monitor observation or successor intent"); + } + return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, status, + changed: status !== "replayed", provider_revision: receipt.provider_revision, cursor: receipt.cursor, + writeback: {...canonicalAuthorityObject(original.writeback, "Monitor writeback"), provider_replayed: status === "replayed"}, + projection_delivery: "pending", projection_source: "committed_authority_journal"}; +} + +function normalize(raw: CoordinationMonitorPollInput): CoordinationMonitorPollInput { + 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")}; + 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", + "next_continuation_policy", "next_target_key", "next_claimed_by", "next_user_todo", "next_user_task_class"]); + for (const key of Object.keys(input.intent)) if (!intentFields.has(key)) throw new Error(`unsupported Monitor successor field: ${key}`); + return input; +} + +function planWriteback(input: CoordinationMonitorPollInput, head: JsonObject) { + const indexed = indexCoordinationProjection(head, input.goal_id); + validateCoordinationTodoReadModel(head, input.goal_id); + const observation = input.observation; + const monitor = selectMonitorTodo([...indexed.todos.values()], + optionalNonEmptyString(observation.todo_id, "todo_id"), optionalNonEmptyString(observation.target_key, "target_key")); + const actor = input.actor_agent_id; + // Observation is claim-neutral, not lease acquisition or execution. Never + // infer a delegated actor or turn an existing execution lease into permission. + if (!actor || !input.registered_agents.includes(actor)) throw new Error("Monitor observation requires a registered actor"); + 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 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}}); + const intent = canonicalAuthorityObject(successorPlan.intent, "Monitor successor intent"); + const route = canonicalAuthorityObject(successorPlan.agent_route, "Monitor successor route"); + const monitorPlan = planMonitorMetadata({schema_version: TODO_MONITOR_METADATA_REQUEST_SCHEMA, + existing: monitor, role: "agent", task_class: "continuous_monitor", enforce_boundedness: false, + observation: {...observation, monitor_effect_id: input.operation_id}}); + const transition = canonicalAuthorityObject(monitorPlan.transition, "Monitor transition"); + if ((intent.next_agent_todo || intent.next_user_todo) && transition.material_change_applied !== true) { + throw new Error("successor authoring requires a new material-change generation; poll without successor options for unchanged evidence"); + } + const metadata = canonicalAuthorityObject(monitorPlan.metadata, "Monitor metadata"); + // Generation is an integer in the persisted Todo contract. The older + // observation tokens remain strings (including "0" and "false"); retain + // their existing wire types so permanent Markdown projection is lossless. + if (metadata.material_change_generation != null) metadata.material_change_generation = Number(metadata.material_change_generation); + const updated: JsonObject = {...monitor, last_actor_agent_id: actor, updated_at: observation.generated_at}; + for (const [key, value] of Object.entries(metadata)) { + if (value === null) delete updated[key]; else updated[key] = value; + } + const reason = optionalNonEmptyString(observation.reason_summary, "reason_summary"); + if (reason) updated.reason = reason; + const nextTodos: JsonObject[] = []; + const plannedTodos = new Map(indexed.todos); + const mutations: {kind: "todo_upsert"; todo: JsonObject}[] = [{kind: "todo_upsert", todo: updated}]; + const readModel = canonicalAuthorityObject(head.todo_read_model, "Todo read model"); + for (const role of ["agent", "user"] as const) { + const text = intent[role === "agent" ? "next_agent_todo" : "next_user_todo"]; + if (!text) continue; + const id = `todo_${canonicalAuthoritySha256({monitor: monitor.todo_id, + generation: transition.material_change_generation, role}).slice(0, 24)}`; + const todo: JsonObject = {schema_version: TODO_DOMAIN_ITEM_SCHEMA, todo_id: id, role, text, + status: "open", done: false, archive_state: "active", + task_class: role === "agent" ? "advancement_task" : intent.next_user_task_class}; + if (role === "agent") { + Object.assign(todo, Object.fromEntries(Object.entries(route).filter(([, value]) => value !== null)), + {unblocks_todo_id: monitor.todo_id}); + } else { + // Reuse public authoring scope: an actor-bound gate, never an inferred + // all-agent/global gate. User actions retain their actor binding too. + const scope = planTodoAuthoringScope({schema_version: TODO_AUTHORING_SCOPE_REQUEST_SCHEMA, + command: "create", role, todo: {}, goal_id: input.goal_id, registered_agents: input.registered_agents, + intent: {task_class: todo.task_class, actor_agent_id: actor}}); + for (const key of ["bound_agent", "blocks_agent", "goal_bound", "global_gate"]) { + if (scope[key] != null && scope[key] !== false) todo[key] = scope[key]; + } + if (todo.task_class === "user_gate") Object.assign(todo, {action_kind: "gate", unblocks_todo_id: monitor.todo_id}); + } + const created = planCoordinationTodoCreate({goal_id: input.goal_id, operation_id: input.operation_id, + actor_agent_id: actor, registered_agents: input.registered_agents, dry_run: input.dry_run, + now: new Date(String(observation.generated_at)), todo}, plannedTodos, readModel.schema_version); + if (created.status !== "planned" && created.status !== "no_change") throw new Error(String(created.reason)); + const record = canonicalAuthorityObject(created.todo, "successor Todo"); + nextTodos.push({...record, todo: record.text, ok: true, dry_run: input.dry_run}); + if (created.status === "planned") mutations.push({kind: "todo_upsert", todo: record}); + plannedTodos.set(String(record.todo_id), record); + } + const receiptFields = ["todo_id", "role", "task_class", "action_kind", "task_repository", + "continuation_policy", "required_capabilities", "claimed_by", "unblocks_todo_id", "target_key"]; + const writeback: JsonObject = {schema_version: "monitor_poll_todo_writeback_v0", dry_run: input.dry_run, + goal_id: input.goal_id, todo_id: monitor.todo_id, monitor_effect_id: input.operation_id, + target_key: transition.target_key || null, result_hash: observation.result_hash, + material_change: observation.material_change, material_change_generation: transition.material_change_generation, + 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}, + 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}; +} + +export async function executeCoordinationMonitorPoll(store: AuthorityStore, + raw: CoordinationMonitorPollInput): Promise { + let input: CoordinationMonitorPollInput; + 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}); + const previous = replay(await store.readReceipt(input.operation_id), input, hash, "replayed"); + if (previous) return previous; + const head = await store.loadAuthority(); + if (head.status !== "loaded") return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...head}; + let plan: ReturnType; + try { plan = planWriteback(input, head.head); } + catch (error) { return failure("monitor_poll_rejected", error instanceof Error ? error.message : String(error)); } + if (input.dry_run) return {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, + status: "planned", changed: true, writeback: plan.writeback, provider_revision: head.provider_revision}; + const commit = prepareCoordinationProjectionCommit({goal_id: input.goal_id, operation_id: input.operation_id, + expected_provider_revision: head.provider_revision, projection: head.head, mutations: plan.mutations}); + commit.receipts = [{schema_version: RECEIPT_SCHEMA, goal_id: input.goal_id, operation_id: input.operation_id, + request_sha256: hash, writeback: plan.writeback}]; + const committed = await store.commitAuthority(commit); + const result = replay(await store.readReceipt(input.operation_id), input, hash, + committed.status === "applied" ? "applied" : "recovered"); + return result ?? (committed.status === "applied" + ? failure("coordination_commit_readback_mismatch", "applied Monitor transaction lacks its durable receipt") + : {schema_version: COORDINATION_MONITOR_POLL_RESULT_SCHEMA, ...committed, changed: false}); +} diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 7af4fba93c..aafdefdd13 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -48,7 +48,7 @@ import { evaluateQuotaVoidCommit } from "./quota/void_commit.ts"; import { readQuotaSettlement } from "./quota/settlement_readback.ts"; import { evaluateTurnEnvelope } from "./quota/turn_envelope.ts"; import { evaluateQuotaMonitorPollCommit } from "./quota/monitor_poll_commit.ts"; -import { planMonitorSuccessor } from "./scheduler/monitor_successor.ts"; +import { planMonitorSuccessor, selectMonitorTodoRequest } from "./scheduler/monitor_successor.ts"; import { evaluateDeliveryWorkspace } from "./agents/delivery_workspace.ts"; import { interpretTurnJournal, @@ -129,6 +129,7 @@ import { createLocalCoordinationTodo, editLocalCoordinationTodo, updateLocalCoordinationTodo, + pollLocalCoordinationMonitor, mutateLocalCoordinationAuthority, listLocalCoordinationTodos, promoteLocalCoordinationAuthority, @@ -455,6 +456,7 @@ export function createEffectRuntimeHandlers( ["coordination.local_authority.todo_claim", claimLocalCoordinationTodo], ["coordination.local_authority.todo_create", createLocalCoordinationTodo], ["coordination.local_authority.todo_update", updateLocalCoordinationTodo], + ["coordination.local_authority.monitor_poll", pollLocalCoordinationMonitor], ["coordination.local_authority.todo_terminal", terminalLifecycleLocalCoordinationTodo], ["coordination.local_authority.todo_archive", archiveLocalCoordinationTodos], ["coordination.local_authority.todo_archive_ack", acknowledgeLocalCoordinationTodoArchive], @@ -473,6 +475,7 @@ export function createEffectRuntimeHandlers( ["task_lease.write_scopes.overlap", evaluateTaskLeaseWriteScopesOverlap], ["quota.monitor_poll.commit", evaluateQuotaMonitorPollCommit], ["scheduler.monitor_successor.plan", planMonitorSuccessor], + ["scheduler.monitor_target.select", selectMonitorTodoRequest], ["coordination.local_authority_shadow.record", recordLocalAuthorityShadow], ["coordination.runtime_shadow.commit_entry", commitLocalAuthorityShadowEntry], ["coordination.runtime_shadow.outbox_read", readLocalAuthorityShadow], diff --git a/loopx/control_plane/quota/monitor_poll.py b/loopx/control_plane/quota/monitor_poll.py index cdd43afee9..4bf5d4b5d4 100644 --- a/loopx/control_plane/quota/monitor_poll.py +++ b/loopx/control_plane/quota/monitor_poll.py @@ -641,11 +641,11 @@ def record_quota_monitor_poll_for_decision( else f"quota-monitor-poll:{goal_id}:{uuid.uuid4().hex}" ) if execute and (safe_todo_id or safe_target_key): - from ..coordination.legacy_writer_fence import ( - require_legacy_coordination_write_allowed, + from ..scheduler.provider_monitor_poll import ( + require_monitor_poll_source_available, ) - require_legacy_coordination_write_allowed( + require_monitor_poll_source_available( runtime_root=runtime_root, goal_id=goal_id, ) diff --git a/loopx/control_plane/quota/monitor_poll_commit.ts b/loopx/control_plane/quota/monitor_poll_commit.ts index e5c05b81ad..ec05ef6182 100644 --- a/loopx/control_plane/quota/monitor_poll_commit.ts +++ b/loopx/control_plane/quota/monitor_poll_commit.ts @@ -633,9 +633,26 @@ function compactProviderWriteback(receipt: JsonObject): JsonObject { ]) { compact[field] = receipt[field] ?? null; } + Object.assign(compact, monitorProjectionDelivery(receipt)); return compact; } +/** Display acknowledgement is diagnostic, never evidence for business commit. + * Omitted on the legacy path to retain its exact v0 response shape. */ +function monitorProjectionDelivery(receipt: JsonObject): JsonObject { + if (receipt.projection_delivery == null) return {}; + const status = optionalString(receipt.projection_delivery, "projection_delivery"); + if (!["delivered", "pending", "not_required"].includes(String(status))) { + throw new EffectRuntimeRequestError("invalid Monitor projection delivery status"); + } + const outbox = requiredObject(receipt.projection_outbox, "projection_outbox"); + const diagnostic: JsonObject = {}; + for (const key of ["schema_version", "status", "reason_code", "retryable", "recommended_action", "retry_business_mutation"]) { + if (outbox[key] != null) diagnostic[key] = outbox[key]; + } + return {projection_delivery: status, projection_outbox: diagnostic}; +} + function buildRecord(request: MonitorRequest): JsonObject { const allowed = admission(request); const material = request.observation.material_change; @@ -1114,6 +1131,7 @@ function validatedProviderReceipt( todo_update: requiredObject(receipt.todo_update, "provider_receipt.todo_update"), next_todos: nextTodos, successor_receipts: successors, + ...monitorProjectionDelivery(receipt), }; } diff --git a/loopx/control_plane/scheduler/monitor_poll_writeback.py b/loopx/control_plane/scheduler/monitor_poll_writeback.py index fc634cbd05..2eae209b89 100644 --- a/loopx/control_plane/scheduler/monitor_poll_writeback.py +++ b/loopx/control_plane/scheduler/monitor_poll_writeback.py @@ -5,13 +5,11 @@ from ..todos.contract import ( TODO_TASK_CLASS_ADVANCEMENT, - TODO_TASK_CLASS_MONITOR, TODO_TASK_CLASS_USER_GATE, normalize_todo_id, ) from ..effect_runtime import EffectRuntimeRejected, effect_runtime_result from ..todos.monitor_metadata import MonitorPollObservation -from .monitor_todo import monitor_todo_task_class def resolve_monitor_todo_item( @@ -35,43 +33,16 @@ def resolve_monitor_todo_item( runtime_root_arg=str(runtime_root) if runtime_root is not None else None, ) items = payload.get("todos") if isinstance(payload.get("todos"), list) else [] - if normalized_todo_id: - matches = [ - item - for item in items - if isinstance(item, dict) - and normalize_todo_id(item.get("todo_id")) == normalized_todo_id - ] - if not matches: - raise ValueError(f"monitor todo_id {normalized_todo_id!r} was not found") - if len(matches) > 1: - raise ValueError(f"monitor todo_id {normalized_todo_id!r} matched multiple todos") - item = matches[0] - item_target_key = str(item.get("target_key") or "").strip() - if safe_target_key and item_target_key and safe_target_key != item_target_key: - raise ValueError( - f"monitor todo_id {normalized_todo_id!r} resolves target_key " - f"{item_target_key!r}, not {safe_target_key!r}" - ) - if monitor_todo_task_class(item) != TODO_TASK_CLASS_MONITOR: - raise ValueError("monitor-poll todo writeback target must be task_class=continuous_monitor") - return item - - matches: list[dict[str, Any]] = [] - for item in items: - if not isinstance(item, dict): - continue - if safe_target_key and str(item.get("target_key") or "").strip() == safe_target_key: - matches.append(item) - if not matches: - target = normalized_todo_id or safe_target_key - raise ValueError(f"monitor todo target {target!r} was not found") - if len(matches) > 1: - raise ValueError(f"monitor target_key {safe_target_key!r} matched multiple todos; pass --todo-id") - item = matches[0] - if monitor_todo_task_class(item) != TODO_TASK_CLASS_MONITOR: - raise ValueError("monitor-poll todo writeback target must be task_class=continuous_monitor") - return item + try: + result = effect_runtime_result("scheduler.monitor_target.select", { + "schema_version": "loopx_monitor_target_request_v0", "items": items, + "todo_id": normalized_todo_id, "target_key": safe_target_key or None, + }) + except EffectRuntimeRejected as exc: + raise ValueError(str(exc)) from exc + if not isinstance(result, dict) or result.get("schema_version") != "loopx_monitor_target_result_v0": + raise TypeError("TypeScript Monitor target selection shape mismatch") + return result["todo"] def write_monitor_poll_todo_state( @@ -117,6 +88,24 @@ def write_monitor_poll_todo_state( if not todo_id and not target_key: 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, + observation={"todo_id": todo_id, "target_key": target_key, + "result_hash": result_hash, "material_change": material_change, + "generated_at": generated_at, "cadence": cadence, + "next_due_at": next_due_at, "reason_summary": reason_summary}, + intent={"next_agent_todo": next_agent_todo, "next_action_kind": next_action_kind, + "next_task_repository": next_task_repository, + "next_required_capabilities": next_required_capabilities or [], + "next_continuation_policy": next_continuation_policy, + "next_target_key": next_target_key, "next_claimed_by": next_claimed_by, + "next_user_todo": next_user_todo, "next_user_task_class": next_user_task_class}, + ) + if canonical is not None: + return canonical if execute: require_legacy_coordination_write_allowed( runtime_root=runtime_root, diff --git a/loopx/control_plane/scheduler/monitor_successor.ts b/loopx/control_plane/scheduler/monitor_successor.ts index 292ee29587..7431728c0b 100644 --- a/loopx/control_plane/scheduler/monitor_successor.ts +++ b/loopx/control_plane/scheduler/monitor_successor.ts @@ -9,6 +9,39 @@ import { compactPythonWhitespace, normalizeTodoAgent, stripPythonWhitespace } fr export const MONITOR_SUCCESSOR_REQUEST_SCHEMA = "loopx_monitor_successor_plan_request_v0"; export const MONITOR_SUCCESSOR_RESULT_SCHEMA = "loopx_monitor_successor_plan_result_v0"; +/** Select from the caller's complete snapshot, never from a compact lane. */ +export function selectMonitorTodo(items: readonly JsonObject[], todoId: string | null, + targetKey: string | null): JsonObject { + if (!todoId && !targetKey) throw new EffectRuntimeRequestError("monitor todo writeback requires --todo-id or --target-key"); + // A completed historical watch must not shadow its active replacement. + // Explicit IDs still resolve first so an inactive target gets a rejection. + const matches = items.filter(item => todoId ? item.todo_id === todoId : + item.target_key === targetKey && item.status !== "done" && item.archive_state !== "archive"); + if (matches.length !== 1) throw new EffectRuntimeRequestError( + matches.length ? "monitor target matched multiple todos; pass --todo-id" : "monitor todo target was not found"); + const item = matches[0]!; + if (targetKey && item.target_key && item.target_key !== targetKey) { + throw new EffectRuntimeRequestError(`monitor todo target_key resolves to '${item.target_key}', not '${targetKey}'`); + } + if (item.role !== "agent" || item.task_class !== "continuous_monitor") { + throw new EffectRuntimeRequestError("monitor-poll todo writeback target must be task_class=continuous_monitor"); + } + if (item.status === "done" || item.archive_state === "archive") { + throw new EffectRuntimeRequestError("monitor-poll requires an active, unfinished Monitor"); + } + return item; +} + +export function selectMonitorTodoRequest(value: unknown): JsonObject { + const request = requireJsonObject(value, "monitor target selection"); + if (request.schema_version !== "loopx_monitor_target_request_v0" || !Array.isArray(request.items)) { + throw new EffectRuntimeRequestError("monitor target selection schema mismatch"); + } + return {schema_version: "loopx_monitor_target_result_v0", todo: selectMonitorTodo( + request.items.map(item => requireJsonObject(item, "monitor target item")), + text(request.todo_id, "todo_id"), text(request.target_key, "target_key"))}; +} + function text(value: unknown, field: string): string | null { const raw = optionalNonEmptyString(value, field); return raw === null ? null : stripPythonWhitespace(raw) || null; diff --git a/loopx/control_plane/scheduler/provider_monitor_poll.py b/loopx/control_plane/scheduler/provider_monitor_poll.py new file mode 100644 index 0000000000..aa7b0db832 --- /dev/null +++ b/loopx/control_plane/scheduler/provider_monitor_poll.py @@ -0,0 +1,56 @@ +"""Provider routing and projection delivery, not a second Monitor planner.""" +from __future__ import annotations + +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from ...agent_registry import registered_agent_ids_from_registry +from ..coordination.local_authority import ( + LocalCoordinationAuthorityUnavailable, + local_authority_is_promoted, + read_canonical_todos_if_promoted, +) +from ..effect_runtime import effect_runtime_result +from ..todos.provider_projection import settle_canonical_todo_projection + + +def require_monitor_poll_source_available(*, runtime_root: Path, goal_id: str) -> None: + """Fail closed on unavailable promoted authority before unrelated quota work.""" + read_canonical_todos_if_promoted(runtime_root=runtime_root, goal_id=goal_id) + + +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], +) -> 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", + "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, + "registered_agents": registered_agent_ids_from_registry(registry_path, goal_id), + "dry_run": not execute, "observation": observation, "intent": intent, + }) + if (not isinstance(result, dict) + or result.get("status") not in {"applied", "replayed", "recovered", "planned"} + or result.get("source_authority") != "file_v0" + or result.get("decision_read_from_provider") is not True + or result.get("legacy_fallback_used") is not False + or not isinstance(result.get("writeback"), dict)): + payload = result if isinstance(result, dict) else {} + raise LocalCoordinationAuthorityUnavailable( + str(payload.get("reason") or "canonical Monitor transaction unavailable"), + code=str(payload.get("reason_code") or "monitor_poll_unavailable"), payload=payload, + ) + # Business commit remains successful even when the separate renderer fails. + payload = {**result, "dry_run": not execute} + settled = settle_canonical_todo_projection(payload=payload, registry_path=registry_path, + runtime_root=runtime_root, goal_id=goal_id) + return {**result["writeback"], "source_authority": "file_v0", + "provider_revision": result.get("provider_revision"), + "projection_delivery": settled.get("projection_delivery"), + "projection_outbox": settled.get("projection_outbox")} diff --git a/loopx/control_plane/testing/model_behavior_packet_safety.py b/loopx/control_plane/testing/model_behavior_packet_safety.py index bc44f225b0..9a2a9d7877 100644 --- a/loopx/control_plane/testing/model_behavior_packet_safety.py +++ b/loopx/control_plane/testing/model_behavior_packet_safety.py @@ -43,11 +43,18 @@ def _decode_facts(chunks: list[str]) -> dict[str, Any]: if sum(map(len, chunks)) > _MAX_ENCODED_CHARS: raise ValueError("encoded boundary") encoded = "".join(chunks) - if not encoded or not re.fullmatch(r"[A-Za-z0-9_-]+", encoded): + if not encoded or not re.fullmatch(r"[A-Za-z0-9_+/-]+", encoded): raise ValueError("encoded alphabet") compressed = base64.b64decode(encoded + "=" * (-len(encoded) % 4), altchars=b"-_", validate=True) - if base64.urlsafe_b64encode(compressed).decode().rstrip("=") != encoded: - raise ValueError("non-canonical base64url") + # The shipped scheduler emits standard Base64 to avoid secret-shaped + # URL-safe chunks. Native CLI also accepts the older URL-safe form. + # Validate either canonical alphabet, not mixed alphabets or pad bits; + # the decoded envelope still receives the full confidentiality scan. + if encoded not in { + base64.b64encode(compressed).decode().rstrip("="), + base64.urlsafe_b64encode(compressed).decode().rstrip("="), + }: + raise ValueError("non-canonical base64") inflater = zlib.decompressobj() raw = inflater.decompress(compressed, _MAX_INFLATED_BYTES + 1) if len(raw) > _MAX_INFLATED_BYTES or not inflater.eof or inflater.unused_data or inflater.unconsumed_tail: diff --git a/tests/control_plane/test_model_behavior_qualification.py b/tests/control_plane/test_model_behavior_qualification.py index ad50a28461..9e4e1d41eb 100644 --- a/tests/control_plane/test_model_behavior_qualification.py +++ b/tests/control_plane/test_model_behavior_qualification.py @@ -176,12 +176,14 @@ def _scheduler_wire(operation: str = "ack") -> bytes: def _scheduler_transport_packet( arm: str, operation: str, *, inline: bool = False, compressed: bytes | None = None, + alphabet: str = "url", ) -> tuple[dict[str, Any], list[str]]: if compressed is None: # Frozen public synthetic wire bytes keep the collision independent of # the platform's zlib encoder. Hex storage is not a secret-shaped value. compressed = _scheduler_wire(operation) - encoded = base64.urlsafe_b64encode(compressed).decode().rstrip("=") + encoder = base64.b64encode if alphabet == "standard" else base64.urlsafe_b64encode + encoded = encoder(compressed).decode().rstrip("=") command = "scheduler-ack-current" if operation == "ack" else "scheduler-fail-current" args = ["quota", command, "--goal-id", "goal-native-followup", "--agent-id", "agent-native-followup"] for offset in range(0, len(encoded), 384): @@ -203,11 +205,13 @@ def _scheduler_transport_packet( ("full_packet", "ack"), ("full_packet", "host_failure"), ("candidate_packet", "ack"), ]) @pytest.mark.parametrize("inline", [False, True]) +@pytest.mark.parametrize("alphabet", ["standard", "url"]) def test_actor_request_scans_decoded_scheduler_facts_without_changing_wire( - arm: str, operation: str, inline: bool, + arm: str, operation: str, inline: bool, alphabet: str, ) -> None: - packet, args = _scheduler_transport_packet(arm, operation, inline=inline) - assert any(SECRET_LIKE_SURFACE_PATTERN.search(arg) for arg in args) + packet, args = _scheduler_transport_packet(arm, operation, inline=inline, alphabet=alphabet) + if alphabet == "url": + assert any(SECRET_LIKE_SURFACE_PATTERN.search(arg) for arg in args) before = json.dumps(packet, sort_keys=True) request = build_model_behavior_actor_request(packet, qualification_id="public-wire-collision", arm=arm) @@ -221,8 +225,9 @@ def test_actor_request_scans_decoded_scheduler_facts_without_changing_wire( ("full_packet", "ack"), ("full_packet", "host_failure"), ("candidate_packet", "ack"), ]) @pytest.mark.parametrize("location", ["before", "host_facts", "extension"]) +@pytest.mark.parametrize("alphabet", ["standard", "url"]) def test_actor_rejects_private_material_inside_encoded_scheduler_facts( - arm: str, operation: str, location: str, + arm: str, operation: str, location: str, alphabet: str, ) -> None: payload = json.loads(zlib.decompress(_scheduler_wire(operation))) if location == "before": @@ -231,7 +236,7 @@ def test_actor_rejects_private_material_inside_encoded_scheduler_facts( payload[location]["note"] = "/" + "Users/example/private.txt" else: payload[location] = [{"nested": ["token" + "=abcdefghijklmnop"]}] - packet, _ = _scheduler_transport_packet(arm, operation, inline=True, + packet, _ = _scheduler_transport_packet(arm, operation, inline=True, alphabet=alphabet, compressed=zlib.compress(json.dumps(payload).encode())) before = json.dumps(packet, sort_keys=True) with pytest.raises(ValueError, match="credential-shaped field|local absolute path|credential-like value"): @@ -243,7 +248,8 @@ def test_actor_rejects_private_material_inside_encoded_scheduler_facts( "bad_json", "utf8", "root_array", "duplicate_key", "nonfinite", "hint_schema", "facts_schema", "before_shape", "current_hint_shape", "truncated", "trailing", "concatenated", "bomb", "encoded_limit", ]) -def test_actor_rejects_malformed_or_unbounded_scheduler_wire(case: str) -> None: +@pytest.mark.parametrize("alphabet", ["standard", "url"]) +def test_actor_rejects_malformed_or_unbounded_scheduler_wire(case: str, alphabet: str) -> None: payload = json.loads(zlib.decompress(_scheduler_wire())) if case == "hint_schema": payload["schema_version"] = "future_hint" @@ -275,13 +281,13 @@ def test_actor_rejects_malformed_or_unbounded_scheduler_wire(case: str) -> None: compressed += zlib.compress(b"{}") elif case == "encoded_limit": compressed = b"x" * 3_073 - packet, _ = _scheduler_transport_packet("full_packet", "ack", compressed=compressed) + packet, _ = _scheduler_transport_packet("full_packet", "ack", compressed=compressed, alphabet=alphabet) with pytest.raises(ValueError, match="scheduler host facts"): build_model_behavior_actor_request(packet, qualification_id="invalid-wire", arm="full_packet") @pytest.mark.parametrize("case", [ - "hint_schema", "scheduler_schema", "command", "missing", "empty_inline", "alphabet", "padding", "pad_bits", + "hint_schema", "scheduler_schema", "command", "missing", "empty_inline", "alphabet", "mixed_alphabet", "padding", "pad_bits", ]) def test_actor_rejects_unrecognized_or_malformed_scheduler_arguments(case: str) -> None: packet, args = _scheduler_transport_packet("full_packet", "ack") @@ -296,7 +302,11 @@ def test_actor_rejects_unrecognized_or_malformed_scheduler_arguments(case: str) elif case == "empty_inline": args.append(_HOST_FACTS_FLAG + "=") elif case == "alphabet": - args[7] = "+invalid" + args[7] = "*invalid" + elif case == "mixed_alphabet": + encoded = "".join(args[7::2]).replace("-", "+", 1) + assert "+" in encoded and "-" in encoded + args[6:] = [_HOST_FACTS_FLAG, encoded] elif case == "padding": args[-1] += "=" else: diff --git a/tests/control_plane/test_native_monitor_poll.py b/tests/control_plane/test_native_monitor_poll.py new file mode 100644 index 0000000000..f6f0ac4629 --- /dev/null +++ b/tests/control_plane/test_native_monitor_poll.py @@ -0,0 +1,114 @@ +"""Actual CLI -> quota -> canonical transaction -> projection -> quota receipt.""" +from __future__ import annotations + +import pytest +import hashlib +import json + +from canonical_authority_fixture import initialize_canonical_authority +from test_monitor_followthrough_contract import _write_fixture, _add_monitor, GOAL_ID, AGENT_ID +from loopx.control_plane.coordination.runtime_shadow import build_todo_runtime_shadow_projection +from loopx.control_plane.coordination.local_authority import read_canonical_todos_if_promoted +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 +from loopx.todos import list_goal_todos + + +def _canonical(tmp_path, native=False): + 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") + 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 + records = [{**{k: v for k, v in item.items() if k in TODO_DOMAIN_RECORD_FIELDS}, + "schema_version": "todo_domain_record_v0"} for item in projection["todos"]] + projection["todos"] = records + projection["todo_read_model"] = {"schema_version": "loopx_todo_domain_read_record_v0", + "todo_count": len(records), "records_sha256": hashlib.sha256(canonical_bytes(records)).hexdigest(), + "contract_fields": list(TODO_DOMAIN_RECORD_FIELDS)} + initialize_canonical_authority(runtime, GOAL_ID, projection, state_path=state) + return registry, runtime, state, monitor + + +@pytest.mark.parametrize("native", [False, True]) +def test_public_cli_native_monitor_settles_with_independent_successor(tmp_path, native): + registry, runtime, state, monitor = _canonical(tmp_path, native=native) + # No Markdown business source survives cutover. The existing renderer may + # recreate its Todo-only view, but neither preflight nor commit requires it. + state.unlink() + result = run_json_cli("quota", "monitor-poll", "--goal-id", GOAL_ID, + "--agent-id", AGENT_ID, "--runtime-profile", "generic_cli", + "--todo-id", monitor["todo_id"], "--result-hash", "revision-a", "--material-change", + "--next-agent-todo", "Validate the observed change", "--next-action-kind", "validate", + "--next-task-repository", "https://github.com/example/repo.git", "--execute", + registry_path=registry, runtime_root=runtime) + writeback = result["todo_writeback"] + assert writeback["material_change_generation"] == 1 + successor = writeback["next_todos"][0] + assert successor["task_class"] == "advancement_task" + assert successor["unblocks_todo_id"] == monitor["todo_id"] + assert successor["continuation_policy"] == "independent_handoff" + readback = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) + assert len(readback["todos"]) == 2 + assert state.exists() + + +def test_native_monitor_retry_and_projection_failure_do_not_repeat_business(tmp_path, monkeypatch): + registry, runtime, state, monitor = _canonical(tmp_path) + import loopx.control_plane.todos.provider_projection as delivery + + def unavailable(**kwargs): + raise OSError("synthetic renderer unavailable") + + monkeypatch.setattr(delivery, "project_current_canonical_todos", unavailable) + before = state.read_bytes() + args = dict(registry_path=registry, runtime_root=runtime, goal_id=GOAL_ID, execute=True, + todo_id=monitor["todo_id"], agent_id=AGENT_ID, monitor_effect_id="isolated-poll-a", + generated_at="2026-09-01T00:00:00Z", result_hash="revision-a", material_change=True, + next_agent_todo="Validate observed change", next_action_kind="validate") + first = write_monitor_poll_todo_state(**args) + repeated = write_monitor_poll_todo_state(**args) + assert first["projection_delivery"] == "pending" + assert repeated["provider_replayed"] is True + assert repeated["next_todos"] == first["next_todos"] + assert state.read_bytes() == before + assert len(read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID)["todos"]) == 2 + + +@pytest.mark.parametrize("role", ["agent", "user"]) +def test_native_monitor_later_same_evidence_cannot_create_more_work(tmp_path, role): + registry, runtime, _state, monitor = _canonical(tmp_path) + args = dict(registry_path=registry, runtime_root=runtime, goal_id=GOAL_ID, execute=True, + todo_id=monitor["todo_id"], agent_id=AGENT_ID, result_hash="revision-a", material_change=True) + write_monitor_poll_todo_state(**args, monitor_effect_id="generation-a", generated_at="2026-09-01T00:00:00Z") + before = read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) + followup = {"next_agent_todo": "Duplicate", "next_action_kind": "validate"} if role == "agent" else { + "next_user_todo": "Duplicate", "next_user_task_class": "user_action"} + with pytest.raises(RuntimeError, match="new material-change generation"): + write_monitor_poll_todo_state(**args, **followup, monitor_effect_id="generation-b", generated_at="2026-09-01T01:00:00Z") + assert read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID) == before + + +def test_cli_preserves_pending_display_diagnostic_after_business_settlement(tmp_path, monkeypatch, capsys): + from loopx.cli import main + import loopx.control_plane.todos.provider_projection as delivery + + registry, runtime, _state, monitor = _canonical(tmp_path, native=True) + + def unavailable(**kwargs): + raise OSError("synthetic renderer unavailable") + + monkeypatch.setattr(delivery, "project_current_canonical_todos", unavailable) + code = main(["--registry", str(registry), "--runtime-root", str(runtime), "--format", "json", + "quota", "monitor-poll", "--goal-id", GOAL_ID, "--agent-id", AGENT_ID, + "--runtime-profile", "generic_cli", "--todo-id", monitor["todo_id"], + "--result-hash", "revision-a", "--material-change", "--next-agent-todo", "Validate observation", + "--next-action-kind", "validate", "--execute"]) + assert code == 0 + result = json.loads(capsys.readouterr().out) + assert result["todo_writeback"]["projection_delivery"] == "pending" + assert result["todo_writeback"]["projection_outbox"]["retry_business_mutation"] is False + assert len(read_canonical_todos_if_promoted(runtime_root=runtime, goal_id=GOAL_ID)["todos"]) == 2 diff --git a/tests/control_plane/test_split_root_todo_writeback_fence.py b/tests/control_plane/test_split_root_todo_writeback_fence.py index 8c35b1daad..dd53e4546a 100644 --- a/tests/control_plane/test_split_root_todo_writeback_fence.py +++ b/tests/control_plane/test_split_root_todo_writeback_fence.py @@ -135,7 +135,7 @@ def _poll_kwargs( } -def test_monitor_poll_writeback_blocked_when_override_root_is_fenced( +def test_monitor_poll_writeback_rejects_unavailable_canonical_override( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, ) -> None: @@ -146,7 +146,7 @@ def test_monitor_poll_writeback_blocked_when_override_root_is_fenced( _fence_check_blocks(monkeypatch) state_before = state.read_text(encoding="utf-8") - with pytest.raises(LegacyCoordinationWriterFenced): + with pytest.raises(LocalCoordinationAuthorityUnavailable): write_monitor_poll_todo_state(**_poll_kwargs(registry, runtime_override)) assert LEGACY_POLL_HASH in state.read_text(encoding="utf-8") @@ -229,7 +229,7 @@ def native(_method: str, request: dict[str, Any]) -> dict[str, Any]: "agent_identity": {"agent_id": AGENT_ID}, } - with pytest.raises(LegacyCoordinationWriterFenced): + with pytest.raises(LocalCoordinationAuthorityUnavailable): monitor_poll.record_quota_monitor_poll_for_decision( before, {"runtime_root": str(runtime_override)}, @@ -246,7 +246,7 @@ def native(_method: str, request: dict[str, Any]) -> dict[str, Any]: assert OVERRIDE_POLL_HASH not in state.read_text(encoding="utf-8") -def test_quota_monitor_poll_cli_preserves_fence_rejection( +def test_quota_monitor_poll_cli_rejects_unavailable_canonical_before_collection( monkeypatch: pytest.MonkeyPatch, tmp_path: Path, capsys: pytest.CaptureFixture[str], @@ -291,17 +291,9 @@ def unexpected_collection(**_kwargs: Any) -> dict[str, Any]: assert exit_code == 1 payload = json.loads(capsys.readouterr().out) - assert payload["error_code"] == "legacy_coordination_writer_fenced" - assert payload["reason"] == ( - "legacy coordination writer is fenced; use the promoted canonical " - f"authority (file_v0) for goal {GOAL_ID}; fence unknown; " - "the primary record was not changed" - ) - assert payload["write_check"] == { - "status": "blocked", - "reason_code": "legacy_coordination_writer_fenced", - "authority_mode": "file_v0", - } + assert payload["error_code"] == "local_authority_todo_list_unavailable" + assert payload["source_authority"] == "file_v0" + assert payload["legacy_fallback_used"] is False assert state.read_bytes() == before diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index f50a956696..6d7fc02291 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -22,6 +22,7 @@ import { import { prepareCoordinationProjectionCommit } from "../../loopx/control_plane/coordination/coordination_projection.ts"; import { executeCoordinationTodoClaim } from "../../loopx/control_plane/coordination/todo_claim.ts"; import { executeCoordinationTodoCreate } from "../../loopx/control_plane/coordination/todo_create.ts"; +import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; import { executeCoordinationTodoUpdate } from "../../loopx/control_plane/coordination/todo_update.ts"; import { listLocalCoordinationTodos, LOCAL_COORDINATION_TODO_LIST_REQUEST_SCHEMA } from "../../loopx/control_plane/coordination/local_authority_runtime.ts"; @@ -257,6 +258,60 @@ 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) => { + const {store, contender} = await factory(t); + const goal = "goal-monitor"; + const fixture = productionScaleCoordinationFixture(goal); + const projection = structuredClone(fixture.projection); + 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)); + assert.ok(monitor, "complex fixture needs a lease-free Monitor"); + Object.assign(monitor, {target_key: "conformance-watch", cadence: "1h", material_change_generation: 4}); + for (const field of ["claimed_by", "bound_agent", "excluded_agents", "last_checked_at", "monitor_effect_id", "task_repository", "result_hash"]) delete monitor[field]; + if (native) for (const record of records) { + record.schema_version = TODO_DOMAIN_ITEM_SCHEMA; delete record.source_section; delete record.index; + } + projection.todo_read_model = {schema_version: native ? TODO_DOMAIN_READ_RECORD_SCHEMA : TODO_CANONICAL_READ_RECORD_SCHEMA, + todo_count: records.length, records_sha256: canonicalAuthoritySha256(records), + contract_fields: [...(native ? TODO_DOMAIN_RECORD_CONTRACT.fields : TODO_CANONICAL_READ_RECORD_FIELDS)]}; + assert.equal((await store.commitAuthority({operation_id: "monitor-seed", expected_provider_revision: null, + events: [], receipts: [], next_projection: projection})).status, "applied"); + const request = {goal_id: goal, operation_id: "monitor-effect", actor_agent_id: "agent-a", + registered_agents: ["agent-a", "agent-b"], dry_run: false, + observation: {todo_id: monitor.todo_id, generated_at: "2099-01-01T00:00:00Z", + result_hash: "new-evidence", material_change: true}, + intent: {next_agent_todo: "Inspect the newly observed change", next_action_kind: "validate"}}; + const before = await store.loadAuthority(); + const invalid = await executeCoordinationMonitorPoll(store, {...request, + intent: {...request.intent, next_claimed_by: "unregistered"}}); + assert.equal(invalid.status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + const [first, second] = await Promise.all([executeCoordinationMonitorPoll(store, request), + executeCoordinationMonitorPoll(contender, request)]); + assert.ok([first, second].some(result => result.status === "applied"), JSON.stringify([first, second])); + assert.equal((await executeCoordinationMonitorPoll(contender, request)).status, "replayed"); + const after = await store.loadAuthority(); + assert.equal(after.status, "loaded"); + if (after.status !== "loaded") return; + const next = after.head.todos as Record[]; + assert.equal(next.length, records.length + 1); + assert.equal(next.find(todo => todo.todo_id === monitor.todo_id)?.material_change_generation, 5); + assert.deepEqual(after.head.leases, projection.leases); + const unrelated = next.filter(todo => todo.todo_id !== monitor.todo_id && records.some(old => old.todo_id === todo.todo_id)); + assert.deepEqual(unrelated, records.filter(todo => todo.todo_id !== monitor.todo_id)); + // Lose the acknowledgement, not the durable write; a new client recovers + // both halves from the existing receipt without advancing generation again. + const lostAck = new Proxy(store, {get(target, property) { + if (property === "commitAuthority") return async (commit: AuthorityStoreCommit) => { + await target.commitAuthority(commit); return {status: "conflict", reason: "synthetic lost acknowledgement"}; + }; + const value = Reflect.get(target, property); return typeof value === "function" ? value.bind(target) : value; + }}); + const recovered = await executeCoordinationMonitorPoll(lostAck, {...request, operation_id: "monitor-next-effect", + observation: {...request.observation, generated_at: "2099-01-01T01:00:00Z", material_change: false}, intent: {}}); + assert.equal(recovered.status, "recovered", JSON.stringify(recovered)); + }); for (const native of [false, true]) test(`${providerName} conformance: governance reads one full Todo/lease snapshot (${native ? "native" : "legacy"})`, async (t) => { const {store} = await factory(t); const goal = "goal-governance"; diff --git a/tests/control_plane_ts/monitor_metadata.test.ts b/tests/control_plane_ts/monitor_metadata.test.ts index 9852ad458a..cdd42903a3 100644 --- a/tests/control_plane_ts/monitor_metadata.test.ts +++ b/tests/control_plane_ts/monitor_metadata.test.ts @@ -138,7 +138,7 @@ test("production-scale snapshot remains immutable while each monitor gets an iso const result = planMonitorMetadata(request({existing: todo, enforce_boundedness: false, observation: observation({target_key: null, cadence: "1h"}), })); - assert.equal(result.metadata.material_change_generation, "1"); + assert.equal(result.metadata.material_change_generation, String(Number(todo.material_change_generation ?? 0) + 1)); assert.equal(result.metadata.consecutive_no_change, "0"); assert.equal(result.metadata.claimed_by, undefined); assert.equal(result.metadata.status, undefined); diff --git a/tests/control_plane_ts/production_scale_coordination_fixture.ts b/tests/control_plane_ts/production_scale_coordination_fixture.ts index 56527b8b5e..c11af5b937 100644 --- a/tests/control_plane_ts/production_scale_coordination_fixture.ts +++ b/tests/control_plane_ts/production_scale_coordination_fixture.ts @@ -89,6 +89,15 @@ function todoRecords( if (role === "agent" && status !== "done" && status !== "deferred") { record.claimed_by = index % 2 === 0 ? "agent-a" : "agent-b"; } + if (record.task_class === "continuous_monitor") { + // Durable mixed-source observation shapes: bounded and watch-only, + // untouched and previously changed, with cadence and retained generation. + Object.assign(record, {target_key: `synthetic-watch-${index}`, cadence: "1h", + last_checked_at: observedAt(index), next_due_at: "2025-02-01T00:00:00Z", + result_hash: `synthetic-result-${index}`, material_change_generation: index % 3, + consecutive_no_change: String(index % 5), material_change: String(index % 3 === 0), + ...(index % 8 === 0 ? {watch_only: "true"} : {max_no_change_before_replan: "5"})}); + } if (role === "agent" && status === "done" && index < 3) { record.successor_todo_ids = [todoId("agent", envelope.completion_target_index + index)]; record.completion_continuation = "successor"; diff --git a/tests/control_plane_ts/todo_monitor_poll.test.ts b/tests/control_plane_ts/todo_monitor_poll.test.ts new file mode 100644 index 0000000000..93a2e18bcf --- /dev/null +++ b/tests/control_plane_ts/todo_monitor_poll.test.ts @@ -0,0 +1,138 @@ +import assert from "node:assert/strict"; +import {mkdtemp} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import test from "node:test"; +import type {JsonObject} from "../../loopx/control_plane/effect_program.ts"; +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {TODO_DOMAIN_ITEM_SCHEMA, TODO_DOMAIN_READ_RECORD_SCHEMA, TODO_DOMAIN_RECORD_CONTRACT} from "../../loopx/control_plane/coordination/coordination_state_contract.ts"; +import {executeCoordinationMonitorPoll} from "../../loopx/control_plane/coordination/todo_monitor_poll.ts"; +import {selectMonitorTodo} from "../../loopx/control_plane/scheduler/monitor_successor.ts"; + +async function seeded(overrides: JsonObject = {}, extras: JsonObject[] = []) { + const store = new FileAuthorityStore(await mkdtemp(join(tmpdir(), "monitor-transaction-")), "goal-a"); + const records = [{schema_version: TODO_DOMAIN_ITEM_SCHEMA, todo_id: "todo_monitor", role: "agent", + status: "open", done: false, text: "Observe external progress", archive_state: "active", + task_class: "continuous_monitor", target_key: "watch-a", cadence: "1h", ...overrides}, ...extras]; + await store.commitAuthority({operation_id: "seed", expected_provider_revision: null, + events: [], receipts: [], next_projection: {goal_id: "goal-a", todos: records, leases: [], + todo_read_model: {schema_version: TODO_DOMAIN_READ_RECORD_SCHEMA, todo_count: records.length, + records_sha256: canonicalAuthoritySha256(records), contract_fields: [...TODO_DOMAIN_RECORD_CONTRACT.fields]}}}); + const request = {goal_id: "goal-a", operation_id: "poll-a", actor_agent_id: "agent-a", + registered_agents: ["agent-a", "agent-b"], dry_run: false, + observation: {todo_id: "todo_monitor", target_key: "watch-a", generated_at: "2026-09-01T00:00:00Z", + result_hash: "result-a", material_change: true}, + intent: {next_agent_todo: "Deliver the newly observed result", next_action_kind: "implementation"}}; + return {store, request}; +} + +test("monitor and independent successor commit together, preview is inert, receipt replays exact intent", async () => { + const {store, request} = await seeded(); + const before = await store.loadAuthority(); + assert.equal((await executeCoordinationMonitorPoll(store, {...request, dry_run: true})).status, "planned"); + assert.deepEqual(await store.loadAuthority(), before); + const result = await executeCoordinationMonitorPoll(store, request); + assert.equal(result.status, "applied", JSON.stringify(result)); + const after = await store.loadAuthority(); + assert.equal(after.status, "loaded"); + if (after.status !== "loaded") return; + const todos = after.head.todos as JsonObject[]; + assert.equal(todos.length, 2); + const monitor = todos.find(todo => todo.todo_id === "todo_monitor")!; + assert.equal(monitor.material_change_generation, 1); + assert.equal(monitor.status, "open"); + assert.equal(monitor.claimed_by, undefined); + const successor = todos.find(todo => todo.todo_id !== monitor.todo_id)!; + assert.equal(successor.task_class, "advancement_task"); + assert.equal(successor.unblocks_todo_id, monitor.todo_id); + assert.equal(successor.continuation_policy, "independent_handoff"); + assert.equal((await executeCoordinationMonitorPoll(store, request)).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), after); + assert.equal((await executeCoordinationMonitorPoll(store, {...request, + intent: {...request.intent, next_agent_todo: "Different intent"}})).reason_code, + "coordination_operation_identity_mismatch"); +}); + +test("invalid successor, inactive target, ownership and exclusion reject before any write", async () => { + for (const overrides of [{status: "done", done: true}, {archive_state: "archive"}, + {claimed_by: "agent-b"}, {excluded_agents: ["agent-a"]}, {bound_agent: "agent-b"}]) { + const {store, request} = await seeded(overrides); + const before = await store.loadAuthority(); + assert.equal((await executeCoordinationMonitorPoll(store, request)).status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal((await store.readReceipt(request.operation_id)).status, "missing"); + } + const {store, request} = await seeded(); + const before = await store.loadAuthority(); + for (const intent of [{next_agent_todo: "Missing routing"}, + {...request.intent, next_claimed_by: "unknown"}]) { + assert.equal((await executeCoordinationMonitorPoll(store, {...request, intent})).status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + } +}); + +test("unchanged observations reschedule without creating delivery work", async () => { + const {store, request} = await seeded(); + const result = await executeCoordinationMonitorPoll(store, {...request, intent: {}, + observation: {...request.observation, material_change: false}}); + assert.equal(result.status, "applied", JSON.stringify(result)); + const receipt = result.writeback as JsonObject; + assert.equal(receipt.consecutive_no_change, 1); + assert.equal(receipt.material_change_generation, 0); + assert.equal(receipt.next_due_at, "2026-09-01T01:00:00Z"); + assert.deepEqual(receipt.next_todos, []); +}); + +test("User successors use public actor-bound authoring scope, never global gates", async () => { + for (const taskClass of ["user_action", "user_gate"]) { + const {store, request} = await seeded(); + const result = await executeCoordinationMonitorPoll(store, {...request, + intent: {next_user_todo: "Review observed change", next_user_task_class: taskClass}}); + assert.equal(result.status, "applied", JSON.stringify(result)); + const todo = ((result.writeback as JsonObject).next_todos as JsonObject[])[0]!; + assert.equal(todo.bound_agent, "agent-a"); + assert.equal(todo.global_gate, undefined); + assert.equal(todo.blocks_agent, taskClass === "user_gate" ? "agent-a" : undefined); + } +}); + +test("same hash cannot author duplicate work; distinct later evidence advances generation", async () => { + const {store, request} = await seeded(); + assert.equal((await executeCoordinationMonitorPoll(store, request)).status, "applied"); + const before = await store.loadAuthority(); + const later = {...request, operation_id: "poll-b", observation: {...request.observation, generated_at: "2026-09-01T01:00:00Z"}}; + assert.equal((await executeCoordinationMonitorPoll(store, later)).status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + const changed = await executeCoordinationMonitorPoll(store, {...later, + observation: {...later.observation, result_hash: "result-b"}, + intent: {...request.intent, next_agent_todo: "Deliver a different change"}}); + assert.equal(changed.status, "applied", JSON.stringify(changed)); + assert.equal((changed.writeback as JsonObject).material_change_generation, 2); +}); + +test("retained lease and stale observation cannot mutate either half", async () => { + const {store, request} = await seeded(); + assert.equal((await executeCoordinationMonitorPoll(store, request)).status, "applied"); + const before = await store.loadAuthority(); + assert.equal(before.status, "loaded"); + if (before.status !== "loaded") return; + assert.equal((await executeCoordinationMonitorPoll(store, {...request, operation_id: "stale", + observation: {...request.observation, generated_at: "2025-01-01T00:00:00Z"}})).status, "failed"); + assert.deepEqual(await store.loadAuthority(), before); + await store.commitAuthority({operation_id: "lease", expected_provider_revision: before.provider_revision, + events: [], receipts: [], next_projection: {...before.head, leases: [{todo_id: "todo_monitor", + owner: "agent-a", status: "active", expires_at: "2099-01-01T00:00:00Z", version: 1, idempotency_key: "execution"}]}}); + const leased = await store.loadAuthority(); + assert.equal((await executeCoordinationMonitorPoll(store, {...request, operation_id: "leased"})).status, "failed"); + assert.deepEqual(await store.loadAuthority(), leased); +}); + +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"}; + const history = {...active, todo_id: "todo_history", status: "done", archive_state: "archive"}; + assert.equal(selectMonitorTodo([history, active], null, "watch").todo_id, active.todo_id); + assert.throws(() => selectMonitorTodo([history, active], history.todo_id, "watch"), /unfinished/); + assert.throws(() => selectMonitorTodo([active, {...active, todo_id: "todo_other"}], null, "watch"), /multiple/); +});