Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
5165626
fix(coordination): retain archive receipts until projection delivery
steven-kid Sep 8, 2026
6977406
Merge main and align archive delivery retirement contract
steven-kid Sep 8, 2026
8d2d3c7
Merge latest main documentation
steven-kid Sep 8, 2026
674d32d
Merge main canonical Markdown recovery
steven-kid Sep 8, 2026
28354cc
fix(todos): allow archive retry to recover missing display
steven-kid Sep 8, 2026
aad5f75
Merge main lifecycle admission unification
steven-kid Sep 8, 2026
ae67052
Merge main documentation updates
steven-kid Sep 8, 2026
2e68039
Merge main release 1.0.2 into archive recovery contribution
steven-kid Sep 9, 2026
468860d
Merge main lifecycle and runtime updates into archive recovery contri…
steven-kid Sep 9, 2026
b7b47d7
Merge main frontier checkpoint fix into archive recovery contribution
steven-kid Sep 9, 2026
cd56660
Merge main diagnostics updates into archive recovery contribution
steven-kid Sep 9, 2026
5791cb4
Merge main independent review fix into archive recovery contribution
steven-kid Sep 9, 2026
638fbcf
Merge main resume and delivery-history migration into archive recover…
steven-kid Sep 9, 2026
2424521
Merge main delivery claim validation and review queue updates
steven-kid Sep 9, 2026
cb6b67d
Merge main canonical delivery response and RFC execution stages
steven-kid Sep 9, 2026
6c575ca
fix(qualification): scan decoded scheduler facts before actor dispatch
steven-kid Sep 9, 2026
b9bd01a
Merge main authoring and monitor settlement updates
steven-kid Sep 10, 2026
37ac787
Merge main Node 24 CI and concluded review updates
steven-kid Sep 10, 2026
138e1a2
Merge remote-tracking branch 'origin/main' into codex/review-4101-int…
huangruiteng Sep 10, 2026
afbd777
Merge remote-tracking branch 'origin/main' into codex/review-4101-int…
huangruiteng Sep 10, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 93 additions & 0 deletions docs/architecture/todo-archive-delivery-recovery.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
# Todo archive delivery recovery

