From 5ba40d0ddd57297da1526090904d66d76f08bebf Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Sun, 13 Sep 2026 09:02:16 -0700 Subject: [PATCH 1/4] feat(coordination): renew canonical leases through provider CAS Signed-off-by: Lihua <1017343802@qq.com> --- .../coordination_state_contract.generated.ts | 4 +- .../coordination_state_contract_generated.py | 4 +- .../coordination_state_contract_v0.json | 3 +- .../coordination/local_authority_runtime.ts | 12 +- .../coordination/local_authority_write.ts | 11 + .../coordination/task_lease_renew.ts | 121 ++++++++++ .../work_items/canonical_task_lease_renew.ts | 60 +++++ .../work_items/task_lease_acquire.ts | 2 +- .../work_items/task_lease_acquire_adapter.py | 67 ++++- .../work_items/task_lease_lifecycle.ts | 93 ++++--- .../task_lease_operation_identity.ts | 44 ++++ .../generate_coordination_state_contract.py | 1 + .../test_canonical_lease_renew.py | 130 ++++++++++ .../canonical_task_lease_renew.test.ts | 228 ++++++++++++++++++ .../canonical_task_lease_renew_process.ts | 30 +++ tsconfig.control-plane.json | 1 + 16 files changed, 741 insertions(+), 70 deletions(-) create mode 100644 loopx/control_plane/coordination/local_authority_write.ts create mode 100644 loopx/control_plane/coordination/task_lease_renew.ts create mode 100644 loopx/control_plane/work_items/canonical_task_lease_renew.ts create mode 100644 loopx/control_plane/work_items/task_lease_operation_identity.ts create mode 100644 tests/control_plane/test_canonical_lease_renew.py create mode 100644 tests/control_plane_ts/canonical_task_lease_renew.test.ts create mode 100644 tests/control_plane_ts/canonical_task_lease_renew_process.ts 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/local_authority_runtime.ts b/loopx/control_plane/coordination/local_authority_runtime.ts index 7ba94c105e..671ef901db 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 @@ function sourceAuthorityFor(store: AuthorityStore): "sqlite_v0" | "file_v0" { return store instanceof SqliteAuthorityStore ? "sqlite_v0" : "file_v0"; } -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..b89ea49565 --- /dev/null +++ b/loopx/control_plane/coordination/task_lease_renew.ts @@ -0,0 +1,121 @@ +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 {leaseEpoch, leaseIsActive, leaseVersion, normalizeAgent, normalizeGoalId, + normalizeHandoffMode, 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 = normalizeHandoffMode(head.head.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.ts b/loopx/control_plane/work_items/task_lease_acquire.ts index 9157f2af37..ed57f45300 100644 --- a/loopx/control_plane/work_items/task_lease_acquire.ts +++ b/loopx/control_plane/work_items/task_lease_acquire.ts @@ -421,7 +421,7 @@ export function decodeTaskLeaseAuthority(value: unknown): AuthorityFacts { }; } -function normalizeHandoffMode(value: unknown): string { +export function normalizeHandoffMode(value: unknown): string { const mode = compact(value) || "legacy"; if (!new Set(["legacy", "soft_claim", "hard_lease"]).has(mode)) { throw new TaskLeaseAcquireError( 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_ts/canonical_task_lease_renew.test.ts b/tests/control_plane_ts/canonical_task_lease_renew.test.ts new file mode 100644 index 0000000000..c28e638ac1 --- /dev/null +++ b/tests/control_plane_ts/canonical_task_lease_renew.test.ts @@ -0,0 +1,228 @@ +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 {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)); +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) { + test(`${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); + }); + test(`${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"))); + }); + test(`${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"))); + }); + test(`${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) { + test(`${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"))); + }); + } + test(`${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); + }); + test(`${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); + }); + test(`${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"))); + }); + test(`${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); + }); + test(`${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); + }); + test(`${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); + }); + test(`${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) { + test(`${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]) { + test(`${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/tsconfig.control-plane.json b/tsconfig.control-plane.json index 71b0b27ac9..4236e5a230 100644 --- a/tsconfig.control-plane.json +++ b/tsconfig.control-plane.json @@ -111,6 +111,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", From 3c45512e6f4b44fecaaf64df799e14f2d08935c0 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Sun, 13 Sep 2026 09:02:45 -0700 Subject: [PATCH 2/4] docs: explain canonical lease renewal and recovery boundaries Signed-off-by: Lihua <1017343802@qq.com> --- docs/reference/canonical-lease-renew.md | 89 ++++++++++++++++++++++++ docs/reference/sqlite-authority-store.md | 4 ++ 2 files changed, 93 insertions(+) create mode 100644 docs/reference/canonical-lease-renew.md 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 a5acafa574..716ac7a4a7 100644 --- a/docs/reference/sqlite-authority-store.md +++ b/docs/reference/sqlite-authority-store.md @@ -190,3 +190,7 @@ payload, 10k/100k commits, and 100 samples per read workload. It emits measured latency percentiles and database bytes, then deletes only its temporary database. These accelerated measurements do not satisfy the separate ten-day soak, retention, disk-exhaustion, restore or promotion gates. + +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. From f218307dfc38f85f737231b177a8385b7fbd7b75 Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Sun, 13 Sep 2026 09:39:38 -0700 Subject: [PATCH 3/4] test(coordination): align fence caller oracle with canonical renew Signed-off-by: Lihua <1017343802@qq.com> --- .../test_shadow_fence_caller_parity_e2e.py | 25 +++++- .../legacy_writer_fence_caller_parity_v0.json | 87 ++++++++++++------- 2 files changed, 81 insertions(+), 31 deletions(-) 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/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", From 0a043839e1833e4cb1c4ca2b6a1a35eb14a1687b Mon Sep 17 00:00:00 2001 From: Lihua <1017343802@qq.com> Date: Sun, 13 Sep 2026 21:17:51 -0700 Subject: [PATCH 4/4] test(coordination): respect SQLite runtime admission Signed-off-by: Lihua <1017343802@qq.com> --- .../canonical_task_lease_renew.test.ts | 36 +++++++++++-------- 1 file changed, 22 insertions(+), 14 deletions(-) diff --git a/tests/control_plane_ts/canonical_task_lease_renew.test.ts b/tests/control_plane_ts/canonical_task_lease_renew.test.ts index c28e638ac1..e42ffb83a4 100644 --- a/tests/control_plane_ts/canonical_task_lease_renew.test.ts +++ b/tests/control_plane_ts/canonical_task_lease_renew.test.ts @@ -8,6 +8,7 @@ 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"; @@ -23,6 +24,10 @@ 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})); @@ -65,13 +70,16 @@ test("renew identity preserves the legacy hash and binds changed TTL as changed }); for (const provider of ["file", "sqlite"] as const) { - test(`${provider} legacy wire requests remain fenced`, async t => { + 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); }); - test(`${provider} canonical renew uses the public lifecycle entrypoint without a legacy lease file`, async t => { + 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)); @@ -83,7 +91,7 @@ for (const provider of ["file", "sqlite"] as const) { if (after.status === "loaded") assert.equal(after.cursor, "2"); await assert.rejects(access(join(runtime, "goals", goal, "task-leases/todo_renew.json"))); }); - test(`${provider} old renewal receipt survives later renewals and expiry without extending again`, async t => { + 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")}); @@ -98,7 +106,7 @@ for (const provider of ["file", "sqlite"] as const) { assert.deepEqual(await store.loadAuthority(), current); await assert.rejects(access(join(runtime, "goals", goal, "task-leases/.lifecycle-operations"))); }); - test(`${provider} canonical facts override stale caller Todo and mode facts`, async t => { + 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"}]}}; @@ -114,7 +122,7 @@ for (const provider of ["file", "sqlite"] as const) { ["missing version", {expected_version: null}, NOW, "version_required"], ["expired", {}, new Date("2026-09-13T10:10:00Z"), "lease_not_active"], ] as const) { - test(`${provider} rejects ${label} without canonical or legacy writes`, async t => { + 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}); @@ -124,13 +132,13 @@ for (const provider of ["file", "sqlite"] as const) { await assert.rejects(access(join(runtime, "goals", goal, "task-leases/.lifecycle-operations"))); }); } - test(`${provider} source revocation before commit is rejected with no write`, async t => { + 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); }); - test(`${provider} maintenance guard and canonical mode remain authoritative`, async t => { + 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", @@ -143,7 +151,7 @@ for (const provider of ["file", "sqlite"] as const) { const blocked = await executeTaskLeaseLifecycle(request, {now: () => NOW}); assert.equal(blocked.error_code, "shadow_management_state_invalid"); assert.deepEqual(await store.loadAuthority(), before); }); - test(`${provider} fenced missing head never initializes a new lease`, async t => { + 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}); @@ -152,14 +160,14 @@ for (const provider of ["file", "sqlite"] as const) { 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"))); }); - test(`${provider} invalid fence stays fail-closed`, async t => { + 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); }); - test(`${provider} a canonical-only request cannot downgrade after losing its fence`, async t => { + 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); @@ -169,7 +177,7 @@ for (const provider of ["file", "sqlite"] as const) { assert.deepEqual(await store.loadAuthority(), before); assert.equal(await import("node:fs/promises").then(fs => fs.readFile(legacyPath, "utf8")), legacyBytes); }); - test(`${provider} canonical renew rejects unsupported mutation controls`, async t => { + 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}); @@ -178,7 +186,7 @@ for (const provider of ["file", "sqlite"] as const) { const transfer = await executeTaskLeaseLifecycle({...request, operation: "transfer"}, {now: () => NOW}); assert.equal(transfer.error_code, "invalid_operation"); assert.deepEqual(await store.loadAuthority(), before); }); - test(`${provider} ambiguous response recovers the exact committed renewal`, async t => { + 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), @@ -189,7 +197,7 @@ for (const provider of ["file", "sqlite"] as const) { 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) { - test(`${provider} real process interruption ${boundary} provider commit preserves renewal identity`, async t => { + 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(); @@ -205,7 +213,7 @@ for (const provider of ["file", "sqlite"] as const) { }); } for (const differentIntent of [false, true]) { - test(`${provider} real processes arbitrate ${differentIntent ? "different" : "identical"} renew intent at one version`, async t => { + 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;