diff --git a/docs/reference/canonical-lease-renew.md b/docs/reference/canonical-lease-renew.md new file mode 100644 index 0000000000..fd3bc89dee --- /dev/null +++ b/docs/reference/canonical-lease-renew.md @@ -0,0 +1,89 @@ +# Canonical lease renewal + +`task-lease renew` supports an already promoted local File or SQLite Goal. +Previously these calls reached the legacy writer fence and were rejected. +Unpromoted Goals keep the existing legacy lease behavior. Selecting a provider +alone does not promote a Goal or grant execution authority. + +Read the current canonical lease before choosing its CAS version: + +```bash +loopx --registry registry.json task-lease inspect \ + --goal-id example-goal --todo-id todo_work --format json +loopx --registry registry.json task-lease renew \ + --goal-id example-goal --todo-id todo_work --owner agent-a \ + --idempotency-key execution-a --expected-version 3 \ + --ttl-seconds 600 --format json +``` + +Use the version and execution key belonging to the current lease. A renewal +increments the lease version, preserves owner/key/epoch/write scopes, and sets +the expiry from the runtime's current clock plus the requested TTL. The default +TTL and its limits remain those of the existing lease command. Expired leases, +ineligible owners, terminal Todos, stale versions and invalid state are rejected. + +The new request schema is canonical-only. Missing/invalid fences or unavailable +providers never downgrade to a legacy lease file. Older runtimes reject the new +schema. Canonical mode, Todo and lease facts come from the provider; stale, +missing or malformed display frontmatter does not decide renewal eligibility. + +## Commit, retry and readback + +The lease projection, event and original receipt commit under one provider CAS. +The operation identity includes the Goal, Todo, owner, execution key and expected +lease version; changing TTL for the same identity is rejected. Keep a frozen +request when the response is lost or ambiguous. Recover the original receipt +before starting another renewal. + +`status=replayed` and `idempotent=true` return the original lease result, even +after a later renewal or expiry. They do **not** extend expiry again or prove +current authority. Use `task-lease inspect` for current state; do not use an old +replay's version as a newly granted execution proof. + +Canonical responses carry `source_authority`, `provider_revision` and `cursor`; +they do not invent a `lease_path`. Renewal writes neither a legacy lease record +nor a second shadow authority. A lease-only change does not require rewriting +Todo Markdown. The existing canonical journal retains the complete historical +projection and receipt. + +## Scope and rollback + +This adds renewal, not canonical transfer/release or the entire actor lifecycle. +It does not migrate legacy lifecycle receipts, change SQLite tables, switch +provider defaults, deploy a PostgreSQL service or authorize an active-Goal +migration. PostgreSQL local routing remains outside this command's supported +provider set. + +Rolling back code retains the canonical state and writer fence; an older client +or runtime may reject renewal again. Do not remove the fence as a rollback or +create a replacement legacy lease. Keep renewal availability and lease expiry in +the operational rollback plan. + +Correctness validation uses real File/SQLite and real CLI/process boundaries. +The SQLite validation profile uses Node 22.22.3 / SQLite 3.51.3. This does not +qualify the full [SQLite D2 capacity/continuity profile](sqlite-authority-store.md), +start its elapsed soak or grant promotion. + +## 中文操作说明 + +已 promoted 的本地 File/SQLite Goal 现在可使用现有 `task-lease renew`;此前该入口 +会被 legacy writer fence 拒绝。未 promoted 的 Goal 继续使用旧 lease 路径;provider +selector 本身不构成 promotion。先用上面的 `task-lease inspect` 读取当前 lease, +再使用该 lease 的 execution key 和 version 调用 renew。 + +合法续租只递增 version,并按 TS runtime 时钟加 TTL 更新 expiry;owner、key、epoch +和 write scopes 保持。到期、owner 不合法、Todo 已终止、version 过时或状态损坏时拒绝。 +canonical mode/Todo/lease 从 provider 读取,不受 stale/missing/损坏显示 frontmatter 控制。 +新请求 schema 绑定 canonical 路径,fence/provider 失效时不回退;旧 runtime 会拒绝该 schema。 + +lease、event 和原 receipt 在一次 provider CAS 中提交。相同 expected version 的同一 +identity 改 TTL 会拒绝;响应丢失时冻结原请求并恢复 receipt。`replayed`/`idempotent=true` +只返回历史结果,不再次延长 expiry,也不授予当前执行权。之后的操作需要重新 inspect +当前状态。返回值提供 provider/revision/cursor,不伪造 lease_path,不写第二份 legacy/shadow +状态,也不为 lease-only 修改制造 Todo Markdown 投影待办。 + +本切片仅增加 renew,不交付 canonical transfer/release、完整 actor lifecycle 或 legacy +receipt 迁移。它不改 SQLite 表结构、默认 provider 或活动 Goal,不部署 PostgreSQL。 +代码回退时保留 canonical state 和 fence,并把旧版本可能再次拒绝续租纳入运维计划; +不能删除 fence 或建立 legacy lease 作为回滚。真实 File/SQLite、CLI 和进程验证不等于 +完整 D2、自然时间 soak 或 promotion 已合格。 diff --git a/docs/reference/sqlite-authority-store.md b/docs/reference/sqlite-authority-store.md index 31f4012267..fa562b7cb3 100644 --- a/docs/reference/sqlite-authority-store.md +++ b/docs/reference/sqlite-authority-store.md @@ -205,6 +205,10 @@ node --no-warnings --experimental-sqlite --experimental-strip-types \ --output .local/sqlite-matched-64k.json ``` +For an already promoted Goal, [canonical lease renewal](canonical-lease-renew.md) +uses the selected provider's CAS and original receipt. It is a command-coverage +slice, not D2 qualification or a provider-default change. + The no-argument default intentionally replaces the former 4 KiB/100k run with a small `rehearsal`; full capacity now requires an explicit profile. The default `rehearsal` creates 100/1,000 commits and checks runner execution, diff --git a/loopx/control_plane/coordination/coordination_state_contract.generated.ts b/loopx/control_plane/coordination/coordination_state_contract.generated.ts index c8849a88a7..7db4d2724f 100644 --- a/loopx/control_plane/coordination/coordination_state_contract.generated.ts +++ b/loopx/control_plane/coordination/coordination_state_contract.generated.ts @@ -79,6 +79,7 @@ export const DELIVERY_WORKSPACE_SNAPSHOT_RESULT_SCHEMA = "loopx_delivery_workspa export const TASK_LEASE_ACQUIRE_REQUEST_SCHEMA = "loopx_task_lease_acquire_native_v0"; export const TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA = "loopx_task_lease_lifecycle_native_v0"; +export const TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA = "loopx_canonical_task_lease_renew_request_v0"; export const CAPABILITY_HOOK_REGISTRATION_SCHEMA = "loopx_capability_hook_registration_v0"; export const CAPABILITY_HOOK_INTERACTION_RESULT_SCHEMA = "loopx_interaction_projection_hook_result_v0"; @@ -293,7 +294,8 @@ export const COORDINATION_STATE_CONTRACT = deepFreeze({ }, "task_lease_protocol": { "acquire_request_schema": TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, - "lifecycle_request_schema": TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA + "lifecycle_request_schema": TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, + "canonical_renew_request_schema": TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA }, "capability_hook_protocol": { "registration_schema": CAPABILITY_HOOK_REGISTRATION_SCHEMA, diff --git a/loopx/control_plane/coordination/coordination_state_contract_generated.py b/loopx/control_plane/coordination/coordination_state_contract_generated.py index 522fff6a95..80798b47dd 100644 --- a/loopx/control_plane/coordination/coordination_state_contract_generated.py +++ b/loopx/control_plane/coordination/coordination_state_contract_generated.py @@ -164,7 +164,8 @@ def _freeze(value: Any) -> Any: 'request_schema': 'loopx_delivery_workspace_request_v0', 'result_schema': 'loopx_delivery_workspace_result_v0'}, 'task_lease_protocol': {'acquire_request_schema': 'loopx_task_lease_acquire_native_v0', - 'lifecycle_request_schema': 'loopx_task_lease_lifecycle_native_v0'}, + 'lifecycle_request_schema': 'loopx_task_lease_lifecycle_native_v0', + 'canonical_renew_request_schema': 'loopx_canonical_task_lease_renew_request_v0'}, 'capability_hook_protocol': {'registration_schema': 'loopx_capability_hook_registration_v0', 'interaction_result_schema': 'loopx_interaction_projection_hook_result_v0', 'turn_start_registration_schema': 'loopx_turn_start_capability_hook_registration_v1', @@ -261,6 +262,7 @@ def _freeze(value: Any) -> Any: TASK_LEASE_ACQUIRE_REQUEST_SCHEMA: Final[str] = 'loopx_task_lease_acquire_native_v0' TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA: Final[str] = 'loopx_task_lease_lifecycle_native_v0' +TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA: Final[str] = 'loopx_canonical_task_lease_renew_request_v0' CAPABILITY_HOOK_REGISTRATION_SCHEMA: Final[str] = 'loopx_capability_hook_registration_v0' CAPABILITY_HOOK_INTERACTION_RESULT_SCHEMA: Final[str] = 'loopx_interaction_projection_hook_result_v0' diff --git a/loopx/control_plane/coordination/coordination_state_contract_v0.json b/loopx/control_plane/coordination/coordination_state_contract_v0.json index f2515388fa..9fbf3bdb3d 100644 --- a/loopx/control_plane/coordination/coordination_state_contract_v0.json +++ b/loopx/control_plane/coordination/coordination_state_contract_v0.json @@ -174,7 +174,8 @@ }, "task_lease_protocol": { "acquire_request_schema": "loopx_task_lease_acquire_native_v0", - "lifecycle_request_schema": "loopx_task_lease_lifecycle_native_v0" + "lifecycle_request_schema": "loopx_task_lease_lifecycle_native_v0", + "canonical_renew_request_schema": "loopx_canonical_task_lease_renew_request_v0" }, "capability_hook_protocol": { "registration_schema": "loopx_capability_hook_registration_v0", diff --git a/loopx/control_plane/coordination/handoff_mode_runtime.ts b/loopx/control_plane/coordination/handoff_mode_runtime.ts index c7c82bc58a..fad914321a 100644 --- a/loopx/control_plane/coordination/handoff_mode_runtime.ts +++ b/loopx/control_plane/coordination/handoff_mode_runtime.ts @@ -3,7 +3,8 @@ import type {JsonObject} from "../effect_program.ts"; import {requireJsonObject} from "../runtime_decode.ts"; import {requireAuthorityStoreId} from "./authority_store_codec.ts"; import {openLocalAuthorityStore, localAuthorityOpenFailure} from "./local_authority_provider.ts"; -import {runtimeRoot, sourceAuthorityFor, withCanonicalWriter} from "./local_authority_runtime.ts"; +import {runtimeRoot, sourceAuthorityFor} from "./local_authority_runtime.ts"; +import {withCanonicalWriter} from "./local_authority_write.ts"; import {ShadowManagementError} from "./shadow_management.ts"; import {executeHandoffModeSet, HANDOFF_MODE_SET_SCHEMA} from "./handoff_mode_transaction.ts"; diff --git a/loopx/control_plane/coordination/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index 0e95a7a3d8..8d05408509 100644 --- a/loopx/control_plane/coordination/local_authority_runtime.ts +++ b/loopx/control_plane/coordination/local_authority_runtime.ts @@ -1,7 +1,7 @@ import {COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA} from "./todo_archive.ts"; import {executeTodoContinuation} from "./todo_continuation.ts"; -import { withFileMutationLock } from "../effect_runtime_io.ts"; -import { ShadowManagementError, requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath } from "./shadow_management.ts"; +import {withCanonicalWriter} from "./local_authority_write.ts"; +import { ShadowManagementError } from "./shadow_management.ts"; import { isAbsolute, join } from "node:path"; import type { JsonObject } from "../effect_program.ts"; @@ -102,14 +102,6 @@ export function sourceAuthorityFor(store: AuthorityStore): "sqlite_v0" | "file_v return store instanceof SqliteAuthorityStore ? "sqlite_v0" : "file_v0"; } -export async function withCanonicalWriter(root: string, goalId: string, dryRun: boolean, write: () => Promise): Promise { - if (dryRun) return await write(); - return await withFileMutationLock(shadowMaintenanceLockPath(root, goalId), async () => { - await requireShadowPrimaryWriteAllowed(root, goalId); - return await write(); - }); -} - /** Monitor observation and successors share the existing writer/fence lifetime. */ export async function pollLocalCoordinationMonitor(value: unknown, dependencies: LocalAuthorityRuntimeDependencies = {}): Promise { diff --git a/loopx/control_plane/coordination/local_authority_write.ts b/loopx/control_plane/coordination/local_authority_write.ts new file mode 100644 index 0000000000..24826035d8 --- /dev/null +++ b/loopx/control_plane/coordination/local_authority_write.ts @@ -0,0 +1,11 @@ +import {withFileMutationLock} from "../effect_runtime_io.ts"; +import {requireShadowPrimaryWriteAllowed, shadowMaintenanceLockPath} from "./shadow_management.ts"; + +/** Canonical command writers share the maintenance guard; provider CAS owns state. */ +export async function withCanonicalWriter(root: string, goalId: string, dryRun: boolean, write: () => Promise): Promise { + if (dryRun) return await write(); + return await withFileMutationLock(shadowMaintenanceLockPath(root, goalId), async () => { + await requireShadowPrimaryWriteAllowed(root, goalId); + return await write(); + }); +} diff --git a/loopx/control_plane/coordination/task_lease_renew.ts b/loopx/control_plane/coordination/task_lease_renew.ts new file mode 100644 index 0000000000..6b93b0da53 --- /dev/null +++ b/loopx/control_plane/coordination/task_lease_renew.ts @@ -0,0 +1,123 @@ +import type {JsonObject} from "../effect_program.ts"; +import type {AuthorityStore, AuthorityStoreCommit} from "./authority_store.ts"; +import {AuthorityStoreProtocolError, canonicalAuthorityObject} from "./authority_store_codec.ts"; +import {CoordinationCommandReceipt, commandReceiptResult} from "./command_receipt.ts"; +import {indexCoordinationProjection, prepareCoordinationProjectionCommit, validateCoordinationTodoReadModel} from "./coordination_projection.ts"; +import {decideTaskLeaseLifecycle} from "../work_items/task_lease_lifecycle_decision.ts"; +import {requireStringLiteral} from "../runtime_decode.ts"; +import {HANDOFF_MODES} from "./handoff_mode_policy.ts"; +import {leaseEpoch, leaseIsActive, leaseVersion, normalizeAgent, normalizeGoalId, + normalizeIdempotencyKey, normalizeOwner, normalizeTodoId, + normalizeTtl, TaskLeaseAcquireError, utcIsoformat, type LeaseRecord} from "../work_items/task_lease_acquire.ts"; +import {taskLeaseOperationIdentity, taskLeaseOperationRequestDigest} from "../work_items/task_lease_operation_identity.ts"; + +export const CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA = "loopx_canonical_task_lease_renew_result_v0"; +export const CANONICAL_TASK_LEASE_RENEW_RECEIPT_SCHEMA = "loopx_canonical_task_lease_renew_receipt_v0"; +export interface CanonicalTaskLeaseRenewInput { + goal_id: string; + todo_id: string; + owner: string; + idempotency_key: string; + expected_version: number | null; + ttl_seconds: number; + registered_agents: readonly string[]; + now: Date; +} +function failed(code: string, reason: string): JsonObject & {schema_version: typeof CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA} { + return {schema_version: CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA, status: "failed", changed: false, + reason_code: code, reason, failure_stage: "validation"}; +} +function leaseRecord(value: JsonObject, input: CanonicalTaskLeaseRenewInput): LeaseRecord { + if ((value.schema_version !== undefined && value.schema_version !== "task_lease_v0") || + (value.goal_id !== undefined && value.goal_id !== input.goal_id) || value.todo_id !== input.todo_id || + (value.status !== "active" && value.status !== "released")) { + throw new AuthorityStoreProtocolError("canonical lease identity or schema is invalid"); + } + if (typeof value.owner !== "string" || typeof value.idempotency_key !== "string" || + normalizeOwner(value.owner) !== value.owner || normalizeIdempotencyKey(value.idempotency_key) !== value.idempotency_key) { + throw new AuthorityStoreProtocolError("canonical lease owner and execution key must be normalized strings"); + } + leaseVersion(value); leaseEpoch(value); + if (value.write_scopes !== undefined && (!Array.isArray(value.write_scopes) || + value.write_scopes.some(scope => typeof scope !== "string"))) { + throw new AuthorityStoreProtocolError("canonical lease write_scopes must be strings"); + } + return value; +} + +/** Renew uses canonical facts and the existing lifecycle decision, never a lease file. */ +export async function executeCanonicalTaskLeaseRenew(store: AuthorityStore, raw: CanonicalTaskLeaseRenewInput, + beforeCommit?: (lease: JsonObject) => Promise): Promise { + let input: CanonicalTaskLeaseRenewInput; + try { + input = {...raw, goal_id: normalizeGoalId(raw.goal_id), todo_id: normalizeTodoId(raw.todo_id), + owner: normalizeOwner(raw.owner), idempotency_key: normalizeIdempotencyKey(raw.idempotency_key), + ttl_seconds: normalizeTtl(raw.ttl_seconds), registered_agents: raw.registered_agents.map(normalizeOwner)}; + if (input.expected_version === null) return failed("version_required", "task lease renew requires the current lease version"); + if (!Number.isSafeInteger(input.expected_version) || input.expected_version < 0) return failed("invalid_expected_version", "lease version must be a non-negative safe integer"); + if (!(input.now instanceof Date) || !Number.isFinite(input.now.valueOf())) return failed("invalid_clock", "renew requires a valid runtime clock"); + } catch (error) { return failed(error instanceof TaskLeaseAcquireError ? error.code : "invalid_canonical_renew", error instanceof Error ? error.message : "invalid canonical renewal"); } + const identityInput = {...input, operation: "renew", new_owner: null, new_idempotency_key: null}; + const operationId = `lease-renew:${taskLeaseOperationIdentity(identityInput)!}`; + const receipt = new CoordinationCommandReceipt({result_schema: CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA, + identity: {schema_version: CANONICAL_TASK_LEASE_RENEW_RECEIPT_SCHEMA, operation_id: operationId, + goal_id: input.goal_id, request_sha256: taskLeaseOperationRequestDigest(identityInput)!}, failure: failed, + decode(original) { + const payload = commandReceiptResult(original); + let lease: LeaseRecord; + try { + lease = leaseRecord(canonicalAuthorityObject(payload.fields.lease, "renew receipt lease"), input); + leaseIsActive(lease, new Date(0)); // Validate the timestamp without requiring current authority. + } catch (error) { + throw new AuthorityStoreProtocolError(error instanceof Error ? error.message : "invalid renewal receipt lease"); + } + if (payload.fields.renewed !== true || payload.changed !== true || lease.owner !== input.owner || + lease.status !== "active" || lease.idempotency_key !== input.idempotency_key || leaseVersion(lease) !== input.expected_version! + 1) { + throw new AuthorityStoreProtocolError("canonical renewal receipt does not match its intent"); + } + return {...payload, fields: {...payload.fields, operation_id: operationId}}; + }}); + let commit: AuthorityStoreCommit; + try { + const replay = await receipt.read(store); + if (replay !== null) return replay; + const head = await store.loadAuthority(); + if (head.status !== "loaded") return {schema_version: CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA, ...head, failure_stage: "validation", changed: false}; + const index = indexCoordinationProjection(head.head, input.goal_id); + validateCoordinationTodoReadModel(head.head, input.goal_id); + const todo = index.todos.get(input.todo_id), rawLease = index.leases.get(input.todo_id); + const lease = rawLease ? leaseRecord(rawLease, input) : null; + const excluded = todo?.excluded_agents ?? []; + if (!Array.isArray(excluded) || excluded.some(value => typeof value !== "string")) return failed("invalid_coordination_projection", "Todo exclusions must be strings"); + const mode = requireStringLiteral(head.head.handoff_mode ?? "legacy", HANDOFF_MODES, "canonical handoff_mode"); + const decision = decideTaskLeaseLifecycle({handoff_mode: mode, registered_agents: input.registered_agents, + todo: todo ? {todo_id: input.todo_id, status: String(todo.status), claimed_by: normalizeAgent(todo.claimed_by), excluded_agents: excluded as string[]} : null, + lease: lease ? {present: true, active: leaseIsActive(lease, input.now), status: String(lease.status), + owner: normalizeOwner(lease.owner), idempotency_key: normalizeIdempotencyKey(lease.idempotency_key), + version: leaseVersion(lease), lease_epoch: leaseEpoch(lease), write_scopes: (lease.write_scopes ?? []) as string[], acquire_ttl_seconds: null} : null, + command: {operation: "renew", owner: input.owner, idempotency_key: input.idempotency_key, + expected_version: input.expected_version, ttl_seconds: input.ttl_seconds, new_owner: null, new_idempotency_key: null}}); + if (decision.outcome !== "apply" || !lease || !decision.next_lease) return {...failed(decision.code, `canonical task lease renew rejected: ${decision.code}`), handoff_mode: mode}; + if (!Number.isSafeInteger(decision.next_lease.version)) return failed("lease_generation_exhausted", "lease version cannot advance safely"); + const next = {...lease, version: decision.next_lease.version, updated_at: utcIsoformat(input.now), + expires_at: utcIsoformat(new Date(input.now.valueOf() + input.ttl_seconds * 1000))}; + commit = prepareCoordinationProjectionCommit({goal_id: input.goal_id, operation_id: operationId, + expected_provider_revision: head.provider_revision, projection: head.head, mutations: [{kind: "lease_upsert", lease: next}]}); + commit.receipts = [{schema_version: CANONICAL_TASK_LEASE_RENEW_RECEIPT_SCHEMA, operation_id: operationId, + goal_id: input.goal_id, request_sha256: taskLeaseOperationRequestDigest(identityInput)!, + result: {changed: true, renewed: true, lease: next, handoff_mode: mode}}]; + await beforeCommit?.(next); + } catch (error) { + return failed(error instanceof TaskLeaseAcquireError ? error.code : "invalid_canonical_renew_state", + error instanceof Error ? error.message : "canonical renewal state could not be read"); + } + try { + const result = await receipt.commit(store, commit); + return result.status === "applied" || result.status === "replayed" || result.status === "recovered" + ? result : {...result, failure_stage: "durable_writeback"}; + } catch { + return {schema_version: CANONICAL_TASK_LEASE_RENEW_RESULT_SCHEMA, status: "ambiguous", changed: false, + reason_code: "canonical_renew_recovery_required", reason: "renewal response is uncertain; recover the same operation", + recovery: {operation_id: operationId, retry_with_same_operation_id: true}, failure_stage: "durable_writeback"}; + } +} diff --git a/loopx/control_plane/work_items/canonical_task_lease_renew.ts b/loopx/control_plane/work_items/canonical_task_lease_renew.ts new file mode 100644 index 0000000000..5e5953eab9 --- /dev/null +++ b/loopx/control_plane/work_items/canonical_task_lease_renew.ts @@ -0,0 +1,60 @@ +import type {JsonObject} from "../effect_program.ts"; +import {canonicalAuthoritySha256} from "../coordination/authority_store_codec.ts"; +import {FileAuthorityStore} from "../coordination/file_authority_store.ts"; +import {SqliteAuthorityStore} from "../coordination/sqlite_authority_store.ts"; +import {openLocalAuthorityStore, localAuthorityOpenFailure} from "../coordination/local_authority_provider.ts"; +import {loadLegacyCoordinationWriterFence} from "../coordination/legacy_writer_fence.ts"; +import {withCanonicalWriter} from "../coordination/local_authority_write.ts"; +import {ShadowManagementError} from "../coordination/shadow_management.ts"; +import {EffectRuntimeLockTimeoutError} from "../effect_runtime_errors.ts"; +import {executeCanonicalTaskLeaseRenew} from "../coordination/task_lease_renew.ts"; +import {revalidateAuthoritySources, TaskLeaseAcquireError, type AuthorityFacts} from "./task_lease_acquire.ts"; + +interface LocalRenewRequest { + runtime_root: string; goal_id: string; todo_id: string; + owner: string | null; idempotency_key: string | null; + expected_version: number | null; ttl_seconds: number | null; authority: AuthorityFacts | null; +} + +/** A canonical-only request requires a durable fence and can never fall back. */ +export async function renewCanonicalTaskLease(request: LocalRenewRequest, + dependencies: {now: () => Date; beforeWrite?: (lease: JsonObject) => void | Promise}): Promise { + const root = request.runtime_root, goalId = request.goal_id; + const initialFence = await loadLegacyCoordinationWriterFence(root, goalId); + const evidence: JsonObject = {source_authority: null, decision_read_from_provider: false, legacy_fallback_used: false}; + const rejected = (code: string, reason: string): JsonObject => ({status: "failed", changed: false, + reason_code: code, reason, failure_stage: "validation", ...evidence}); + if (initialFence.status === "missing") return rejected("canonical_renew_fence_missing", "canonical renewal fence is missing; legacy fallback is forbidden"); + if (initialFence.status === "failed") return rejected(initialFence.reason_code, initialFence.reason); + try { + return await withCanonicalWriter(root, goalId, false, async () => { + const verifyFence = async () => { + const current = await loadLegacyCoordinationWriterFence(root, goalId); + if (current.status !== "loaded" || canonicalAuthoritySha256(current.fence) !== canonicalAuthoritySha256(initialFence.fence)) { + throw new TaskLeaseAcquireError("canonical renewal writer fence changed; inspect authority before retrying", "canonical_renew_fence_changed", {goal_id: goalId}); + } + }; + await verifyFence(); + const store = await openLocalAuthorityStore(root, goalId); + // Current local opening supports exactly these built-in providers. Reject + // any future kind until its typed opening contract is adopted here. + if (store instanceof SqliteAuthorityStore) evidence.source_authority = "sqlite_v0"; + else if (store instanceof FileAuthorityStore) evidence.source_authority = "file_v0"; + else return rejected("canonical_renew_provider_unsupported", "canonical renewal provider is unsupported"); + if (!request.owner || !request.idempotency_key || !request.authority) return rejected("authority_required", "canonical renewal needs registered actor context"); + evidence.decision_read_from_provider = true; + const result = await executeCanonicalTaskLeaseRenew(store, {goal_id: goalId, todo_id: request.todo_id, + owner: request.owner, idempotency_key: request.idempotency_key, expected_version: request.expected_version, + ttl_seconds: request.ttl_seconds!, registered_agents: request.authority.registered_agents, now: dependencies.now()}, async lease => { + await dependencies.beforeWrite?.(lease); + await verifyFence(); + await revalidateAuthoritySources(request.authority!.source_receipts); + }); + return {...result, ...evidence}; + }); + } catch (error) { + if (error instanceof ShadowManagementError || error instanceof TaskLeaseAcquireError) return {...rejected(error.code, error.message), ...error.payload}; + if (error instanceof EffectRuntimeLockTimeoutError) return rejected("lock_acquire_timeout", "canonical renewal maintenance lock timed out"); + return {...rejected("canonical_renew_route_failed", error instanceof Error ? error.message : "canonical renewal route failed"), ...localAuthorityOpenFailure(error)}; + } +} diff --git a/loopx/control_plane/work_items/task_lease_acquire_adapter.py b/loopx/control_plane/work_items/task_lease_acquire_adapter.py index c84e64f458..074a6b01c8 100644 --- a/loopx/control_plane/work_items/task_lease_acquire_adapter.py +++ b/loopx/control_plane/work_items/task_lease_acquire_adapter.py @@ -9,14 +9,14 @@ from pathlib import Path from typing import Any +from ...history import load_registry +from ...paths import resolve_runtime_root from ..coordination.coordination_state_contract_generated import ( LOCAL_AUTHORITY_SHADOW_BINDING_SCHEMA, TASK_LEASE_ACQUIRE_REQUEST_SCHEMA, + TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA, TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, ) - -from ...history import load_registry -from ...paths import resolve_runtime_root from ..coordination.runtime_shadow import resolve_coordination_runtime_shadow_config from ..goals.active_state_event_projection import ( state_event_log_candidates as _state_event_log_candidates, @@ -54,7 +54,9 @@ def _attach_local_authority_shadow( str(lease.get("updated_at") or lease.get("released_at") or "unknown"), ) ) - from ..coordination.local_authority_shadow_observation import observe_local_authority_commit + from ..coordination.local_authority_shadow_observation import ( + observe_local_authority_commit, + ) evidence = observe_local_authority_commit( registry_path=registry_path, @@ -334,6 +336,27 @@ def _require_native_acquire_shape(payload: object) -> dict[str, Any]: return payload +def _canonical_renew_authority_facts(registry_path: Path, goal_id: str) -> dict[str, Any]: + """Only registration is external to the canonical head; bind its source.""" + for _attempt in range(TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS): + before = _authority_source_receipt("registry", registry_path) + goal = _registry_goal(load_registry(registry_path), goal_id) + after = _authority_source_receipt("registry", registry_path) + if before != after: + continue + if goal is None: + raise TaskLeaseError("canonical renewal goal is not registered", code="goal_not_found") + return { + # This neutral transport value never decides a canonical mode. + "handoff_mode": HANDOFF_MODE_LEGACY, + "registered_agent_candidates": _raw_registered_agent_candidates(goal), + "todos": [], + "todo_projection_error": None, + "source_receipts": [after], + } + raise TaskLeaseError("canonical renewal registry changed; retry", code="authority_source_changed") + + def _finalize_native_acquire_result( payload: dict[str, Any], *, @@ -591,9 +614,16 @@ def execute_native_task_lease_lifecycle( fence_operation_id = secrets.token_hex(32) for attempt in range(TASK_LEASE_AUTHORITY_SNAPSHOT_ATTEMPTS): authority: dict[str, Any] | None = None + canonical_renew = False + if normalized_operation == "renew": + from ..coordination.local_authority import local_authority_is_promoted + + canonical_renew = local_authority_is_promoted( + runtime_root=runtime_root, goal_id=goal_id + ) if registry_path is not None and (needs_authority or normalized_operation == "release"): try: - authority = task_lease_acquire_authority_facts( + authority = _canonical_renew_authority_facts(registry_path, goal_id) if canonical_renew else task_lease_acquire_authority_facts( registry_path=registry_path, goal_id=str(goal_id or ""), todo_id=str(todo_id or ""), @@ -613,7 +643,10 @@ def execute_native_task_lease_lifecycle( # optional source projection is unavailable. authority = None request: dict[str, Any] = { - "schema_version": TASK_LEASE_LIFECYCLE_NATIVE_SCHEMA_VERSION, + "schema_version": ( + TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA + if canonical_renew else TASK_LEASE_LIFECYCLE_NATIVE_SCHEMA_VERSION + ), "operation": normalized_operation, "runtime_root": str(runtime_root), "goal_id": goal_id, @@ -652,7 +685,12 @@ def execute_native_task_lease_lifecycle( # CLI argument or persisted authority fact. "current_time": _now.isoformat() if _now is not None else None, } - if registry_path is not None: + if canonical_renew: + request = {key: request[key] for key in ( + "schema_version", "operation", "runtime_root", "goal_id", "todo_id", + "owner", "idempotency_key", "expected_version", "ttl_seconds", "authority", "current_time", + )} + if registry_path is not None and not canonical_renew: registry = load_registry(registry_path) goal = _registry_goal(registry, str(goal_id)) if resolve_coordination_runtime_shadow_config(goal).enabled: @@ -661,7 +699,7 @@ def execute_native_task_lease_lifecycle( "provider": "file_v0", } compacted_todo = _compact_lifecycle_todo(todo, todo_id=str(todo_id)) - if compacted_todo is not None: + if compacted_todo is not None and not canonical_renew: request["todo"] = compacted_todo from ..effect_runtime import effect_runtime_result @@ -686,6 +724,19 @@ def execute_native_task_lease_lifecycle( owner=owner, ) result = dict(payload) + if canonical_renew and "source_authority" not in result: + raise RuntimeError("canonical renewal omitted provider evidence") + if "source_authority" in result: + if ( + normalized_operation != "renew" + or result.get("source_authority") not in ("file_v0", "sqlite_v0") + or result.get("decision_read_from_provider") is not True + or result.get("legacy_fallback_used") is not False + ): + raise RuntimeError("canonical task-lease result has invalid provider evidence") + # The canonical transaction already owns its lease/event/receipt. + # Do not attach a second legacy lease-file or shadow write. + return result if authority is not None and authority.get("handoff_mode") and "handoff_mode" not in result: result["handoff_mode"] = authority["handoff_mode"] if _legacy_provider_projection: diff --git a/loopx/control_plane/work_items/task_lease_lifecycle.ts b/loopx/control_plane/work_items/task_lease_lifecycle.ts index 81752d187e..36e816ba8f 100644 --- a/loopx/control_plane/work_items/task_lease_lifecycle.ts +++ b/loopx/control_plane/work_items/task_lease_lifecycle.ts @@ -1,4 +1,8 @@ import {leaseOwnerRejection as ownerRejection} from "./task_lease_eligibility.ts"; +import {renewCanonicalTaskLease} from "./canonical_task_lease_renew.ts"; +import {taskLeaseStableValue as stableValue, taskLeaseDigest as digest, + taskLeaseOperationIdentity as operationIdentity, + taskLeaseOperationRequestDigest as operationRequestDigest} from "./task_lease_operation_identity.ts"; import { ShadowManagementError, requireShadowPrimaryWriteAllowed } from "../coordination/shadow_management.ts"; import { LegacyCoordinationWriteError, requireLegacyCoordinationPrimaryWriteAllowed } from "../coordination/legacy_writer_fence.ts"; import { createHash, randomUUID } from "node:crypto"; @@ -65,7 +69,7 @@ import { type LeaseOutboxCapture, type LocalAuthorityShadowBinding, } from "../coordination/local_authority_shadow_outbox.ts"; -import { TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA } from "../coordination/coordination_state_contract.generated.ts"; +import { TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA } from "../coordination/coordination_state_contract.generated.ts"; export const TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION = TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA; @@ -87,6 +91,7 @@ type LifecycleStage = "validation" | "durable_writeback"; interface LifecycleRequest { schema_version: typeof TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION; + canonical_renew: boolean; operation: TaskLeaseLifecycleOperation; runtime_root: string; goal_id: string; @@ -367,13 +372,21 @@ function decodeOperation(value: unknown): TaskLeaseLifecycleOperation { function decodeRequest(value: unknown): LifecycleRequest { const input = requireJsonObject(value, "task lease lifecycle request"); - if (input.schema_version !== TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION) { + const canonicalRenew = input.schema_version === TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA; + if (input.schema_version !== TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION && !canonicalRenew) { throw new TaskLeaseLifecycleError( "Task-lease lifecycle request schema mismatch", "schema_mismatch", ); } const operation = decodeOperation(input.operation); + if (canonicalRenew && operation !== "renew") throw new TaskLeaseLifecycleError("canonical renewal schema only accepts renew", "invalid_operation"); + if (canonicalRenew) { + const fields = new Set(["schema_version", "operation", "runtime_root", "goal_id", "todo_id", "owner", + "idempotency_key", "expected_version", "ttl_seconds", "authority", "current_time"]); + const unsupported = Object.keys(input).find(key => !fields.has(key)); + if (unsupported) throw new TaskLeaseLifecycleError(`canonical renewal does not accept ${unsupported}`, "invalid_canonical_renew_request"); + } let goalId: string; let todoId: string; try { @@ -398,7 +411,11 @@ function decodeRequest(value: unknown): LifecycleRequest { let authority: AuthorityFacts | null = null; if (input.authority !== undefined && input.authority !== null) { - authority = decodeTaskLeaseAuthority(input.authority); + // A canonical-only request cannot fall back. Its real mode comes from the + // provider; stale display frontmatter is not a lease decision input. + authority = decodeTaskLeaseAuthority(canonicalRenew + ? {...requireJsonObject(input.authority, "authority"), handoff_mode: "legacy"} + : input.authority); } if ( (operation === "renew" || operation === "transfer" || operation === "terminal_verify" || operation === "holder_verify") && @@ -453,6 +470,7 @@ function decodeRequest(value: unknown): LifecycleRequest { })(); return { schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION, + canonical_renew: canonicalRenew, operation, runtime_root: typeof input.runtime_root === "string" ? input.runtime_root : (() => { throw new TaskLeaseLifecycleError("runtime_root must be a string", "invalid_runtime_root"); @@ -556,51 +574,6 @@ function effectIdFor(request: LifecycleRequest): string | null { }).effect_id; } -function stableValue(value: unknown): unknown { - if (Array.isArray(value)) return value.map(stableValue); - if (typeof value !== "object" || value === null) return value; - return Object.fromEntries( - Object.entries(value as Record) - .sort(([left], [right]) => left.localeCompare(right)) - .map(([key, child]) => [key, stableValue(child)]), - ); -} - -function digest(value: unknown): string { - return createHash("sha256") - .update(JSON.stringify(stableValue(value)), "utf8") - .digest("hex"); -} - -function operationIdentity(request: LifecycleRequest): string | null { - if (!request.idempotency_key) return null; - return digest({ - schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION, - operation: request.operation, - goal_id: request.goal_id, - todo_id: request.todo_id, - owner: request.owner, - idempotency_key: request.idempotency_key, - expected_version: request.expected_version, - }); -} - -function operationRequestDigest(request: LifecycleRequest): string | null { - if (!request.idempotency_key) return null; - return digest({ - schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION, - operation: request.operation, - goal_id: request.goal_id, - todo_id: request.todo_id, - owner: request.owner, - idempotency_key: request.idempotency_key, - expected_version: request.expected_version, - ttl_seconds: request.ttl_seconds, - new_owner: request.new_owner, - new_idempotency_key: request.new_idempotency_key, - }); -} - function operationReceiptPath(request: LifecycleRequest): string | null { const id = operationIdentity(request); if (!id) return null; @@ -2733,6 +2706,24 @@ async function fenceClose( } } +function canonicalRenewEnvelope(request: LifecycleRequest, result: JsonObject): JsonObject { + const evidence = Object.fromEntries(["source_authority", "decision_read_from_provider", "legacy_fallback_used", + "provider_revision", "cursor", "current_provider_revision", "current_cursor", "handoff_mode", + "operation_id", "recovery", "commit_status", "receipt_status"].filter(key => result[key] !== undefined) + .map(key => [key, result[key]])); + if (result.status === "applied" || result.status === "replayed" || result.status === "recovered") { + const replayed = result.status !== "applied"; + return {ok: true, schema_version: "task_lease_v0", action: "renew", status: result.status, + renewed: true, idempotent: replayed, lease: result.lease, original_receipt: result.original_receipt, + ...evidence, settlement: lifecycleSettlement(request, replayed ? "replayed" : "committed")}; + } + const envelope = failureEnvelope(request, {code: String(result.reason_code ?? result.conflict_kind ?? "canonical_renew_failed"), + message: String(result.reason ?? "canonical renewal could not complete; inspect the result before retrying"), + payload: {...evidence, status: result.status}, stage: result.failure_stage === "durable_writeback" ? "durable_writeback" : "validation"}); + delete envelope.lease_path; + return envelope; +} + export async function executeTaskLeaseLifecycle( value: unknown, dependencies: LifecycleDependencies = {}, @@ -2740,6 +2731,12 @@ export async function executeTaskLeaseLifecycle( let request: LifecycleRequest | null = null; try { request = decodeRequest(value); + if (request.canonical_renew) { + const canonical = await renewCanonicalTaskLease(request, { + now: () => lifecycleNow(request!, dependencies), beforeWrite: dependencies.beforeWrite, + }); + return canonicalRenewEnvelope(request, canonical); + } if (request.operation === "fence_close") return await fenceClose(request, dependencies); if (request.operation === "terminal_verify" || request.operation === "holder_verify") { const fence = await fenceVerify(request, dependencies); diff --git a/loopx/control_plane/work_items/task_lease_operation_identity.ts b/loopx/control_plane/work_items/task_lease_operation_identity.ts new file mode 100644 index 0000000000..40ac55c41a --- /dev/null +++ b/loopx/control_plane/work_items/task_lease_operation_identity.ts @@ -0,0 +1,44 @@ +import {createHash} from "node:crypto"; +import {TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA} from "../coordination/coordination_state_contract.generated.ts"; + +export interface TaskLeaseOperationIdentityInput { + operation: string; + goal_id: string; + todo_id: string; + owner: string | null; + idempotency_key: string | null; + expected_version: number | null; + ttl_seconds: number | null; + new_owner: string | null; + new_idempotency_key: string | null; +} + +/** Preserve the legacy v0 digest encoding; changing it changes replay identity. */ +export function taskLeaseStableValue(value: unknown): unknown { + if (Array.isArray(value)) return value.map(taskLeaseStableValue); + if (typeof value !== "object" || value === null) return value; + return Object.fromEntries(Object.entries(value as Record) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([key, child]) => [key, taskLeaseStableValue(child)])); +} + +export function taskLeaseDigest(value: unknown): string { + return createHash("sha256").update(JSON.stringify(taskLeaseStableValue(value)), "utf8").digest("hex"); +} + +export function taskLeaseOperationIdentity(request: TaskLeaseOperationIdentityInput): string | null { + if (!request.idempotency_key) return null; + return taskLeaseDigest({schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, + operation: request.operation, goal_id: request.goal_id, todo_id: request.todo_id, + owner: request.owner, idempotency_key: request.idempotency_key, + expected_version: request.expected_version}); +} + +export function taskLeaseOperationRequestDigest(request: TaskLeaseOperationIdentityInput): string | null { + if (!request.idempotency_key) return null; + return taskLeaseDigest({schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA, + operation: request.operation, goal_id: request.goal_id, todo_id: request.todo_id, + owner: request.owner, idempotency_key: request.idempotency_key, + expected_version: request.expected_version, ttl_seconds: request.ttl_seconds, + new_owner: request.new_owner, new_idempotency_key: request.new_idempotency_key}); +} diff --git a/scripts/generate_coordination_state_contract.py b/scripts/generate_coordination_state_contract.py index dc4e92196c..3bdedbd7f7 100644 --- a/scripts/generate_coordination_state_contract.py +++ b/scripts/generate_coordination_state_contract.py @@ -108,6 +108,7 @@ TASK_LEASE_PROTOCOL_KEYS = ( "acquire_request_schema", "lifecycle_request_schema", + "canonical_renew_request_schema", ) CAPABILITY_HOOK_PROTOCOL_KEYS = ( "registration_schema", diff --git a/tests/control_plane/test_canonical_lease_renew.py b/tests/control_plane/test_canonical_lease_renew.py new file mode 100644 index 0000000000..67efa971c8 --- /dev/null +++ b/tests/control_plane/test_canonical_lease_renew.py @@ -0,0 +1,130 @@ +"""Real CLI renewal of provider-owned leases, including stale/missing display.""" +import hashlib +import json +import subprocess +import sys +import tempfile +from datetime import UTC, datetime, timedelta +from pathlib import Path + +import pytest +from canonical_authority_fixture import ( + initialize_canonical_authority, + isolate_sqlite_runtime, +) + +from loopx.control_plane.coordination.coordination_state_contract import ( + TODO_DOMAIN_READ_RECORD_SCHEMA_VERSION, + TODO_DOMAIN_RECORD_FIELDS, +) +from loopx.control_plane.coordination.local_authority_shadow_projection import ( + canonical_bytes, +) +from loopx.control_plane.coordination.runtime_shadow import ( + build_todo_runtime_shadow_projection, +) +from loopx.control_plane.work_items import task_lease_acquire_adapter +from loopx.control_plane.work_items.task_lease import renew_task_lease + +REPO = Path(__file__).resolve().parents[2] + + +@pytest.mark.parametrize("provider", ["file", "sqlite"]) +@pytest.mark.parametrize("display", ["missing", "invalid_mode", "malformed_yaml"]) +def test_public_renew_preserves_canonical_authority_and_historical_receipt(tmp_path, monkeypatch, provider, display): + isolate_sqlite_runtime(tmp_path, monkeypatch) + monkeypatch.setattr(tempfile, "tempdir", str(tmp_path)) + runtime, state, registry = tmp_path / "runtime", tmp_path / "state.md", tmp_path / "registry.json" + goal = "renew-goal" + state.write_text("# Synthetic renewal\n\n## Agent Todo\n") + registry.write_text(json.dumps({"common_runtime_root": str(runtime), "goals": [{ + "id": goal, "repo": str(tmp_path), "state_file": state.name, + "coordination": {"registered_agents": ["agent-a", "agent-b"], + "runtime_shadow": {"enabled": True, "provider": "file_v0"}}, + }]})) + now = datetime.now(UTC).replace(microsecond=0) + def iso(value): + return value.isoformat().replace("+00:00", "Z") + lease = {"schema_version": "task_lease_v0", "goal_id": goal, "todo_id": "todo_renew", + "owner": "agent-a", "idempotency_key": "execution-a", "version": 1, "lease_epoch": 7, + "status": "active", "write_scopes": ["src/**"], "acquire_ttl_seconds": 600, + "acquired_at": iso(now - timedelta(seconds=60)), "updated_at": iso(now - timedelta(seconds=60)), + "expires_at": iso(now + timedelta(seconds=300))} + projection = build_todo_runtime_shadow_projection(goal_id=goal, handoff_mode="hard_lease", leases=[lease], todos=[{ + "schema_version": "todo_item_v0", "todo_id": "todo_renew", "role": "agent", "status": "open", "done": False, + "text": "Canonical renewal", "archive_state": "active", "source_section": "Agent Todo", "index": 1, + "claimed_by": "agent-a", "task_class": "advancement_task", + }]) + for todo in projection["todos"]: + todo["schema_version"] = "todo_domain_record_v0" + todo.pop("index") + todo.pop("source_section") + projection["todo_read_model"] = {"schema_version": TODO_DOMAIN_READ_RECORD_SCHEMA_VERSION, + "contract_fields": list(TODO_DOMAIN_RECORD_FIELDS), "todo_count": 1, + "records_sha256": hashlib.sha256(canonical_bytes(projection["todos"])).hexdigest()} + initialize_canonical_authority(runtime, goal, projection, state_path=state, provider=provider) + if display == "missing": + state.unlink() + elif display == "malformed_yaml": + state.write_text("---\nhandoff_mode: [broken\n---\n\n## Agent Todo\n") + else: + state.write_text("---\nhandoff_mode: invalid_stale_display\n---\n\n## Agent Todo\n") + display_before = state.read_bytes() if state.exists() else None + + def read_head(): + module = (REPO / "loopx/control_plane/coordination/local_authority_provider.ts").as_uri() + script = f"import {{openLocalAuthorityStore}} from {json.dumps(module)};const s=await openLocalAuthorityStore(process.argv[1],process.argv[2]);console.log(JSON.stringify(await s.loadAuthority()));" + child = subprocess.run(["node", "--no-warnings", "--experimental-strip-types", "--input-type=module", "-e", script, str(runtime), goal], + capture_output=True, text=True, timeout=30, check=True) + return json.loads(child.stdout) + + def outside_provider(): + return {str(path.relative_to(runtime)): path.read_bytes() for path in runtime.rglob("*") + if path.is_file() and path.relative_to(runtime).parts[:2] not in { + ("authority", "file-v0"), ("authority", "sqlite-v0")}} + + def renew(version, ttl=600, expected_exit=0): + child = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), "--format", "json", + "task-lease", "renew", "--goal-id", goal, "--todo-id", "todo_renew", "--owner", "agent-a", + "--idempotency-key", "execution-a", "--expected-version", str(version), "--ttl-seconds", str(ttl)], + capture_output=True, text=True, timeout=60, check=False) + assert child.returncode == expected_exit, child.stdout + child.stderr + return json.loads(child.stdout) + + before = outside_provider() + try: + first = renew(1) + assert first["ok"] and first["source_authority"] == provider + "_v0" + assert first["decision_read_from_provider"] is True and first["legacy_fallback_used"] is False + assert "lease_path" not in first and "projection_delivery" not in first and "authority_shadow" not in first + assert first["lease"]["version"] == 2 and first["lease"]["lease_epoch"] == 7 + assert first["lease"]["write_scopes"] == lease["write_scopes"] + second = renew(2) + assert second["lease"]["version"] == 3 + def forbidden_legacy_path(*args, **kwargs): + raise AssertionError("canonical renew used a legacy capture/write path") + with monkeypatch.context() as guard: + guard.setattr(task_lease_acquire_adapter, "task_lease_acquire_authority_facts", forbidden_legacy_path) + guard.setattr(task_lease_acquire_adapter, "_attach_local_authority_shadow", forbidden_legacy_path) + third = renew_task_lease(registry_path=registry, runtime_root=runtime, goal_id=goal, + todo_id="todo_renew", owner="agent-a", idempotency_key="execution-a", expected_version=3, ttl_seconds=600) + assert third["lease"]["version"] == 4 + current = read_head() + inspected = subprocess.run([sys.executable, "-m", "loopx.cli", "--registry", str(registry), "--format", "json", + "task-lease", "inspect", "--goal-id", goal, "--todo-id", "todo_renew"], capture_output=True, text=True, timeout=60, check=True) + observed = json.loads(inspected.stdout) + assert observed["source_authority"] == provider + "_v0" and observed["active"] is True + assert observed["lease"] == third["lease"] and observed["lease_path"] is None + replay = renew(1) + assert replay["idempotent"] is True and replay["status"] == "replayed" + assert replay["lease"] == first["lease"] and replay["original_receipt"] == first["original_receipt"] + assert read_head() == current + changed = renew(1, 900, expected_exit=1) + assert changed["error_code"] == "coordination_operation_identity_mismatch" + assert read_head() == current + assert outside_provider() == before + assert (state.read_bytes() if state.exists() else None) == display_before + assert current["cursor"] == "4" and current["head"]["todos"] == projection["todos"] + finally: + subprocess.run([sys.executable, "-c", "from loopx.control_plane.effect_runtime import effect_runtime_result; effect_runtime_result('runtime.shutdown',{},retry_safe=False)"], + capture_output=True, text=True, timeout=30, check=True) diff --git a/tests/control_plane/test_shadow_fence_caller_parity_e2e.py b/tests/control_plane/test_shadow_fence_caller_parity_e2e.py index 6de744494b..e23093ca3f 100644 --- a/tests/control_plane/test_shadow_fence_caller_parity_e2e.py +++ b/tests/control_plane/test_shadow_fence_caller_parity_e2e.py @@ -232,7 +232,9 @@ def observe_row(ws: Workspace, row: dict) -> dict: def test_fence_caller_parity(workspaces: Callable[[str], Workspace], row: dict) -> None: if row["id"] == "fixture-missing": pytest.fail(f"parity fixture is missing: {FIXTURE}") - observed = observe_row(workspaces(row["workspace"]), row) + ws = workspaces(row["workspace"]) + canonical_before = native(ws.w, "read", {}) if row["caller"] == "task_lease_renew" else None + observed = observe_row(ws, row) assert observed["exit"] == row["exit"], observed if row.get("match") == "subset": assert {key: observed["envelope"].get(key) for key in row["expect"]} == row["expect"], observed @@ -243,6 +245,27 @@ def test_fence_caller_parity(workspaces: Callable[[str], Workspace], row: dict) # operation ID and lease expiry are not a literal legacy-writer envelope. assert observed["envelope"].get("claimed_todos") or observed["envelope"].get("active_leases"), observed assert observed["envelope"].get("provider_revision"), observed + if row["caller"] == "task_lease_renew": + # The public CLI now uses a canonical-only request. Keep the old wire + # fence rows intact, and prove the new route changes only its real lease. + assert "lease_path" not in observed["envelope"], observed + canonical_after = native(ws.w, "read", {}) + before = canonical_before["head"] + after = canonical_after["head"] + assert int(after["cursor"]) == int(before["cursor"]) + 1 + assert after["head"]["todos"] == before["head"]["todos"] + old_lease = next(item for item in before["head"]["leases"] if item["todo_id"] == ws.ids["todo_b"]) + renewed = next(item for item in after["head"]["leases"] if item["todo_id"] == ws.ids["todo_b"]) + assert renewed["version"] == old_lease["version"] + 1 + for field in ("owner", "idempotency_key", "lease_epoch", "write_scopes"): + assert renewed[field] == old_lease[field] + # A later completion must refresh its proof after this successful renew. + # Explicitly preserve the stale-proof rejection before updating the + # fixture's current version for the existing valid-preview row. + stale = ws.call(*row_args(ws, "todo_complete_dry_run_leased")) + assert stale["ok"] is False and stale["error_code"] == "version_mismatch", stale + assert native(ws.w, "read", {}) == canonical_after + ws.lease_version = renewed["version"] assert observed["effect"] == row["effect"], observed assert observed["outbox_added"] == row.get("outbox_added", []), observed diff --git a/tests/control_plane_ts/canonical_task_lease_renew.test.ts b/tests/control_plane_ts/canonical_task_lease_renew.test.ts new file mode 100644 index 0000000000..e42ffb83a4 --- /dev/null +++ b/tests/control_plane_ts/canonical_task_lease_renew.test.ts @@ -0,0 +1,236 @@ +import assert from "node:assert/strict"; +import {createHash} from "node:crypto"; +import {mkdtemp, writeFile, rm, access} from "node:fs/promises"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; +import {spawn, spawnSync} from "node:child_process"; +import {fileURLToPath} from "node:url"; +import test, {type TestContext} from "node:test"; +import {executeTaskLeaseLifecycle, TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION} from "../../loopx/control_plane/work_items/task_lease_lifecycle.ts"; +import {SqliteAuthorityStore} from "../../loopx/control_plane/coordination/sqlite_authority_store.ts"; +import {sqliteAuthorityRuntime} from "../../loopx/control_plane/coordination/sqlite_runtime.ts"; +import {FileAuthorityStore} from "../../loopx/control_plane/coordination/file_authority_store.ts"; +import {selectLocalSqliteAuthority} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; +import {engageLegacyCoordinationWriterFence} from "../../loopx/control_plane/coordination/legacy_writer_fence.ts"; +import {canonicalAuthoritySha256} from "../../loopx/control_plane/coordination/authority_store_codec.ts"; +import {authorityProjectionFixture} from "./authority_projection_fixture.ts"; +import {executeCanonicalTaskLeaseRenew} from "../../loopx/control_plane/coordination/task_lease_renew.ts"; +import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; +import {taskLeaseOperationIdentity, taskLeaseOperationRequestDigest} from "../../loopx/control_plane/work_items/task_lease_operation_identity.ts"; +import {legacyCoordinationWriterFencePath} from "../../loopx/control_plane/coordination/legacy_writer_fence.ts"; +import {TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA} from "../../loopx/control_plane/coordination/coordination_state_contract.generated.ts"; +import {shadowManagementStatePath} from "../../loopx/control_plane/coordination/shadow_management.ts"; +import {atomicWriteJson} from "../../loopx/control_plane/effect_runtime_io.ts"; + +const NOW = new Date("2026-09-13T10:05:00Z"); +const CHILD = fileURLToPath(new URL("./canonical_task_lease_renew_process.ts", import.meta.url)); +function sqliteSkipReason(): string | undefined { + try { sqliteAuthorityRuntime(); return undefined; } + catch { return "requires a WAL-fixed SQLite runtime with finalized statements"; } +} +async function fixture(t: TestContext, provider: "file" | "sqlite") { + const root = await mkdtemp(join(tmpdir(), "loopx-canonical-renew-")); + t.after(() => rm(root, {recursive: true, force: true})); + const runtime = join(root, "runtime"), goal = "renew-goal", state = join(root, "state.md"); + const store = provider === "sqlite" + ? new SqliteAuthorityStore(join(runtime, "authority/sqlite-v0"), goal) + : new FileAuthorityStore(join(runtime, "authority/file-v0"), goal); + if (provider === "sqlite") assert.equal((await selectLocalSqliteAuthority(runtime, goal, true)).ok, true); + const todo = {todo_id: "todo_renew", role: "agent", status: "open", done: false, + text: "Renew the canonical lease", archive_state: "active", claimed_by: "agent-a", task_class: "advancement_task"}; + const lease = {schema_version: "task_lease_v0", goal_id: goal, todo_id: "todo_renew", owner: "agent-a", + idempotency_key: "execution-a", version: 1, lease_epoch: 7, status: "active", write_scopes: ["src/**"], + acquire_ttl_seconds: 600, acquired_at: "2026-09-13T10:00:00Z", updated_at: "2026-09-13T10:00:00Z", expires_at: "2026-09-13T10:10:00Z"}; + const projection = authorityProjectionFixture(goal, [todo], [lease], "native", {handoff_mode: "hard_lease"}); + assert.equal((await store.commitAuthority({expected_provider_revision: null, operation_id: "seed", + next_projection: projection, events: [], receipts: []})).status, "applied"); + const head = await store.loadAuthority(); assert.equal(head.status, "loaded"); + if (head.status !== "loaded") throw new Error("fixture head missing"); + const content = "# Synthetic renewal\n\n## Agent Todo\n"; + await writeFile(state, content); + assert.equal((await engageLegacyCoordinationWriterFence({schema_version: "loopx_legacy_coordination_writer_fence_engage_request_v0", + runtime_root: runtime, goal_id: goal, state_path: state, fence: {schema_version: "loopx_legacy_coordination_writer_fence_v0", + state: "engaged", goal_id: goal, fence_id: "renew-fixture", source_version: "state:1", + source_projection_sha256: canonicalAuthoritySha256(projection), expected_shadow_provider_revision: head.provider_revision}})).status, "applied"); + const request = {schema_version: TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA, operation: "renew", runtime_root: runtime, + goal_id: goal, todo_id: "todo_renew", owner: "agent-a", idempotency_key: "execution-a", expected_version: 1, ttl_seconds: 600, + authority: {handoff_mode: "hard_lease", registered_agent_candidates: [["agent-a", "agent-b"]], todos: [todo], + todo_projection_error: null, source_receipts: [{source_id: "state", path: state, state: "file", sha256: createHash("sha256").update(content).digest("hex")}]}}; + return {root, runtime, goal, state, store, request, lease, projection}; +} + +test("renew identity preserves the legacy hash and binds changed TTL as changed intent", () => { + const request = {operation: "renew", goal_id: "renew-goal", todo_id: "todo_renew", owner: "agent-a", + idempotency_key: "execution-a", expected_version: 1, ttl_seconds: 600, new_owner: null, new_idempotency_key: null}; + // Independent SHA-256 of the documented v0 sorted JSON field sets. + assert.equal(taskLeaseOperationIdentity(request), "4faaad76883d5e7f260cee986f0331b3f004ced3b981d6d935ac3a657325d04c"); + assert.equal(taskLeaseOperationRequestDigest(request), "954e868545c3b78fbce5a1aea466d097ebeef5071c7913e95606b8ed04b33d06"); + assert.equal(taskLeaseOperationIdentity({...request, ttl_seconds: 900}), taskLeaseOperationIdentity(request)); + assert.notEqual(taskLeaseOperationRequestDigest({...request, ttl_seconds: 900}), taskLeaseOperationRequestDigest(request)); +}); + +for (const provider of ["file", "sqlite"] as const) { + const skip = provider === "sqlite" ? sqliteSkipReason() : undefined; + const providerTest = (name: string, body: (t: TestContext) => Promise) => + test(name, {skip}, body); + providerTest(`${provider} legacy wire requests remain fenced`, async t => { + const {store, request} = await fixture(t, provider); const before = await store.loadAuthority(); + const result = await executeTaskLeaseLifecycle({...request, schema_version: TASK_LEASE_LIFECYCLE_REQUEST_SCHEMA_VERSION}, {now: () => NOW}); + assert.equal(result.error_code, "legacy_coordination_writer_fenced"); + assert.deepEqual(await store.loadAuthority(), before); + }); + providerTest(`${provider} canonical renew uses the public lifecycle entrypoint without a legacy lease file`, async t => { + const {runtime, goal, store, request, lease} = await fixture(t, provider); + const result = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(result.ok, true, JSON.stringify(result)); + assert.equal(result.source_authority, `${provider}_v0`); + assert.equal(result.legacy_fallback_used, false); + assert.equal(result.lease_path, undefined); + assert.deepEqual(result.lease, {...lease, version: 2, updated_at: NOW.toISOString().replace(".000Z", "Z"), expires_at: "2026-09-13T10:15:00Z"}); + const after = await store.loadAuthority(); assert.equal(after.status, "loaded"); + if (after.status === "loaded") assert.equal(after.cursor, "2"); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/todo_renew.json"))); + }); + providerTest(`${provider} old renewal receipt survives later renewals and expiry without extending again`, async t => { + const {store, request, runtime, goal} = await fixture(t, provider); + const first = await executeTaskLeaseLifecycle(request, {now: () => NOW}); assert.equal(first.ok, true); + const second = await executeTaskLeaseLifecycle({...request, expected_version: 2}, {now: () => new Date("2026-09-13T10:06:00Z")}); + assert.equal(second.ok, true); assert.equal((second.lease as Record).version, 3); + const current = await store.loadAuthority(); + const replay = await executeTaskLeaseLifecycle(request, {now: () => new Date("2026-09-14T10:00:00Z")}); + assert.equal(replay.ok, true); assert.equal(replay.idempotent, true); assert.equal(replay.status, "replayed"); + assert.deepEqual(replay.lease, first.lease); assert.deepEqual(replay.original_receipt, first.original_receipt); + assert.deepEqual(await store.loadAuthority(), current); + const reuse = await executeTaskLeaseLifecycle({...request, ttl_seconds: 900}, {now: () => NOW}); + assert.equal(reuse.error_code, "coordination_operation_identity_mismatch"); + assert.deepEqual(await store.loadAuthority(), current); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/.lifecycle-operations"))); + }); + providerTest(`${provider} canonical facts override stale caller Todo and mode facts`, async t => { + const {store, request} = await fixture(t, provider); + const stale = {...request, authority: {...request.authority, handoff_mode: "soft_claim", + todos: [{...request.authority.todos[0]!, status: "done", claimed_by: "agent-b"}]}}; + const result = await executeTaskLeaseLifecycle(stale, {now: () => NOW}); + assert.equal(result.ok, true, JSON.stringify(result)); assert.equal(result.handoff_mode, "hard_lease"); + const head = await store.loadAuthority(); if (head.status !== "loaded") throw new Error("missing head"); + assert.equal((head.head.todos as Record[])[0]!.status, "open"); + }); + for (const [label, changes, now, code] of [ + ["wrong owner", {owner: "agent-b"}, NOW, "owner_conflicts_with_claim"], + ["wrong execution", {idempotency_key: "wrong"}, NOW, "lease_cas_mismatch"], + ["stale version", {expected_version: 0}, NOW, "version_mismatch"], + ["missing version", {expected_version: null}, NOW, "version_required"], + ["expired", {}, new Date("2026-09-13T10:10:00Z"), "lease_not_active"], + ] as const) { + providerTest(`${provider} rejects ${label} without canonical or legacy writes`, async t => { + const {store, request, runtime, goal} = await fixture(t, provider); + const before = await store.loadAuthority(); + const result = await executeTaskLeaseLifecycle({...request, ...changes}, {now: () => now}); + assert.equal(result.ok, false); assert.equal(result.error_code, code, JSON.stringify(result)); + assert.deepEqual(await store.loadAuthority(), before); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/todo_renew.json"))); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/.lifecycle-operations"))); + }); + } + providerTest(`${provider} source revocation before commit is rejected with no write`, async t => { + const {store, state, request} = await fixture(t, provider); const before = await store.loadAuthority(); + const result = await executeTaskLeaseLifecycle(request, {now: () => NOW, beforeWrite: async () => {await writeFile(state, "changed source");}}); + assert.equal(result.ok, false); assert.equal(result.error_code, "authority_source_changed"); + assert.deepEqual(await store.loadAuthority(), before); + }); + providerTest(`${provider} maintenance guard and canonical mode remain authoritative`, async t => { + const {store, request, runtime, goal, projection} = await fixture(t, provider); + let head = await store.loadAuthority(); if (head.status !== "loaded") throw new Error("missing head"); + await store.commitAuthority({expected_provider_revision: head.provider_revision, operation_id: "mode-change", + next_projection: {...projection, handoff_mode: "soft_claim"}, events: [], receipts: []}); + const before = await store.loadAuthority(); + const refused = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(refused.error_code, "handoff_mode_forbids_lease"); assert.equal(refused.handoff_mode, "soft_claim"); + assert.deepEqual(await store.loadAuthority(), before); + await atomicWriteJson(shadowManagementStatePath(runtime, goal), {}); + const blocked = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(blocked.error_code, "shadow_management_state_invalid"); assert.deepEqual(await store.loadAuthority(), before); + }); + providerTest(`${provider} fenced missing head never initializes a new lease`, async t => { + const {store, request, runtime, goal} = await fixture(t, provider); + await rm(store.path); + const result = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(result.ok, false); assert.equal(result.legacy_fallback_used, false); + await assert.rejects(access(store.path)); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/todo_renew.json"))); + await assert.rejects(access(join(runtime, "goals", goal, "task-leases/.lifecycle-operations"))); + }); + providerTest(`${provider} invalid fence stays fail-closed`, async t => { + const {store, request, runtime, goal} = await fixture(t, provider); const before = await store.loadAuthority(); + await writeFile(legacyCoordinationWriterFencePath(runtime, goal), "{"); + const result = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(result.ok, false); assert.equal(result.error_code, "legacy_writer_fence_read_failed"); + assert.deepEqual(await store.loadAuthority(), before); + }); + providerTest(`${provider} a canonical-only request cannot downgrade after losing its fence`, async t => { + const {store, request, runtime, goal, lease} = await fixture(t, provider); const before = await store.loadAuthority(); + const legacyPath = join(runtime, "goals", goal, "task-leases/todo_renew.json"); + const legacyBytes = JSON.stringify(lease); await writeFile(legacyPath, legacyBytes); + await rm(legacyCoordinationWriterFencePath(runtime, goal)); + const result = await executeTaskLeaseLifecycle({...request, schema_version: TASK_LEASE_CANONICAL_RENEW_REQUEST_SCHEMA}, {now: () => NOW}); + assert.equal(result.ok, false); assert.equal(result.error_code, "canonical_renew_fence_missing"); + assert.deepEqual(await store.loadAuthority(), before); + assert.equal(await import("node:fs/promises").then(fs => fs.readFile(legacyPath, "utf8")), legacyBytes); + }); + providerTest(`${provider} canonical renew rejects unsupported mutation controls`, async t => { + const {store, request} = await fixture(t, provider); const before = await store.loadAuthority(); + for (const fields of [{dry_run: true}, {new_owner: "agent-b"}, {release_lease: true}]) { + const result = await executeTaskLeaseLifecycle({...request, ...fields}, {now: () => NOW}); + assert.equal(result.error_code, "invalid_canonical_renew_request"); assert.deepEqual(await store.loadAuthority(), before); + } + const transfer = await executeTaskLeaseLifecycle({...request, operation: "transfer"}, {now: () => NOW}); + assert.equal(transfer.error_code, "invalid_operation"); assert.deepEqual(await store.loadAuthority(), before); + }); + providerTest(`${provider} ambiguous response recovers the exact committed renewal`, async t => { + const {store, request, lease} = await fixture(t, provider); + const wrapper: AuthorityStore = {storeIdentity: () => store.storeIdentity(), loadAuthority: () => store.loadAuthority(), + readReceipt: id => store.readReceipt(id), scanCommitted: (...args) => store.scanCommitted(...args), + commitAuthority: async input => {assert.equal((await store.commitAuthority(input)).status, "applied"); throw new Error("response lost");}}; + const result = await executeCanonicalTaskLeaseRenew(wrapper, {...request, registered_agents: ["agent-a", "agent-b"], now: NOW}); + assert.equal(result.status, "recovered", JSON.stringify(result)); + assert.deepEqual(result.lease, {...lease, version: 2, updated_at: "2026-09-13T10:05:00Z", expires_at: "2026-09-13T10:15:00Z"}); + const head = await store.loadAuthority(); assert.equal(head.status, "loaded"); if (head.status === "loaded") assert.equal(head.cursor, "2"); + }); + for (const boundary of ["before", "after"] as const) { + providerTest(`${provider} real process interruption ${boundary} provider commit preserves renewal identity`, async t => { + const {root, store, request} = await fixture(t, provider); + const config = join(root, "child.json"); await writeFile(config, JSON.stringify({...request, now: NOW.toISOString()})); + const before = await store.loadAuthority(); + const child = spawnSync(process.execPath, ["--no-warnings", "--experimental-sqlite", "--experimental-strip-types", CHILD, config, boundary, "600"], {encoding: "utf8", timeout: 30000}); + assert.equal(child.signal, "SIGKILL", child.stderr); + if (boundary === "before") assert.deepEqual(await store.loadAuthority(), before); + const result = await executeTaskLeaseLifecycle(request, {now: () => NOW}); + assert.equal(result.ok, true); assert.equal(result.status, boundary === "before" ? "applied" : "replayed"); + const final = await store.loadAuthority(); if (final.status !== "loaded") throw new Error("missing head"); + assert.equal(final.cursor, "2"); assert.equal((result.lease as Record).version, 2); + assert.equal((await executeTaskLeaseLifecycle(request, {now: () => NOW})).status, "replayed"); + assert.deepEqual(await store.loadAuthority(), final); + }); + } + for (const differentIntent of [false, true]) { + providerTest(`${provider} real processes arbitrate ${differentIntent ? "different" : "identical"} renew intent at one version`, async t => { + const {root, store, request} = await fixture(t, provider); + const config = join(root, "race.json"); await writeFile(config, JSON.stringify({...request, now: NOW.toISOString()})); + const children: ReturnType[] = []; let readyCount = 0; + t.after(() => {for (const child of children) child.kill();}); + const results = await Promise.all([600, differentIntent ? 900 : 600].map(ttl => new Promise>((resolve, reject) => { + const child = spawn(process.execPath, ["--no-warnings", "--experimental-sqlite", "--experimental-strip-types", CHILD, config, "race", String(ttl)], {stdio: ["ignore", "pipe", "pipe", "ipc"]}); + children.push(child); let output = "", error = ""; + const timeout = setTimeout(() => {child.kill("SIGKILL"); reject(new Error("renew race timed out"));}, 30000); + child.stdout!.on("data", data => {output += String(data);}); child.stderr!.on("data", data => {error += String(data);}); + child.on("message", () => {if (++readyCount === 2) for (const peer of children) peer.send!({go: true});}); + child.on("error", reject); + child.on("close", code => {clearTimeout(timeout); if (code !== 0) reject(new Error(error)); else resolve(JSON.parse(output));}); + }))); + assert.deepEqual(results.map(r => r.status).sort(), differentIntent ? ["applied", "failed"] : ["applied", "recovered"]); + if (differentIntent) assert.equal(results.find(r => r.status === "failed")!.reason_code, "coordination_operation_identity_mismatch"); + const head = await store.loadAuthority(); if (head.status !== "loaded") throw new Error("missing head"); + assert.equal(head.cursor, "2"); assert.equal((head.head.leases as Record[])[0]!.version, 2); + }); + } +} diff --git a/tests/control_plane_ts/canonical_task_lease_renew_process.ts b/tests/control_plane_ts/canonical_task_lease_renew_process.ts new file mode 100644 index 0000000000..876dfaf6e8 --- /dev/null +++ b/tests/control_plane_ts/canonical_task_lease_renew_process.ts @@ -0,0 +1,30 @@ +/** Disposable child for renewal race and lost-response integration tests. */ +import {readFileSync} from "node:fs"; +import {openLocalAuthorityStore} from "../../loopx/control_plane/coordination/local_authority_provider.ts"; +import {executeCanonicalTaskLeaseRenew} from "../../loopx/control_plane/coordination/task_lease_renew.ts"; +import type {AuthorityStore} from "../../loopx/control_plane/coordination/authority_store.ts"; +const [path, mode, ttl] = process.argv.slice(2); +const request = JSON.parse(readFileSync(path!, "utf8")); +const store = await openLocalAuthorityStore(request.runtime_root, request.goal_id); +let release: () => void = () => {}; +const barrier = new Promise(resolve => {release = resolve;}); +process.on("message", () => release()); +const measured: AuthorityStore = { + storeIdentity: () => store.storeIdentity(), readReceipt: id => store.readReceipt(id), + scanCommitted: (...args) => store.scanCommitted(...args), + loadAuthority: async () => { + const head = await store.loadAuthority(); + if (mode === "race") {process.send!({ready: true}); await barrier;} + return head; + }, + commitAuthority: async input => { + if (mode === "before") process.kill(process.pid, "SIGKILL"); + const result = await store.commitAuthority(input); + if (mode === "after" && result.status === "applied") process.kill(process.pid, "SIGKILL"); + return result; + }, +}; +const result = await executeCanonicalTaskLeaseRenew(measured, {...request, + registered_agents: ["agent-a", "agent-b"], now: new Date(request.now), ttl_seconds: Number(ttl)}); +process.stdout.write(JSON.stringify(result)); +if (process.connected) process.disconnect(); diff --git a/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json b/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json index db8950644b..abfb99a5df 100644 --- a/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json +++ b/tests/fixtures/control_plane/legacy_writer_fence_caller_parity_v0.json @@ -1131,27 +1131,23 @@ } }, { - "id": "cli-task_lease_renew-engaged", + "id": "cli-task_lease_renew-canonical", "surface": "cli", "workspace": "w1", "caller": "task_lease_renew", "fence_state": "engaged", - "exit": 1, + "exit": 0, "expect": { - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_b}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "blocked", - "reason_code": "legacy_coordination_writer_fenced", - "authority_mode": "file_v0", - "fence_id": "caller-fixture" - }, - "handoff_mode": "hard_lease", - "ok": false, + "ok": true, "schema_version": "task_lease_v0", "action": "renew", - "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", - "error_code": "legacy_coordination_writer_fenced" + "status": "applied", + "renewed": true, + "idempotent": false, + "handoff_mode": "hard_lease", + "source_authority": "file_v0", + "decision_read_from_provider": true, + "legacy_fallback_used": false }, "effect": { "added": [], @@ -1172,8 +1168,25 @@ "error": "native task-lease lifecycle schema mismatch", "error_code": "RuntimeError", "schema_version": "task_lease_v0" + }, + "9231d5b3d": { + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_b}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "blocked", + "reason_code": "legacy_coordination_writer_fenced", + "authority_mode": "file_v0", + "fence_id": "caller-fixture" + }, + "handoff_mode": "hard_lease", + "ok": false, + "schema_version": "task_lease_v0", + "action": "renew", + "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", + "error_code": "legacy_coordination_writer_fenced" } - } + }, + "match": "subset" }, { "id": "cli-task_lease_transfer-engaged", @@ -1717,27 +1730,23 @@ "baseline": null }, { - "id": "cli-task_lease_renew-engaged-capture", + "id": "cli-task_lease_renew-canonical-capture", "surface": "cli", "workspace": "w2", "caller": "task_lease_renew", "fence_state": "engaged", - "exit": 1, + "exit": 0, "expect": { - "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_b}.json", - "write_check": { - "schema_version": "loopx_legacy_coordination_write_check_result_v0", - "status": "blocked", - "reason_code": "legacy_coordination_writer_fenced", - "authority_mode": "file_v0", - "fence_id": "caller-fixture" - }, - "handoff_mode": "hard_lease", - "ok": false, + "ok": true, "schema_version": "task_lease_v0", "action": "renew", - "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", - "error_code": "legacy_coordination_writer_fenced" + "status": "applied", + "renewed": true, + "idempotent": false, + "handoff_mode": "hard_lease", + "source_authority": "file_v0", + "decision_read_from_provider": true, + "legacy_fallback_used": false }, "effect": { "added": [], @@ -1745,7 +1754,25 @@ "changed": [] }, "outbox_added": [], - "baseline": null + "baseline": { + "9231d5b3d": { + "lease_path": "{runtime_root}/goals/observable/task-leases/{todo_b}.json", + "write_check": { + "schema_version": "loopx_legacy_coordination_write_check_result_v0", + "status": "blocked", + "reason_code": "legacy_coordination_writer_fenced", + "authority_mode": "file_v0", + "fence_id": "caller-fixture" + }, + "handoff_mode": "hard_lease", + "ok": false, + "schema_version": "task_lease_v0", + "action": "renew", + "error": "legacy coordination writer is fenced; use the promoted canonical authority (file_v0) for goal observable; fence caller-fixture; the primary record was not changed", + "error_code": "legacy_coordination_writer_fenced" + } + }, + "match": "subset" }, { "id": "cli-task_lease_transfer-engaged-capture", diff --git a/tsconfig.control-plane.json b/tsconfig.control-plane.json index 29b84286df..5d2b76be3b 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -115,6 +115,7 @@ "tests/control_plane_ts/todo_completion_validation_plan.test.ts", "tests/control_plane_ts/todo_next_action.test.ts", "tests/control_plane_ts/task_lease_acquire.test.ts", + "tests/control_plane_ts/canonical_task_lease_renew.test.ts", "tests/control_plane_ts/task_lease_acquire_cli.test.ts", "tests/control_plane_ts/legacy_writer_fence_caller_parity_support.ts", "tests/control_plane_ts/legacy_writer_fence_caller_parity.test.ts",