This contract hardens the promoted archive transaction introduced by
[#4053](https://github.com/huangruiteng/loopx/pull/4053), under the
[TypeScript control-plane migration RFC](./rfcs/typescript-control-plane-migration-v0.md).
Archive selection, standing-decision retention, CAS, and historical receipts
remain owned by the native transaction.

## Pending delivery and retry identity

The local TypeScript adapter durably records one pending archive attempt per
goal and role before executing the transaction. The record binds the operation
ID, retention limit, observed provider revision, and authority store identity.
It is correlation state, not a second Todo authority or a copy of the receipt.
The existing per-goal maintenance lock serializes attempt creation and retirement.

While that attempt is pending, a repeated archive request for the same role and
limit reuses its identity. Receipt replay precedes current-head revision checks
and archive selection, so a committed result survives a process exit before
Markdown projection, a projection write failure, and unrelated head advancement.
The replay reports the original receipt, moved IDs, count, and revision without
creating another canonical commit. A different retention limit is rejected
until the pending delivery is settled.

Fresh Python requests bind the revision they actually observed. If that head
changes before selection, the native transaction rejects the stale request.
Older v0 callers that omit the optional revision retain their original receipt
hash. A fresh attempt preserves an explicitly supplied revision; a pending
retry keeps the original attempt's binding rather than today's observed head.

After the Python projection adapter delivers or verifies the current canonical
view, it acknowledges the exact operation through
`coordination.local_authority.todo_archive_ack`. The native owner checks the
receipt and store identity before retiring the attempt. A stale acknowledgement
cannot retire a newer attempt. Transport failure preserves the already committed
result; uncertain transaction outcomes retain the pending identity for retry.
Proven pre-commit rejection releases it so a corrected request can proceed.

Preview bypasses pending attempts and historical receipt replay and performs no
writes. A fresh empty archive creates no canonical receipt and does not advance
the provider revision. Local correlation storage is bounded to two slots per
goal, independent of archive history size.

## Delivery boundary and migration cost

The recovery boundary ends at successful projection acknowledgement. A later
identical CLI invocation is a new archive request and can process a new batch.
This does not promise recovery of a CLI response lost after acknowledgement or
discover attempts made before correlation records existed. Callers requiring
that stronger guarantee need an explicit caller-retained request identity.

Missing Markdown can be rebuilt from canonical Todos during archive retry;
successful recovery still requires the same exact-attempt acknowledgement.
An explicit `todo project-markdown` rebuild repairs the display without
acknowledging an archive attempt. The next archive with the same role and limit
first replays and acknowledges that prior batch; it does not archive a new one.

A changed promoted archive uses four Python/TypeScript request-responses:
authority read, archive transaction, projection readback, and acknowledgement.
The previous path used three; empty and preview paths add no acknowledgement.
The extra durable correlation writes and acknowledgement are a correctness cost,
not a claimed migration speedup. Archive domain rules stay in TypeScript.
Retire the separate Python ACK bridge when the canonical journal consumer owns
both projection delivery and exact-attempt acknowledgement. Remove facade parts
as their concrete callers converge, retaining required Python projection
adapters. A native CLI permits full transport removal; it is not a prerequisite
for retiring redundant boundaries. Keep provider receipts and independent
crash/retry conformance coverage.

## 中文契约

本变更修复 promoted archive 在 canonical commit 后、Markdown 投影交付前失败时的
重试身份丢失。TypeScript 在执行前持久记录每个 goal/role 的未交付 operation ID、
保留数量、观察到的 provider revision 和 store identity。相同 role/limit 的重试先
恢复原 receipt,再考虑当前 head;不会重复提交,也不会把原归档数量错误地报告为零。
未完成交付时修改保留数量会被拒绝。

新 Python 请求绑定实际读到的 revision;无历史 receipt 且 head 已变化时拒绝执行。
旧 v0 请求省略 revision 时保持原有请求 hash。预览不读取未交付记录、不重放历史结果、
不写状态;空归档不推进 canonical revision。确认仅在投影交付成功后进行,且必须匹配
原 operation、receipt 和 store。过期确认不能清除新操作;结果不确定时保留重试身份,
确定未提交的拒绝允许后续修正请求继续。

Markdown 缺失时,归档重试可从 canonical Todos 重建投影,成功后仍须确认原 attempt。
单独运行 `todo project-markdown` 只修复显示,不确认归档 attempt;之后相同 role/limit
的归档调用会先重放并确认旧批次,不会归档新一批。

恢复保证截至投影确认成功;之后相同 CLI 调用代表新一批归档,不承诺恢复确认之后才
丢失的 stdout,也不追溯发现旧版本未记录的调用。发生归档的路径从三次跨运行时调用
增加为四次,额外确认和持久化是明确的正确性成本。canonical journal 消费者同时接管
投影交付与精确 attempt 确认后,删除独立 Python ACK 桥接;按具体调用者收敛逐步删除
facade,保留仍必需的 Python 投影 adapter。原生 CLI 允许彻底移除 transport,但不是
删除冗余边界的普遍前提。保留 provider receipt 和独立崩溃恢复测试。
170 changes: 170 additions & 0 deletions loopx/control_plane/coordination/local_archive_attempt.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,170 @@
import { readFile } from "node:fs/promises";
import { join } from "node:path";

import type { JsonObject } from "../effect_program.ts";
import { durableWriteJson } from "../effect_runtime_io.ts";
import type { AuthorityStore } from "./authority_store.ts";
import {
canonicalAuthorityObject,
canonicalAuthoritySha256,
hasExactAuthorityKeys,
requireAuthorityStoreId,
} from "./authority_store_codec.ts";
import { shadowManagementDirectory } from "./shadow_management.ts";
import {
COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA,
COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA,
executeCoordinationTodoArchiveCompleted,
type CoordinationTodoArchiveInput,
} from "./todo_terminal_lifecycle.ts";

const ATTEMPT_SCHEMA = "loopx_local_todo_archive_attempt_v0";
export const LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA =
"loopx_local_coordination_todo_archive_ack_result_v0";

type Role = CoordinationTodoArchiveInput["role"];
interface ArchiveAttempt {
operation_id: string;
expected_provider_revision: string | null;
max_active_done: number;
store_identity: string;
}

function attemptPath(root: string, goalId: string, role: Role): string {
return join(shadowManagementDirectory(root, goalId), `archive-${role}.json`);
}

async function readAttempt(root: string, goalId: string, role: Role): Promise<ArchiveAttempt | null> {
let value: unknown;
try {
value = JSON.parse(await readFile(attemptPath(root, goalId, role), "utf8"));
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return null;
throw error;
}
const slot = canonicalAuthorityObject(value, "archive attempt slot");
if (!hasExactAuthorityKeys(slot, ["schema_version", "goal_id", "role", "attempt"]) ||
slot.schema_version !== ATTEMPT_SCHEMA || slot.goal_id !== goalId || slot.role !== role) {
throw new Error("archive attempt slot identity mismatch");
}
if (slot.attempt === null) return null;
const attempt = canonicalAuthorityObject(slot.attempt, "archive attempt");
if (!hasExactAuthorityKeys(attempt, ["operation_id", "expected_provider_revision",
"max_active_done", "store_identity"]) ||
!Number.isSafeInteger(attempt.max_active_done) || Number(attempt.max_active_done) < 0) {
throw new Error("invalid archive attempt");
}
return {
operation_id: requireAuthorityStoreId(attempt.operation_id, "archive operation id"),
expected_provider_revision: attempt.expected_provider_revision === null ? null :
requireAuthorityStoreId(attempt.expected_provider_revision, "archive expected provider revision"),
max_active_done: Number(attempt.max_active_done),
store_identity: requireAuthorityStoreId(attempt.store_identity, "archive store identity"),
};
}

async function writeAttempt(
root: string, goalId: string, role: Role, attempt: ArchiveAttempt | null,
): Promise<void> {
// A durable null tombstone avoids a delete/parent-fsync recovery window.
// There are at most two bounded slots per goal, regardless of history size.
await durableWriteJson(attemptPath(root, goalId, role), {
schema_version: ATTEMPT_SCHEMA, goal_id: goalId, role,
attempt: attempt === null ? null : {...attempt},
});
}

/** Caller holds the goal's shadow-maintenance mutex throughout this operation. */
export async function executeLocalArchiveAttempt(
store: AuthorityStore, root: string, input: CoordinationTodoArchiveInput,
): Promise<JsonObject> {
if (input.dry_run) return executeCoordinationTodoArchiveCompleted(store, input);
let attempt = await readAttempt(root, input.goal_id, input.role);
if (attempt === null) {
// A retained v0 receipt may predate observed-head binding. Preserve the
// actual caller fields; adding today's revision would change its identity.
const receipt = await store.readReceipt(input.operation_id);
if (receipt.status === "missing") {
const head = await store.loadAuthority();
if (head.status !== "loaded") return {
schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...head, changed: false,
};
} else if (receipt.status !== "found") return {
schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...receipt, changed: false,
};
}
const identity = await store.storeIdentity();
if (identity.status !== "available") return {
schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA, ...identity, changed: false,
};
if (attempt !== null && attempt.store_identity !== identity.store_identity) {
throw new Error("archive attempt belongs to a different authority store");
}
if (attempt !== null && attempt.max_active_done !== input.max_active_done) return {
schema_version: COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA,
status: "conflict", changed: false, conflict_kind: "archive_attempt_pending",
operation_id: attempt.operation_id,
reason: "complete the pending archive projection before changing the retention limit",
};
if (attempt === null) {
attempt = {
operation_id: input.operation_id,
expected_provider_revision: input.expected_provider_revision ?? null,
max_active_done: input.max_active_done, store_identity: identity.store_identity,
};
// Persist correlation before the provider can commit. The provider alone
// owns selection, request identity, CAS, and historical result receipts.
await writeAttempt(root, input.goal_id, input.role, attempt);
}
const result = await executeCoordinationTodoArchiveCompleted(store, {
...input, operation_id: attempt.operation_id,
expected_provider_revision: attempt.expected_provider_revision ?? undefined,
});
const identityRejected = result.status === "failed" &&
result.reason_code === "coordination_operation_identity_mismatch";
const beforeCommitRejected = result.status === "conflict" || (
result.status === "failed" && ["invalid_coordination_todo_archive",
"invalid_coordination_projection"].includes(String(result.reason_code))
);
if (result.status === "no_change" || identityRejected || (beforeCommitRejected &&
(await store.readReceipt(attempt.operation_id)).status === "missing")) {
// No committed mutation exists to project. In particular, retain the
// established zero-revision-change semantics of no-change archive calls.
await writeAttempt(root, input.goal_id, input.role, null);
}
return {...result, operation_id: attempt.operation_id};
}

/** A successful projection acknowledges exactly one attempt, never a newer one. */
export async function acknowledgeLocalArchiveAttempt(
store: AuthorityStore, root: string, goalId: string, role: Role, operationId: string,
): Promise<JsonObject> {
const result = {schema_version: LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA,
goal_id: goalId, role, operation_id: operationId};
const attempt = await readAttempt(root, goalId, role);
if (attempt === null) return {...result, status: "no_change", changed: false};
if (attempt.operation_id !== operationId) return {...result, status: "stale", changed: false};
const identity = await store.storeIdentity();
if (identity.status !== "available") return {...result, ...identity, changed: false};
if (attempt.store_identity !== identity.store_identity) {
throw new Error("archive acknowledgement belongs to a different authority store");
}
const receipt = await store.readReceipt(operationId);
const original = receipt.status === "found" ? receipt.receipts[0] : undefined;
if (receipt.status !== "found" || receipt.receipts.length !== 1 ||
original?.schema_version !== COORDINATION_TODO_ARCHIVE_RECEIPT_SCHEMA ||
original.operation_id !== operationId || original.goal_id !== goalId ||
original.request_sha256 !== canonicalAuthoritySha256({
goal_id: goalId, role, max_active_done: attempt.max_active_done,
dry_run: false,
...(attempt.expected_provider_revision === null ? {} : {
expected_provider_revision: attempt.expected_provider_revision,
}),
})) return {
...result, status: "failed", changed: false,
reason_code: "archive_ack_receipt_unavailable",
reason: "archive projection acknowledgement requires its committed receipt",
};
await writeAttempt(root, goalId, role, null);
return {...result, status: "acknowledged", changed: true};
}
58 changes: 53 additions & 5 deletions loopx/control_plane/coordination/local_authority_runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,9 +55,13 @@ import {
import {
COORDINATION_TODO_ARCHIVE_RESULT_SCHEMA,
COORDINATION_TODO_TERMINAL_LIFECYCLE_RESULT_SCHEMA,
executeCoordinationTodoArchiveCompleted,
executeCoordinationTodoTerminalLifecycle,
} from "./todo_terminal_lifecycle.ts";
import {
acknowledgeLocalArchiveAttempt,
executeLocalArchiveAttempt,
LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA,
} from "./local_archive_attempt.ts";
import { editCoordinationTodo, TODO_COMPATIBILITY_EDIT_RESULT_SCHEMA } from "./todo_compatibility_edit.ts";
import {
normalizeIdempotencyKey,
Expand All @@ -73,6 +77,8 @@ export const LOCAL_COORDINATION_TODO_TERMINAL_LIFECYCLE_REQUEST_SCHEMA =
"loopx_local_coordination_todo_terminal_lifecycle_request_v0";
export const LOCAL_COORDINATION_TODO_ARCHIVE_REQUEST_SCHEMA =
"loopx_local_coordination_todo_archive_request_v0";
export const LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA =
"loopx_local_coordination_todo_archive_ack_request_v0";
export {
LOCAL_COORDINATION_MUTATION_REQUEST_SCHEMA,
LOCAL_COORDINATION_MUTATION_RESULT_SCHEMA,
Expand Down Expand Up @@ -155,6 +161,11 @@ function requiredNonNegativeSafeInteger(value: unknown, label: string): number {
return value as number;
}

function archiveRole(value: unknown): "agent" | "user" {
if (value !== "agent" && value !== "user") throw new TypeError("unsupported archive role");
return value;
}

function optionalNonNegativeSafeInteger(value: unknown, label: string): number | null {
return value === null || value === undefined
? null
Expand Down Expand Up @@ -882,16 +893,23 @@ export async function archiveLocalCoordinationTodos(
input.max_active_done,
"max_active_done",
);
const role = archiveRole(input.role);
const operationId = requireAuthorityStoreId(input.operation_id, "operation id");
const expectedRevision = input.expected_provider_revision === undefined ? undefined :
requireAuthorityStoreId(input.expected_provider_revision, "expected provider revision");
if (typeof input.dry_run !== "boolean") throw new TypeError("dry_run must be a boolean");
const now = claimObservedAt(input.observed_at);
return await withCanonicalWriter(root, goalId, input.dry_run === true, async () => {
const store = dependencies.createStore?.(authorityDirectory(root), goalId) ??
new FileAuthorityStore(authorityDirectory(root), goalId);
return {...await executeCoordinationTodoArchiveCompleted(store, {
return {...await executeLocalArchiveAttempt(store, root, {
goal_id: goalId,
role: requireAuthorityStoreId(input.role, "role") as "agent" | "user",
role,
max_active_done: maxActiveDone,
operation_id: requireAuthorityStoreId(input.operation_id, "operation id"),
operation_id: operationId,
expected_provider_revision: expectedRevision,
dry_run: input.dry_run as boolean,
now: claimObservedAt(input.observed_at),
now,
}), ...providerEvidence};
});
} catch (error) {
Expand All @@ -904,6 +922,36 @@ export async function archiveLocalCoordinationTodos(
}
}

/** Retire one local retry identity only after its compatibility projection succeeds. */
export async function acknowledgeLocalCoordinationTodoArchive(
value: unknown,
dependencies: LocalAuthorityRuntimeDependencies = {},
): Promise<JsonObject> {
try {
const input = requireJsonObject(value, "local Todo archive acknowledgement");
if (input.schema_version !== LOCAL_COORDINATION_TODO_ARCHIVE_ACK_REQUEST_SCHEMA) {
throw new TypeError("local Todo archive acknowledgement schema mismatch");
}
const root = runtimeRoot(input.runtime_root);
const goalId = requireAuthorityStoreId(input.goal_id, "goal id");
const role = archiveRole(input.role);
const operationId = requireAuthorityStoreId(input.operation_id, "operation id");
return await withCanonicalWriter(root, goalId, false, async () => {
const store = dependencies.createStore?.(authorityDirectory(root), goalId) ??
new FileAuthorityStore(authorityDirectory(root), goalId);
return acknowledgeLocalArchiveAttempt(store, root, goalId, role, operationId);
});
} catch (error) {
return {
schema_version: LOCAL_TODO_ARCHIVE_ACK_RESULT_SCHEMA,
status: "failed", changed: false,
reason_code: error instanceof ShadowManagementError ? error.reason_code :
"invalid_local_coordination_todo_archive_ack_request",
reason: error instanceof Error ? error.message : "invalid archive acknowledgement",
};
}
}

/** Embedded file adapter; no Markdown input or projection write is accepted. */
export async function editLocalCoordinationTodo(
value: unknown,
Expand Down
Loading
Loading