diff --git a/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.md b/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.md new file mode 100644 index 000000000..92fa0f293 --- /dev/null +++ b/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.md @@ -0,0 +1,91 @@ +# Design follow-up: Goal continuity across restart and replacement + +- **Status:** Deferred design note; non-normative, not an accepted new runtime contract. +- **Origin:** Retains useful questions from [Duang777's #5169](https://github.com/loopx-project/loopx/pull/5169), with its implementation and evidence claims narrowed during review. +- **Language mirror:** [中文](goal-immutability-coherence-defense-v0.zh-CN.md). +- **Owning contracts:** [Goal instance identity and orphan recovery](goal-instance-identity-and-orphan-recovery-v0.md), [Goal direction baseline](goal-direction-baseline-v0.md), [governed amendment](shared-goal-alignment-and-governed-amendment-v0.md), [semantic handoff](capable-manager-semantic-handoff-v0.md), and [shared authority](shared-goal-authority-state-provider-v0.md). + +This preserves follow-up design value from the original “Goal Immutability as +Coherence Defense” draft. It creates no second roadmap, state owner, acceptance +gate or delivery claim. The owning RFCs decide activation, authority and rollout. +The questions below are proposed qualification scenarios, not evidence of a +remaining defect in every named path. + +## What is worth retaining + +Durable commitments should outlive a model's working context. A restarted or +replaced Agent should recover the authorized Goal, constraints, accepted work +and unresolved obligations from their existing owners. A same-name replacement +Goal must not inherit an old instance's authority merely because names match. +These are useful long-horizon failure scenarios even when individual storage +and command tests pass. + +Three distinctions prevent an overly broad “coherence” guarantee: + +| Fact | What it can establish | What it cannot establish | +| --- | --- | --- | +| Exact GoalRef and source-owned instance fence | Which Goal lifetime may admit an action | Whether that action is useful or its output correct | +| Provider revision / CAS | Whether a new write still has its expected storage basis | Current Goal authority if the caller resolves/rebinds the wrong instance; revision tokens are opaque, not ordered counters | +| Operation identity and verified original receipt | Which operation already committed and its original result | Permission to repeat an external effect or attach the result to a replacement Goal | +| Authorized intent / acceptance basis | Which constraints and completion criteria govern this work | Model compliance or outcome correctness without independent evidence | + +Immutable **instance identity** does not mean immutable **Goal intent**. Authorized +amendments must remain possible and versioned through their existing owner. +A model can produce a wrong change against a perfectly current CAS revision. +Prompt/context improvements, typed constraints and outcome validation complement +storage fences; none substitutes for all the others. + +## Proposed follow-up slices under existing owners + +| Slice and owner | Real caller scenario | Decisive acceptance, including recovery | +| --- | --- | --- | +| Instance-qualified continuity — Goal instance RFC; related collaboration/session consumers in [#5106](https://github.com/loopx-project/loopx/pull/5106) and [#5130](https://github.com/loopx-project/loopx/pull/5130) | Retire A through the authorized lifecycle, create same-name B, then deliver A's delayed Todo/result, claim renewal and plan confirmation through their real entrypoints | No mutation or execution authority leaks into B. Typed stale-instance outcomes remain observable; B's legitimate work succeeds. Historical A receipts remain attributable to A where retention/access policy permits. Registry activation and its legacy/off behavior follow the owning RFC. | +| Constraint continuity — direction-baseline and governed-amendment RFCs; roadmap R4 | Resume/rebind an Agent after context loss with stale material/acceptance basis, then repeat with an authorized amendment and refreshed basis | Original constraints and accepted work are recovered from canonical owners; re-evaluation remains Agent-scoped. Unrelated work is not globally blocked. The legitimate amendment can progress; no implicit freeze of all Goal intent. | +| Recoverable late-result disposition — handoff and Effect recovery owners; roadmap R3 | An old request's result arrives after requester/instance replacement or after an external effect has committed but its response was lost | Preserve original request/result lineage and the external effect's durable evidence. Reconcile at the owning ledger; do not silently discard evidence, automatically rebind to B or rerun the effect. An authorized recovery path returns a result or records an explicit terminal disposition. | + +Before implementing a slice, inventory current main, related PRs and existing +fixtures. Extend the current owner's missing cases rather than creating a +parallel “semantic certificate” or generic coherence engine. Shared decisions +belong in existing typed TS owners; provider adapters supply physical evidence. +As of this note's 2026-09-27 review, #5106, #5130 and the related App Turn recovery +[#5139](https://github.com/loopx-project/loopx/pull/5139) are open; their merge or +isolated tests alone would not certify the combined journeys above. + +## Qualification method and unresolved decisions + +Use disposable runtimes, synthetic public-safe Goals and actual supported +backends. Derive expected outcomes from the owning contract before executing: + +- Exercise same-instance restart, same-name replacement, authorized amendment, + delayed input, overlapping invalid conditions and valid post-recovery work. + A stale rejection alone is not restored progress. +- Exercise both accidental scope expansion and escape: covered old bindings, + unrelated current work and newly created subjects follow the declared scope. +- For concurrent creation, preserve the registry's declared uniqueness and + linearization contract. Do not assume that every racing request must create + a separate writable Goal; distinct successful lifetimes must never share an ID. +- Inject one fault at a time and demonstrate oracle sensitivity: wrong GoalRef, + missing commit fence, dropped handoff constraint or duplicated effect. Keep + real effect evidence distinct from simulated adapters and model evaluations. + +Open decisions belong to the existing owners: whether each producer already +captures sufficient immutable instance/basis evidence; how a stale requester +receives a recoverable outcome; and whether an explicit new producer/schema is +needed. Never derive the original instance from a mutable “current Goal” lookup. +If a format change is necessary, qualify backup/migration and mixed writers; +this note does not predeclare “no migration needed.” + +## Delivered boundary and evidence status + +[#5169's operation replay](../../reference/authority-operation-replay.md) verifies +matching historical File/SQLite commits without rewinding current state. It does +not implement or qualify the three follow-up journeys. Tests of stale provider +revisions are not tests of Goal replacement, compaction or semantic correctness. + +The original draft's quantitative research/experiment tables are not retained as +accepted evidence: this PR does not supply an independently reviewable public +harness, oracle and source provenance for them. Future evidence must name the +exact revision, real entrypoint/backend, fault, independent oracle and recovery +readback. No numeric success rate or claim that “CAS prevents coherence collapse” +is carried forward. This note neither changes File/SQLite defaults nor adds a +new prerequisite to their existing D1–D3 qualification. diff --git a/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.zh-CN.md b/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.zh-CN.md new file mode 100644 index 000000000..db65a9a19 --- /dev/null +++ b/docs/architecture/rfcs/goal-immutability-coherence-defense-v0.zh-CN.md @@ -0,0 +1,74 @@ +# 后续设计:重启与实例替换中的 Goal 连续性 + +- **状态:** 延后设计记录;非规范性内容,不是新接受的运行时契约。 +- **来源:** 保留 [Duang777 在 #5169 中提出的问题](https://github.com/loopx-project/loopx/pull/5169),评审时收窄其实现与证据宣称。 +- **语言镜像:** [English](goal-immutability-coherence-defense-v0.md)。 +- **所属契约:** [Goal 实例身份与孤儿恢复](goal-instance-identity-and-orphan-recovery-v0.zh-CN.md)、[Goal 方向基线](goal-direction-baseline-v0.zh-CN.md)、[受治理的修改](shared-goal-alignment-and-governed-amendment-v0.md)、[语义交接](capable-manager-semantic-handoff-v0.zh-CN.md)、[共享权威](shared-goal-authority-state-provider-v0.md)。 + +这里保留原“Goal 不可变性作为一致性防御”草稿中的后续设计价值,不增加第二份 +roadmap、状态 owner、验收门禁或交付宣称。激活、权限和上线由所属 RFC 决定。 +以下是建议补充的验收场景,不表示每条路径当前都存在尚未修复的缺陷。 + +## 值得保留的部分 + +持久工作承诺应当比模型的工作上下文活得更久。Agent 重启或被替换后,应从已有 +owner 恢复获得授权的 Goal、约束、已接受的工作和未完成义务。同名的新 Goal +不能仅因名字相同就继承旧实例的权限。即便单个存储和命令测试已经通过,这些仍是 +有价值的长程运行反例。 + +必须区分以下事实,避免泛化为“语义一致性保证”: + +| 事实 | 能证明什么 | 不能证明什么 | +| --- | --- | --- | +| 精确 GoalRef 与源 registry 拥有的实例 fence | 哪次 Goal 生命周期可以接纳动作 | 动作是否有用、输出是否正确 | +| Provider revision / CAS | 新写入是否仍基于预期的存储版本 | 调用方错误解析或重绑定实例时的当前 Goal 权限;revision 是不透明 token,不是可排序计数器 | +| 操作身份与经过验证的原回执 | 哪个操作已经提交及其原结果 | 重复执行外部效果、或把结果挂到替代 Goal 的权限 | +| 已授权的意图/验收基线 | 当前工作应遵守哪些约束和完成条件 | 缺乏独立证据时,模型确实遵守了约束或结果正确 | + +**实例身份不可变**不等于 **Goal 意图不可修改**。获得授权的修改必须仍能通过现有 +owner 留下版本并生效。模型完全可能基于最新 CAS 版本提交错误改动。 +Prompt/上下文优化、类型化约束与结果验证和存储 fence 相互补充,没有一项可以 +替代全部其他机制。 + +## 归入现有 owner 的后续切片 + +| 切片与 owner | 真实调用场景 | 决定性验收,包含恢复 | +| --- | --- | --- | +| 按实例确认连续性——Goal 实例 RFC;相关 collaboration/session 消费者见 [#5106](https://github.com/loopx-project/loopx/pull/5106)、[#5130](https://github.com/loopx-project/loopx/pull/5130) | 通过获授权的生命周期退役 A,建立同名 B,再经真实入口提交 A 的迟到 Todo/结果、claim 续约、计划确认 | B 不受到错误写入或执行权限污染;过期实例结果可观察,B 的合法工作仍能推进。保留/访问策略允许时,A 的历史回执仍归属于 A。Registry 激活及 legacy/off 行为遵循所属 RFC。 | +| 约束连续性——方向基线、受治理修改 RFC;roadmap R4 | Agent 上下文丢失后以过期材料/验收基线恢复或重绑定,再以已获授权的修改和新基线重复执行 | 从 canonical owner 恢复原约束和已接受工作;重新评估局限于相关 Agent,不把无关工作全部阻塞。合法修改可以推进,不隐式冻结全部 Goal 意图。 | +| 迟到结果的可恢复处置——handoff 与 Effect recovery owner;roadmap R3 | 请求方/实例替换后收到旧请求结果,或外部效果已经提交但响应丢失 | 保留原请求/结果关系和外部效果的持久证据;在所属 ledger 对账,不能静默丢弃证据、自动改挂到 B 或重跑效果。通过获授权的恢复路径返回结果,或明确记录终止处置。 | + +实施任一切片前,先核对当前 main、相关 PR 和已有 fixture。补齐当前 owner 的缺口, +不另建“semantic certificate”或通用 coherence 引擎。共享决策放在已有类型化 TS +owner,provider adapter 提供物理存储证据。本记录于 2026-09-27 核对时,#5106、 +#5130 及相关 App Turn 恢复 [#5139](https://github.com/loopx-project/loopx/pull/5139) +仍开放;仅合并它们或通过各自孤立测试,不代表上述组合链路已经验收。 + +## 验证方法与未决事项 + +使用可丢弃 runtime、公开安全的合成 Goal 和真实受支持 backend;执行前从所属 +契约独立推导预期结果: + +- 覆盖同实例重启、同名替换、获授权的修改、迟到输入、重叠非法条件,以及恢复后 + 的合法工作。只有 stale 拒绝不等于恢复推进。 +- 同时覆盖范围误扩大与逃逸:旧绑定、无关当前工作和新增对象须遵守声明的范围。 +- 并发创建应遵守 registry 的唯一性和线性化契约,不假定每个竞争请求都应建立一份 + 独立可写 Goal;不同的成功生命周期不能共用实例 ID。 +- 每次注入一个故障并验证 oracle 的敏感性:错误 GoalRef、缺失提交 fence、交接约束 + 丢失或重复效果。区分真实效果证据、模拟 adapter 与模型评测。 + +未决事项交给已有 owner:各 producer 是否已捕获足够的不可变实例/基线证据; +过期请求方如何收到可恢复结果;是否确需新的 producer/schema。不能通过可变的 +“当前 Goal”查询倒推出原实例。如果需要格式修改,必须验收备份/迁移和混合版本 +writer;本记录不预先宣称“无需迁移”。 + +## 本次交付边界与证据状态 + +[#5169 的操作重放](../../reference/authority-operation-replay.md) 验证历史 File/SQLite +提交的完整意图,且不会回退当前状态;它没有实现或验收以上三组后续链路。 +过期 provider revision 的测试不能冒充 Goal 替换、上下文压缩或语义正确性测试。 + +原草稿的研究百分比和实验成绩表不作为已接受证据保留:本 PR 没有提供可独立审查的 +公开 harness、oracle 和来源依据。后续证据须明确精确 revision、真实入口/backend、 +故障、独立 oracle 和恢复读回。不沿用数值成功率,也不沿用“CAS 防止语义崩塌”的 +结论。本记录不改变 File/SQLite 默认值,也不给现有 D1–D3 验收添加新前置条件。 diff --git a/docs/reference/authority-operation-replay.md b/docs/reference/authority-operation-replay.md new file mode 100644 index 000000000..1bc7e8b5d --- /dev/null +++ b/docs/reference/authority-operation-replay.md @@ -0,0 +1,81 @@ +# Authority operation replay + +File and SQLite accept a direct retry of an already committed operation when +its complete canonical body matches the original transaction. The body is +`next_projection`, `events` and `receipts`; JSON object key order is not intent. +`operation_id` selects that transaction within the opened Goal store. + +This supports retry after a lost response. It does not create another state +transition, append events again, restore an old head or authorize another +external effect. If A committed, then B committed, replaying A returns A's +original cursor/provider revision while B remains the current head. + +## Commit and recovery contract + +| Case | File / SQLite | PostgreSQL / NoKV | +| --- | --- | --- | +| New operation, current CAS basis | Commit atomically | Commit atomically | +| New operation, stale CAS basis | `provider_revision_mismatch` | `provider_revision_mismatch` | +| Existing operation, identical body | `applied` with original cursor/revision | Ordinary commit remains a conflict; recover via `readReceipt` | +| Existing operation, different body | `operation_id_exists` | Conflict; normal revision-check precedence remains | + +For a verified historical retry, File/SQLite do not require the caller's CAS +basis to remain current: they are returning a historical fact, not admitting a +new write. Validation, store identity/existing-only admission and the current +store integrity checks still apply. An invalid request fails before replay. +A matching operation with a changed projection, event or receipt is never +acknowledged as the original commit. + +File reconstructs the original projection using the retained journal. SQLite +verifies the original checkpoint/delta window in the **same write transaction** +before acknowledging a replay. Comparing the caller to a stored digest alone +is insufficient: a damaged retained receipt/event must fail its own proof. +No full-history audit is added to ordinary retry. Digests detect inconsistent +bytes, not an administrator who rewrites both data and proof. + +`CoordinationCommandReceipt` remains the provider-neutral business recovery +owner. It reads and validates the original command receipt even when a provider +returns `applied`, and reconciles conflict or ambiguous responses. Do not remove +that readback: `applied` can describe a historical commit, and the other +providers retain their existing direct-commit behavior. NoKV's ambiguous-write +readback recovery is distinct from its ordinary `commitAuthority` contract. + +This changes File/SQLite's previous duplicate-commit rejection behavior. +Concurrent identical lease renewals can now both report `applied`, while only +one renewal/version transition is persisted. Business recovery still validates +request identity; a historical receipt never grants current lease authority. +There is no feature flag, new request field, storage migration or frontend +configuration. Reverting requires reverting the provider behavior, not merely +removing tests. Existing durable receipts keep their original format. + +## Scope and verification + +The shared-authority RFC owns this storage/recovery boundary. Goal lifetime +identity, lease epochs and provider revisions remain separate contracts. +This change does not implement Goal replacement isolation, semantic correctness +of model output, default provider activation or long-horizon qualification. +The [deferred Goal continuity note](../architecture/rfcs/goal-immutability-coherence-defense-v0.md) +preserves related restart, instance-replacement and constraint-recovery scenarios +under their existing RFC owners; those scenarios are not qualified by this PR. + +`authority_operation_replay_conformance.ts` runs on both File and SQLite. It +checks full body drift, canonical key ordering, historical replay and concurrent +same-operation attempts against complete head, receipt and history readback. +SQLite adds retained-row corruption cases that must refuse replay without +changing durable rows. Existing real-process lease renewal and command receipt +suites cover the business entrypoint above these providers. + +## 中文说明 + +File/SQLite 现在可以直接重试已成功提交的操作:操作 ID 相同,而且完整的 +projection、events、receipts 相同,才返回原 cursor/revision。A 提交后 B 又提交, +重试 A 只确认 A 的历史结果,不把当前状态退回 A,也不再追加事件。 + +这不是绕过新写入的 CAS。新操作仍须匹配当前版本;历史重试则须证明原事务。 +SQLite 在同一事务内校验对应 checkpoint/delta 窗口,不能只比较数据库中保存的 +摘要字段。历史回执或事件损坏时必须拒绝,不能报告成功。 + +上层 `CoordinationCommandReceipt` 仍须读回并校验业务回执,处理响应丢失及不确定 +提交;PostgreSQL/NoKV 的普通重复提交仍返回冲突。并发同意图续约可能从过去的 +`applied/recovered` 变成 `applied/applied`,但实际只写入一次。历史成功不授予当前 +lease 执行权。此变更没有新配置或存储格式,也不证明 Goal 实例隔离或模型语义正确。 diff --git a/docs/reference/sqlite-authority-store.md b/docs/reference/sqlite-authority-store.md index 1452f7279..4ae0c8781 100644 --- a/docs/reference/sqlite-authority-store.md +++ b/docs/reference/sqlite-authority-store.md @@ -63,6 +63,10 @@ rotation, corruption repair, or network-filesystem sharing is supported. ## Read integrity +Direct retries follow the [authority operation replay contract](authority-operation-replay.md): +matching full intent returns its verified original position without a new write. + + Authority reads share one SQLite snapshot, and writes run the same live proof inside their transaction before publishing a new commit row. The proof is layered so that each layer pays only for what it returns: diff --git a/loopx/control_plane/coordination/file_authority_store.ts b/loopx/control_plane/coordination/file_authority_store.ts index 1bd95ef72..0fc1be38c 100644 --- a/loopx/control_plane/coordination/file_authority_store.ts +++ b/loopx/control_plane/coordination/file_authority_store.ts @@ -367,6 +367,39 @@ export class FileAuthorityStore implements AuthorityStore { if (this.existingOnly && current === null) { return { status: "failed", reason_code: "existing_authority_missing", reason: "existing-only store cannot bootstrap a missing authority" }; } + // Content-aware idempotency: same operation_id with matching full body + // ({events, receipts, projection}) returns the original receipt + // (retry-after-crash); the check must precede the revision gate so a + // stuck writer can't block the already-committed replay. A different + // body (including a projection-only drift) remains a conflict. + if (current !== null) { + const existing = current.receipt(normalized.operation_id); + if (existing) { + const cursorIndex = Number(existing.cursor) - 1; + const [historical] = current.scan(cursorIndex, 1); + if (!historical) { + return {status: "failed", reason_code: "provider_protocol_violation", + reason: `receipt entry cursor ${existing.cursor} missing from journal scan`}; + } + const intendedBody = {events: normalized.events, receipts: normalized.receipts, + projection: normalized.next_projection}; + const existingBody = {events: existing.events, receipts: existing.receipts, + projection: historical.projection}; + if (canonicalAuthorityBytes(intendedBody).equals(canonicalAuthorityBytes(existingBody))) { + return { + status: "applied", + provider_revision: existing.provider_revision, + cursor: existing.cursor, + }; + } + return { + status: "conflict", + conflict_kind: "operation_id_exists", + current_provider_revision: current.provider_revision, + current_cursor: current.cursor, + }; + } + } if ((current?.provider_revision ?? null) !== normalized.expected_provider_revision) { return { status: "conflict", @@ -375,14 +408,6 @@ export class FileAuthorityStore implements AuthorityStore { current_cursor: current?.cursor ?? null, }; } - if (current?.receipt(normalized.operation_id)) { - return { - status: "conflict", - conflict_kind: "operation_id_exists", - current_provider_revision: current.provider_revision, - current_cursor: current.cursor, - }; - } const document = FileAuthorityJournal.append(current, this.goalId, identity, normalized, (previous, transaction) => fileAuthorityRevision(this.goalId, identity, previous, transaction)); const {cursor, provider_revision: revision} = document; diff --git a/loopx/control_plane/coordination/sqlite_authority_store.ts b/loopx/control_plane/coordination/sqlite_authority_store.ts index ceabac52f..d7746c50e 100644 --- a/loopx/control_plane/coordination/sqlite_authority_store.ts +++ b/loopx/control_plane/coordination/sqlite_authority_store.ts @@ -463,10 +463,26 @@ export class SqliteAuthorityStore implements AuthorityStore { const cursor = current?.state.cursor ?? null; const revision = current?.provider_revision ?? null; let conflict: "provider_revision_mismatch" | "operation_id_exists" | null = null; - if (revision !== normalized.expected_provider_revision) conflict = "provider_revision_mismatch"; - else if (db.prepare("SELECT 1 FROM commits WHERE operation_id = ?").get(normalized.operation_id)) { + // Replay proves the retained transaction in this same write snapshot. + // A matching stored digest alone is not evidence that its row is intact. + const existingRow = db.prepare(`SELECT ${COMMIT_COLUMNS} FROM commits WHERE operation_id = ?`) + .get(normalized.operation_id); + if (existingRow) { + const retained = this.decodeCommitRow(existingRow); + const window = this.verifiedRange(db, retained.cursor, retained.cursor); + const original = window.transactions[0]; + if (!original || original.operation_id !== normalized.operation_id) { + protocol("SQLite replay is not part of its retained window"); + } + const digest = commitDigest(window.identity, retained.cursor, normalized.operation_id, + normalized.next_projection, normalized.events, normalized.receipts); + if (digest === retained.commit_digest) { + db.exec("ROLLBACK"); transactionOpen = false; + return {status: "applied", provider_revision: original.provider_revision, cursor: original.cursor}; + } conflict = "operation_id_exists"; } + if (!conflict && revision !== normalized.expected_provider_revision) conflict = "provider_revision_mismatch"; if (conflict) { db.exec("ROLLBACK"); transactionOpen = false; return {status: "conflict", conflict_kind: conflict, current_provider_revision: revision, diff --git a/tests/control_plane_ts/authority_operation_replay_conformance.ts b/tests/control_plane_ts/authority_operation_replay_conformance.ts new file mode 100644 index 000000000..f6f4db868 --- /dev/null +++ b/tests/control_plane_ts/authority_operation_replay_conformance.ts @@ -0,0 +1,65 @@ +/** Direct local-provider replay: full intent matches a verified historical commit. + * This is not Goal-instance isolation or permission to execute another effect. */ +import assert from "node:assert/strict"; +import test from "node:test"; +import type {AuthorityStoreConformanceFactory} from "./authority_store_conformance.ts"; +import {authorityStoreCommitFixture as commit} from "./authority_store_conformance.ts"; + +export function registerAuthorityOperationReplayConformance( + provider: string, factory: AuthorityStoreConformanceFactory, +): void { + test(`${provider}: historical operation replay preserves full intent and later state`, async t => { + const {store, contender} = await factory(t); + const input = commit(null, "original", 1, 1); + input.next_projection.metadata = {labels: ["retained", "中"], nested: {value: 7}}; + const first = await store.commitAuthority(input); + assert.equal(first.status, "applied"); if (first.status !== "applied") return; + const second = await store.commitAuthority(commit(first.provider_revision, "later", 2, 2)); + assert.equal(second.status, "applied"); if (second.status !== "applied") return; + const snapshot = async () => ({head: await contender.loadAuthority(), + receipt: await contender.readReceipt("original"), history: await contender.scanCommitted(null, 10)}); + const before = await snapshot(); + for (const basis of [null, second.provider_revision]) { + for (const change of ["none", "key_order", "projection", "events", "receipts"] as const) { + const replay = structuredClone(input); + replay.expected_provider_revision = basis; + if (change === "key_order") replay.next_projection = Object.fromEntries(Object.entries(replay.next_projection).reverse()); + if (change === "projection") replay.next_projection.metadata = {labels: ["different"]}; + if (change === "events") replay.events = [{type: "different"}]; + if (change === "receipts") replay.receipts = [{result: "different"}]; + const result = await contender.commitAuthority(replay); + if (change === "none" || change === "key_order") assert.deepEqual(result, first); + else { + assert.equal(result.status, "conflict", `${change}-only drift must be rejected`); + if (result.status === "conflict") assert.equal(result.conflict_kind, "operation_id_exists"); + } + assert.deepEqual(await snapshot(), before, `${change} with basis ${basis} changed committed state`); + } + } + }); + + for (const matching of [true, false]) { + test(`${provider}: concurrent ${matching ? "matching" : "different"} intent never double commits`, async t => { + const {store, contender} = await factory(t); + const first = commit(null, "racing-operation", 1, 1); + const second = structuredClone(first); + if (!matching) second.next_projection.authority_revision = 2; + const results = await Promise.all([store.commitAuthority(first), contender.commitAuthority(second)]); + const applied = results.filter(result => result.status === "applied"); + assert.equal(applied.length, matching ? 2 : 1); + if (matching) assert.deepEqual(results[0], results[1]); + else { + const rejected = results.find(result => result.status === "conflict"); + assert.equal(rejected?.conflict_kind, "operation_id_exists"); + } + assert.deepEqual(await store.loadAuthority(), await contender.loadAuthority()); + const page = await contender.scanCommitted(null, 10); + assert.equal(page.status, "page"); if (page.status !== "page") return; + assert.equal(page.transactions.length, 1); + const winner = results[0]!.status === "applied" ? first : second; + assert.deepEqual(page.transactions[0]!.projection, winner.next_projection); + assert.deepEqual(page.transactions[0]!.receipts, winner.receipts); + assert.deepEqual(page.transactions[0]!.events, winner.events); + }); + } +} diff --git a/tests/control_plane_ts/authority_store.test.ts b/tests/control_plane_ts/authority_store.test.ts index eaf4c7b1c..318743a33 100644 --- a/tests/control_plane_ts/authority_store.test.ts +++ b/tests/control_plane_ts/authority_store.test.ts @@ -13,6 +13,7 @@ import { authorityStoreCommitFixture as commit, registerAuthorityStoreConformance, } from "./authority_store_conformance.ts"; +import {registerAuthorityOperationReplayConformance} from "./authority_operation_replay_conformance.ts"; async function fixture(t: test.TestContext, goalId = "goal-a") { const root = await mkdtemp(join(tmpdir(), "loopx-authority-store-")); @@ -23,6 +24,11 @@ async function fixture(t: test.TestContext, goalId = "goal-a") { registerAuthorityStoreConformance("file provider", async (t) => { const { root, store } = await fixture(t); return { store, contender: new FileAuthorityStore(root, "goal-a") }; +}, "applied"); + +registerAuthorityOperationReplayConformance("file provider", async (t) => { + const { root, store } = await fixture(t); + return { store, contender: new FileAuthorityStore(root, "goal-a") }; }); test("file provider persists object keys in deterministic Unicode order", async (t) => { diff --git a/tests/control_plane_ts/authority_store_conformance.ts b/tests/control_plane_ts/authority_store_conformance.ts index d47e5b62b..eb4dc2b25 100644 --- a/tests/control_plane_ts/authority_store_conformance.ts +++ b/tests/control_plane_ts/authority_store_conformance.ts @@ -250,6 +250,7 @@ async function withConcurrentAuthorityReads(backends: readonly AuthorityStore export function registerAuthorityStoreConformance( providerName: string, factory: AuthorityStoreConformanceFactory, + matchingReplayStatus: "applied" | "conflict" = "conflict", ): void { test(`${providerName} conformance: captured complete source survives provider reopen and readback`, async (t) => { const {store, contender} = await factory(t); @@ -274,7 +275,7 @@ export function registerAuthorityStoreConformance( if (result.status !== "applied") return; const replay = await store.commitAuthority({expected_provider_revision: result.provider_revision, operation_id: "capture-source", events: [], next_projection: captured, receipts: []}); - assert.equal(replay.status, "conflict"); + assert.equal(replay.status, matchingReplayStatus); if (replay.status === "conflict") assert.equal(replay.conflict_kind, "operation_id_exists"); const retained = await contender.readReceipt("capture-source"); assert.equal(retained.status, "found"); 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 143966299..37f81694a 100644 --- a/tests/control_plane_ts/canonical_task_lease_renew.test.ts +++ b/tests/control_plane_ts/canonical_task_lease_renew.test.ts @@ -383,7 +383,7 @@ for (const provider of ["file", "sqlite"] as const) { 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"]); + assert.deepEqual(results.map(r => r.status).sort(), differentIntent ? ["applied", "failed"] : ["applied", "applied"]); 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/sqlite_authority_store.test.ts b/tests/control_plane_ts/sqlite_authority_store.test.ts index cb603deff..6d2f9dd62 100644 --- a/tests/control_plane_ts/sqlite_authority_store.test.ts +++ b/tests/control_plane_ts/sqlite_authority_store.test.ts @@ -11,12 +11,15 @@ import { AUTHORITY_STATE_CHECKPOINT_INTERVAL } from "../../loopx/control_plane/c import { canonicalAuthorityBytes } from "../../loopx/control_plane/coordination/authority_store_codec.ts"; import { authorityStoreCommitFixture, registerAuthorityStoreConformance } from "./authority_store_conformance.ts"; +import {registerAuthorityOperationReplayConformance} from "./authority_operation_replay_conformance.ts"; + async function fixture(t: test.TestContext) { const directory = await mkdtemp(join(tmpdir(), "sqlite-authority-")); t.after(() => rm(directory, {recursive: true, force: true})); return {store: new SqliteAuthorityStore(directory, "goal"), contender: new SqliteAuthorityStore(directory, "goal")}; } -registerAuthorityStoreConformance("SQLite", fixture); +registerAuthorityStoreConformance("SQLite", fixture, "applied"); +registerAuthorityOperationReplayConformance("SQLite", fixture); test("SQLite commits and reads back every JSON object key", {timeout: 30000}, async t => { const {store} = await fixture(t); @@ -418,3 +421,32 @@ test("SQLite receipt batch bounds are checked before opening storage", async t = const {existsSync} = await import("node:fs"); assert.equal(existsSync(store.path), false); }); + +for (const cursor of [1, 2]) for (const field of ["events", "receipts"] as const) { + test(`SQLite replay refuses corrupt retained ${field} at cursor ${cursor}`, async t => { + const {store} = await fixture(t); + let revision: string | null = null; + const inputs = []; + for (let index = 1; index <= 3; index++) { + const input = authorityStoreCommitFixture(revision, `replay-${index}`, index, index); + inputs.push(input); + const committed = await store.commitAuthority(input); + assert.equal(committed.status, "applied"); if (committed.status !== "applied") return; + revision = committed.provider_revision; + } + const {DatabaseSync} = createRequire(import.meta.url)("node:sqlite"); + const db = new DatabaseSync(store.path); + try { + db.prepare(`UPDATE commits SET ${field}=? WHERE cursor=?`).run('[{"forged":true}]', cursor); + const before = db.prepare("SELECT * FROM commits ORDER BY cursor").all(); + assert.equal((await store.loadAuthority()).status, "loaded", "current head remains valid"); + assert.equal((await store.readReceipt(`replay-${cursor}`)).status, "failed"); + const replay = await store.commitAuthority(inputs[cursor - 1]!); + assert.equal(replay.status, "failed", "stored digest alone cannot prove the retained transaction"); + if (replay.status === "failed") assert.equal(replay.reason_code, "provider_protocol_violation"); + assert.deepEqual(db.prepare("SELECT * FROM commits ORDER BY cursor").all(), before); + const head = await store.loadAuthority(); + assert.equal(head.status, "loaded"); if (head.status === "loaded") assert.equal(head.cursor, "3"); + } finally { db.close(); } + }); +}