diff --git a/docs/architecture/runtime-managed-workspace-m2-extraction-ledger.zh-CN.md b/docs/architecture/runtime-managed-workspace-m2-extraction-ledger.zh-CN.md new file mode 100644 index 0000000000..01879eb2e0 --- /dev/null +++ b/docs/architecture/runtime-managed-workspace-m2-extraction-ledger.zh-CN.md @@ -0,0 +1,48 @@ +# Managed Workspace M2 Extraction Ledger + +- 新基线:`upstream/main@32e3cbbd0` +- 历史实现来源:`codex/managed-workspace-mutation-authority-m2@d9ba64697` +- 当前重建分支:`codex/m2-1-main-rebuild`(验证后替换正式 Draft 分支) +- 原则:历史分支只作为测试与实现来源;最终 diff 直接建立在已合入 M1.3 的最新主线上,不带入旧集成栈提交 + +## M2.1 文件归属 + +| 文件 | 归属不变量 | 处理 | +|---|---|---| +| `packages/core/src/workspace-version-authority.ts` | successor fact、strict decoder、causal scanner/head | 从历史提交迁移并按当前 subpath API 重建 | +| `packages/core/src/__tests__/workspace-version-authority.test.ts` | public scanner 行为 | 先迁 RED,再迁最小实现 | +| `packages/storage/src/sqlite-runtime-schema.ts` | successor projection schema | 历史 migration 11 重编号为当前 schema 13;保留 main 的 11/12 | +| `packages/storage/src/sqlite-runtime-store.ts` | T2 + successor + projection + head CAS bundle | 迁移并对齐当前 tool-ledger/continuation imports | +| `packages/storage/src/workspace-version-authority-internal.ts` | store-bound 唯一 writer capability | 迁移;不加入 public package exports | +| `packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts` | atomicity、retry、rollback、corruption、migration | 迁移历史测试并补 stale CAS 与 populated 12→13 fixture | +| `packages/storage/src/__tests__/sqlite-runtime-crash.test.ts` | 真实子进程中 successor transaction 的 kill/reopen 边界 | 扩展现有 crash harness,证明 rollback 与 exact retry | +| `docs/architecture/runtime-managed-workspace-mutation-version-authority-v1.zh-CN.md` | M2.1 实现合同 | 按当前 schema、最新 main 与 PR 切片重写 | +| `docs/architecture/runtime-resume-phase3-phase4-workspace-checkpoint-design.zh-CN.md` | 总路线 | 只更新 M2 子切片,不迁入旧依赖图 | + +## 明确不进入 M2.1 + +| 历史或未来内容 | 归属 | +|---|---| +| `managed-dependency-environment`、bundled npm、worker bridge | 已合入 main 的 M1.3 前置,不从旧 M2 分支重复携带 | +| Git candidate ref/commit、delta/path policy、orphan GC | M2.2 | +| T1 mutation profile/base freeze、owner-bound mutating admission | M2.3 | +| 真实 Write/Edit composition、Host crash/reopen | M2.4 | +| continuation boundary 绑定 workspace version | M3 | +| restore、rebaseline、publish、undo、replication | M4 | + +## Commit 映射 + +| 当前提交 | 来源 | 说明 | +|---|---|---| +| `feat(core): define causal workspace successor facts` | `d9ba64697` 的 core 片段 | 去掉已删除的 root index export,只保留 subpath contract | +| `feat(storage): atomically accept workspace successor versions` | `d9ba64697` 的 storage 片段 | 解决当前接口冲突并把 migration 顺延到 13 | +| 文档提交 | 现行工程实践 | 增加 extraction ledger、平台/失败边界和 M2.2–M2.4 切片 | + +## Diff 审核门槛 + +最终 PR 必须同时满足: + +1. 相对最新 `upstream/main` 只出现本表列出的 M2.1 文件; +2. `git range-diff` 能说明历史 `d9ba64697` 的核心逻辑去向; +3. 不包含 bundled npm、Desktop、CLI、Git candidate 或 Write/Edit production composition; +4. 最终 Draft 必须以最新 `upstream/main` 为直接祖先,并运行 schema/rebuild/crash/concurrency suites。 diff --git a/docs/architecture/runtime-managed-workspace-mutation-version-authority-v1.zh-CN.md b/docs/architecture/runtime-managed-workspace-mutation-version-authority-v1.zh-CN.md new file mode 100644 index 0000000000..871c17abec --- /dev/null +++ b/docs/architecture/runtime-managed-workspace-mutation-version-authority-v1.zh-CN.md @@ -0,0 +1,211 @@ +# Managed Workspace Mutation Version Authority v1 + +- 状态:M2.1 切片实现完成;已从最新 `upstream/main` 平铺重建,保持 Draft 等待 M2.4 生产消费者 +- 更新日期:2026-08-17 +- 主要不变量:工具成功结果、successor workspace fact、version projection 与 canonical head 推进只能全可见或全不可见 +- canonical source:immutable RuntimeEvents +- 写入 owner:storage-internal successor bundle writer +- 原子性边界:单个 SQLite `BEGIN IMMEDIATE ... COMMIT` + +## 1. 本切片解决什么 + +M0 只允许一个 epoch 接受 baseline。M2.1 把 authority spine 扩展成可推进的版本链: + +```text +epoch_opened(seq=1) +baseline_accepted(seq=2, head revision=1) +version_accepted(seq=3, head revision=2) +version_accepted(seq=4, head revision=3) +... +``` + +一次 Write/Edit 只有在以下四部分位于同一个 SQLite transaction 时才可对外宣称成功: + +```text +tool function_response (T2) ++ maka.workspace.version_accepted RuntimeEvent ++ runtime_workspace_versions successor projection ++ runtime_workspace_heads compare-and-set +``` + +任何一项失败都回滚全部写入。不能先提交 T2 再“尽力”更新 workspace head,也不能先推进 head再补工具结果。 + +本切片故意不接真实 Write/Edit:调用者还不能自行声称某个 Git commit/tree 是可信 candidate。这个证明属于 +M2.2 的 Git candidate owner;M2.1 只建立其唯一持久化出口。 + +## 2. Owner、失败状态与回滚 + +| 项目 | 决策 | +|---|---| +| fact contract / pure scanner | `@maka/core/workspace-version-authority` | +| bundle writer | `@maka/storage` 内部 WeakMap capability;不从 package root 导出 | +| canonical evidence | workspace authority facts + successor 引用的 tool call/dispatch/response facts | +| disposable state | `runtime_workspace_versions`、`runtime_workspace_heads` | +| 并发裁判 | SQLite write transaction + base version/event/revision 三元 CAS | +| exact retry | immutable outcome 与 successor fact 精确一致时返回该 operation 原先接受的 head;后续合法 head 推进不改变历史结果 | +| stale writer | base head 任一字段不一致即拒绝,且对应 tool operation 保持 `prepared` | +| corruption | malformed fact、断链、重复 version identity、tool evidence 不匹配全部 fail closed | +| 运行时回滚 | transaction 未提交时 T2/fact/projection/head 全部回滚 | +| 数据升级 | schema 12 baseline rows 原样升级到 schema 13;schema 13 不支持向旧 binary 降级 | + +“唯一 writer”不是注释:普通 RuntimeEvent writer 仍拒绝 workspace fact;successor writer 仅能通过 +`workspace-version-authority-internal.ts` 中与具体 store 实例绑定的 capability 调用。 + +## 3. v1 successor fact + +```ts +{ + kind: 'maka.workspace.version_accepted', + version: 1, + payload: { + protocol: 'workspace_version_accepted_v1', + repositoryId, + workspaceId, + workspaceEpochId, + workspaceVersionId, + objectFormat, + parents: [parentWorkspaceVersionId], + origin: { + kind: 'tool_mutation', + operationId, + dispatchEventId, + outcomeEventId + }, + baseAcceptedEventId, + baseHeadRevision, + commitOid, + treeOid, + policyHash, + treeDeltaDigest, + changedFileCount, + deletedFileCount, + executionProfileDigest + } +} +``` + +Strict decoder 拒绝额外字段、未知 kind/version/protocol、非法 ID/OID/digest、空 parent、多 parent、非安全计数和 +自指 version。 + +Scanner 还必须证明: + +- successor 与 epoch 的 repository/workspace/epoch/object format/policy 完全一致; +- `parents[0]`、`baseAcceptedEventId`、`baseHeadRevision` 同时指向扫描时的 current head; +- authority `event_seq` 连续,不能跨过或重放旧 revision; +- workspace version identity 在整个 authority 中唯一; +- origin 引用的是同一条无 corruption 的 Write/Edit `reconcile` operation; +- `dispatchEventId` 与 `outcomeEventId` 精确指向该 operation 的 immutable T1/T2 facts; +- T2 必须是成功的 `function_response`;`isError: true` 不能创建 successor,也不能通过 rebuild。 + +最后两项在 SQLite canonical reader/rebuild 中对同一 snapshot 的 RuntimeEvents 运行 tool-ledger scanner 后交叉 +验证;不能由 tool projection 或 caller 自报代替。 + +## 4. 原子 bundle 时序 + +```mermaid +sequenceDiagram + participant O as Future Managed Mutation Owner + participant S as SQLite Successor Writer + participant T as Tool RuntimeEvents + participant W as Workspace RuntimeEvents + participant P as Workspace Projections + + O->>S: commitWorkspaceSuccessorInternal(T2, successor) + S->>S: BEGIN IMMEDIATE + S->>S: scan immutable workspace + tool evidence + alt exact bundle already committed + S-->>O: created=false + else stale base / identity drift / corruption + S->>S: ROLLBACK + S-->>O: fail closed + else new successor + S->>T: append function_response and settle operation + S->>W: append version_accepted at next seq + S->>S: rescan canonical facts and evidence + S->>P: insert successor version + S->>P: CAS head(base version/event/revision -> successor) + S->>S: compare projections with canonical scan + S->>S: COMMIT + S-->>O: created=true + new head + end +``` + +## 5. Schema 13 + +Schema 13 保留 schema 12 baseline rows,并扩展 `runtime_workspace_versions`: + +- `origin_kind` 支持 `baseline | tool_mutation`; +- mutation row 持久化 operation/dispatch/outcome、base revision 与 execution profile digest; +- CHECK 约束禁止 baseline row 夹带 mutation 字段,也禁止 mutation row 缺少因果字段; +- `runtime_workspace_heads` 重新建立到新 version table 的复合外键; +- populated schema 12 fixture 证明 baseline/head 在升级后保持可读,并可继续接受 revision 2。 + +RuntimeEvents 仍是唯一事实源。projection 删除后可以重建;canonical fact 或 tool evidence 损坏时 rebuild 不得先 +清空旧 projection。 + +## 6. Crash / concurrency matrix + +| 故障点 | 可见状态 | 恢复结论 | +|---|---|---| +| T2 写入前 | operation 仍为 `prepared`,head 不变 | 可由未来 mutation owner 重试 | +| successor event insert 后异常 | T2、fact、projection、head 全回滚 | reopen 仍看到旧 head | +| projection insert 后异常 | transaction 回滚 | 不存在“fact 有、head 没有”的半状态 | +| head CAS 前已有另一 successor | stale writer 被拒绝 | 不结算 stale operation | +| COMMIT 成功、响应丢失 | exact retry 返回 `created=false` | 不重复写 T2/fact/version | +| 更晚 successor 已推进 head 后重试旧 operation | 返回旧 operation 原先接受的 successor | current head 由独立读取返回,不篡改历史重试结果 | +| fact/tool evidence 被篡改 | canonical reader fail closed | projection 不能掩盖 corruption | + +SQLite transaction 提供三平台一致的数据库原子性;本切片不声称 Git ref、目录 rename 或 workspace 文件内容已经 +与该 transaction 原子绑定,那是 M2.2/M2.4 的责任。 + +## 7. 后续 PR 切片 + +### M2.2 — Git mutation candidate owner + +主要不变量:只有 Git artifact owner 能把一个 operation-bound candidate 证明为 base head 的合法 successor。 + +- owner:managed Git workspace service; +- 原子边界:candidate ref/commit publication 与 durable candidate receipt; +- 失败状态:base drift、额外路径、ignored mutation、artifact missing、unknown metadata 全部 park/fail closed; +- 回滚:未被 SQLite 接受的 candidate 是 orphan,可按 receipt/ref 证明后回收; +- 不做:不写 T2,不推进 workspace head,不接 Desktop/CLI。 + +该 PR 在 M2.4 消费者存在前保持 Draft。 + +### M2.3 — Mutation execution admission + +主要不变量:T1 前冻结 base workspace head、execution profile、operation identity 和 candidate lease;T1 后禁止切换 +attached/managed mode 或退回旧直接写路径。 + +- owner:Runtime Host managed mutation admission; +- 原子边界:T1 durable dispatch 选择 `managed_mutation_v1`; +- 失败状态:能力缺失、scope 失效、base/head/profile 不一致均在工具副作用前拒绝; +- 回滚:T1 前失败走标准 tool error;T1 后不确定状态保持 unsettled,交给 M2.4 收敛。 + +该 PR 不新增第二套 workspace writer,只能消费 M2.1/M2.2 的 opaque capabilities。 + +### M2.4 — Write/Edit production composition + +主要不变量:真实 Write/Edit 的成功只能来自“owned candidate 已验证 + M2.1 bundle 已提交”。 + +- 把 owner-bound worker、candidate capture 与 successor bundle 串成唯一生产路径; +- 增加真实 Host kill/reopen crash matrix; +- 保持现有 tool result、permission、sandbox 与 cancellation 语义; +- 作为 M2.2/M2.3 的首个生产消费者,完成后才允许前三个切片转 Ready/依序合并。 + +M2.4 不接 workspace-bound continuation;那属于 M3。 + +## 8. 当前验证 + +- strict successor decode/scanner 与 causal head advancement; +- exact retry 不重复写 fact/version/head; +- exact retry 在后续 head 推进后仍返回 immutable 的原接受结果; +- 失败 Write/Edit outcome 在 writer 与 canonical rebuild 两处均被拒绝; +- stale successor 不结算对应 prepared operation; +- 真实 child process 在 successor transaction 内被杀后,reopen 证明 T2/fact/projection/head 全回滚; +- 真实 child process 在 COMMIT 后被杀,reopen exact retry 收敛到同一 successor; +- populated schema 12 → 13 数据升级; +- canonical origin 被篡改后 reader fail closed; +- schema、SQLite multi-process 与既有 recovery 定向 suites 保持通过。 + +这一组验证只证明 M2.1 的 persistence authority,不代表 M2 整体完成,也不提升当前用户可见 resume 能力。 diff --git a/docs/architecture/runtime-resume-phase3-phase4-workspace-checkpoint-design.zh-CN.md b/docs/architecture/runtime-resume-phase3-phase4-workspace-checkpoint-design.zh-CN.md index 1201c7d9ad..65a62960ee 100644 --- a/docs/architecture/runtime-resume-phase3-phase4-workspace-checkpoint-design.zh-CN.md +++ b/docs/architecture/runtime-resume-phase3-phase4-workspace-checkpoint-design.zh-CN.md @@ -592,6 +592,20 @@ M1.1 合同见 这一步取代旧的“先做通用 checkpoint contract,再接 observe-only Git carrier”。不得同时维护两套 managed workspace version writer。 +M2 按 owner 与原子边界拆成四个 stacked slices: + +1. **M2.1 mutation version persistence authority(当前切片)**:定义 successor fact/scanner,并在一个 + SQLite transaction 中原子提交 T2、`version_accepted`、version projection 和 head CAS;schema 13 + 从 populated schema 12 保留 baseline/head;不接真实 Write/Edit; +2. **M2.2 Git mutation candidate owner**:唯一拥有 operation-bound candidate ref/commit、delta/path + policy、artifact receipt 与 orphan GC;不写 T2; +3. **M2.3 mutation execution admission**:T1 前冻结 base head、execution profile、operation identity 与 + candidate lease;T1 后禁止 silent fallback; +4. **M2.4 Write/Edit production composition**:把 owner-bound worker、candidate capture 与 M2.1 bundle + 串成唯一生产路径,并用真实 Host kill/reopen crash test 完成前三片的生产消费者。 + +M2.2/M2.3 在 M2.4 消费者存在前保持 Draft。M2.4 不绑定 continuation boundary;该能力仍属于 M3。 + ### M3 — Workspace-bound continuation / resume Continuation boundary 同时绑定 immutable RuntimeEvent cursor 与 accepted workspace version/epoch: @@ -627,9 +641,12 @@ M0.1 Git artifact owner (merged) └─> M1.1 execution scope admission (current) └─> M1.2 runtime-host composition └─> M1.3 explicit environment provisioning - └─> M2 mutation version acceptance - └─> M3 workspace-bound continuation - └─> M4 restore / rebaseline / publish / replication + └─> M2.1 mutation version persistence authority + └─> M2.2 Git mutation candidate owner + └─> M2.3 mutation execution admission + └─> M2.4 Write/Edit production composition + └─> M3 workspace-bound continuation + └─> M4 restore / rebaseline / publish / replication Independent maintenance gates before broad production enablement: - legacy non-empty DB root-binding adoption diff --git a/packages/core/src/__tests__/workspace-version-authority.test.ts b/packages/core/src/__tests__/workspace-version-authority.test.ts index 8ec50fb25e..e8934f9865 100644 --- a/packages/core/src/__tests__/workspace-version-authority.test.ts +++ b/packages/core/src/__tests__/workspace-version-authority.test.ts @@ -3,10 +3,12 @@ import { describe, it } from 'node:test'; import { decodeRuntimeEvent, type RuntimeEvent } from '../runtime-event.js'; import { buildWorkspaceBaselineAuthorityEvents, + buildWorkspaceSuccessorAuthorityEvent, scanWorkspaceBaselineAuthority, validateWorkspaceFactEventLane, workspaceAuthorityIdentity, type WorkspaceBaselineAuthorityInput, + type WorkspaceSuccessorAuthorityInput, } from '../workspace-version-authority.js'; describe('workspace version authority contract', () => { @@ -172,8 +174,65 @@ describe('workspace version authority contract', () => { assert.equal(orphan.hasCorruption, true); assert.equal(orphan.issues[0]?.code, 'orphan_baseline_version'); }); + + it('advances one canonical head through a causal successor fact', () => { + const baseline = buildWorkspaceBaselineAuthorityEvents(baselineInput()); + const successor = buildWorkspaceSuccessorAuthorityEvent(successorInput()); + assert.deepEqual(decodeRuntimeEvent(successor), successor); + + const scan = scanWorkspaceBaselineAuthority([ + { event: baseline.epochOpenedEvent, eventSeq: 1 }, + { event: baseline.baselineAcceptedEvent, eventSeq: 2 }, + { event: successor, eventSeq: 3 }, + ]); + + assert.equal(scan.hasCorruption, false); + assert.equal(scan.successors.length, 1); + assert.deepEqual(scan.heads, [ + { + repositoryId: baselineInput().epoch.repositoryId, + workspaceId: baselineInput().epoch.workspaceId, + workspaceEpochId: baselineInput().epoch.workspaceEpochId, + workspaceVersionId: successorInput().successor.workspaceVersionId, + acceptedEventId: successorInput().acceptedEventId, + commitOid: successorInput().successor.commitOid, + treeOid: successorInput().successor.treeOid, + revision: 2, + }, + ]); + }); }); +function successorInput(): WorkspaceSuccessorAuthorityInput { + const baseline = baselineInput(); + return { + acceptedEventId: 'workspace-successor-event-1', + committedAt: baseline.committedAt + 1, + successor: { + repositoryId: baseline.epoch.repositoryId, + workspaceId: baseline.epoch.workspaceId, + workspaceEpochId: baseline.epoch.workspaceEpochId, + workspaceVersionId: 'version_77777777777777777777777777777777', + objectFormat: baseline.epoch.objectFormat, + parentWorkspaceVersionId: baseline.baseline.workspaceVersionId, + baseAcceptedEventId: baseline.baselineAcceptedEventId, + baseHeadRevision: 1, + commitOid: '7'.repeat(40), + treeOid: '8'.repeat(40), + policyHash: baseline.epoch.policyHash, + treeDeltaDigest: `sha256:${'9'.repeat(64)}`, + changedFileCount: 1, + deletedFileCount: 0, + executionProfileDigest: `sha256:${'a'.repeat(64)}`, + }, + origin: { + operationId: 'operation-successor-1', + dispatchEventId: 'dispatch-successor-1', + outcomeEventId: 'outcome-successor-1', + }, + }; +} + function baselineInput( overrides: Partial = {}, ): WorkspaceBaselineAuthorityInput { diff --git a/packages/core/src/workspace-version-authority.ts b/packages/core/src/workspace-version-authority.ts index 16f2a99c54..05f4315a93 100644 --- a/packages/core/src/workspace-version-authority.ts +++ b/packages/core/src/workspace-version-authority.ts @@ -2,6 +2,7 @@ import type { RuntimeEvent } from './runtime-event.js'; export const WORKSPACE_EPOCH_OPENED_FACT_KIND = 'maka.workspace.epoch_opened' as const; export const WORKSPACE_BASELINE_ACCEPTED_FACT_KIND = 'maka.workspace.baseline_accepted' as const; +export const WORKSPACE_VERSION_ACCEPTED_FACT_KIND = 'maka.workspace.version_accepted' as const; export const WORKSPACE_FACT_VERSION = 1 as const; export const WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1 = 'runtime_workspace_version_authority_v1' as const; @@ -53,6 +54,52 @@ export interface WorkspaceBaselineAcceptedV1 extends WorkspaceBaselineDescriptor policyHash: `sha256:${string}`; } +export interface WorkspaceSuccessorDescriptorV1 { + repositoryId: string; + workspaceId: string; + workspaceEpochId: string; + workspaceVersionId: string; + objectFormat: WorkspaceGitObjectFormat; + parentWorkspaceVersionId: string; + baseAcceptedEventId: string; + baseHeadRevision: number; + commitOid: string; + treeOid: string; + policyHash: `sha256:${string}`; + treeDeltaDigest: `sha256:${string}`; + changedFileCount: number; + deletedFileCount: number; + executionProfileDigest: `sha256:${string}`; +} + +export interface WorkspaceMutationOriginV1 { + operationId: string; + dispatchEventId: string; + outcomeEventId: string; +} + +export interface WorkspaceVersionAcceptedV1 { + protocol: 'workspace_version_accepted_v1'; + repositoryId: string; + workspaceId: string; + workspaceEpochId: string; + workspaceVersionId: string; + objectFormat: WorkspaceGitObjectFormat; + parents: readonly [string]; + origin: { kind: 'tool_mutation' } & WorkspaceMutationOriginV1; + baseAcceptedEventId: string; + baseHeadRevision: number; + commitOid: string; + treeOid: string; + policyHash: `sha256:${string}`; + treeDeltaDigest: `sha256:${string}`; + changedFileCount: number; + deletedFileCount: number; + executionProfileDigest: `sha256:${string}`; +} + +export type WorkspaceAcceptedVersionV1 = WorkspaceBaselineAcceptedV1 | WorkspaceVersionAcceptedV1; + export type RuntimeEventWorkspaceFactEnvelope = | { kind: typeof WORKSPACE_EPOCH_OPENED_FACT_KIND; @@ -63,6 +110,11 @@ export type RuntimeEventWorkspaceFactEnvelope = kind: typeof WORKSPACE_BASELINE_ACCEPTED_FACT_KIND; version: typeof WORKSPACE_FACT_VERSION; payload: WorkspaceBaselineAcceptedV1; + } + | { + kind: typeof WORKSPACE_VERSION_ACCEPTED_FACT_KIND; + version: typeof WORKSPACE_FACT_VERSION; + payload: WorkspaceVersionAcceptedV1; }; export interface WorkspaceBaselineAuthorityInput { @@ -73,6 +125,13 @@ export interface WorkspaceBaselineAuthorityInput { baseline: WorkspaceBaselineDescriptorV1; } +export interface WorkspaceSuccessorAuthorityInput { + acceptedEventId: string; + committedAt: number; + successor: WorkspaceSuccessorDescriptorV1; + origin: WorkspaceMutationOriginV1; +} + export interface WorkspaceAuthorityIdentity { sessionId: typeof WORKSPACE_AUTHORITY_SESSION_ID; invocationId: string; @@ -100,16 +159,24 @@ export interface ScannedWorkspaceBaselineAuthority { authority: WorkspaceAuthorityIdentity; } +export interface ScannedWorkspaceSuccessorAuthority { + successor: WorkspaceVersionAcceptedV1; + acceptedEventId: string; + acceptedAt: number; + eventSeq: number; + authority: WorkspaceAuthorityIdentity; +} + export interface WorkspaceEpochRecordV1 extends WorkspaceEpochOpenedV1 { epochOpenedEventId: string; authority: WorkspaceAuthorityIdentity; committedAt: number; } -export interface WorkspaceVersionRecordV1 extends WorkspaceBaselineAcceptedV1 { - baselineAcceptedEventId: string; +export type WorkspaceVersionRecordV1 = WorkspaceAcceptedVersionV1 & { + acceptedEventId: string; committedAt: number; -} +}; export interface WorkspaceHeadRecordV1 { readonly repositoryId: string; @@ -143,7 +210,9 @@ export type WorkspaceAuthorityIssueCode = | 'missing_baseline_version' | 'orphan_baseline_version' | 'event_order_conflict' - | 'baseline_contract_conflict'; + | 'baseline_contract_conflict' + | 'successor_contract_conflict' + | 'workspace_head_conflict'; export interface WorkspaceAuthorityIssue { code: WorkspaceAuthorityIssueCode; @@ -153,6 +222,8 @@ export interface WorkspaceAuthorityIssue { export interface WorkspaceBaselineAuthorityScanResult { baselines: ScannedWorkspaceBaselineAuthority[]; + successors: ScannedWorkspaceSuccessorAuthority[]; + heads: WorkspaceHeadRecordV1[]; issues: WorkspaceAuthorityIssue[]; hasCorruption: boolean; } @@ -249,6 +320,47 @@ export function buildWorkspaceBaselineAuthorityEvents( return { epochOpenedEvent, baselineAcceptedEvent }; } +export function buildWorkspaceSuccessorAuthorityEvent( + input: WorkspaceSuccessorAuthorityInput, +): RuntimeEvent { + assertWorkspaceSuccessorAuthorityInput(input); + const identity = workspaceAuthorityIdentity(input.successor.workspaceEpochId); + const payload: WorkspaceVersionAcceptedV1 = { + protocol: 'workspace_version_accepted_v1', + repositoryId: input.successor.repositoryId, + workspaceId: input.successor.workspaceId, + workspaceEpochId: input.successor.workspaceEpochId, + workspaceVersionId: input.successor.workspaceVersionId, + objectFormat: input.successor.objectFormat, + parents: [input.successor.parentWorkspaceVersionId], + origin: { kind: 'tool_mutation', ...input.origin }, + baseAcceptedEventId: input.successor.baseAcceptedEventId, + baseHeadRevision: input.successor.baseHeadRevision, + commitOid: input.successor.commitOid, + treeOid: input.successor.treeOid, + policyHash: input.successor.policyHash, + treeDeltaDigest: input.successor.treeDeltaDigest, + changedFileCount: input.successor.changedFileCount, + deletedFileCount: input.successor.deletedFileCount, + executionProfileDigest: input.successor.executionProfileDigest, + }; + return { + id: input.acceptedEventId, + ...identity, + ts: input.committedAt, + partial: false, + role: 'system', + author: 'system', + actions: { + workspaceFact: { + kind: WORKSPACE_VERSION_ACCEPTED_FACT_KIND, + version: WORKSPACE_FACT_VERSION, + payload, + }, + }, + }; +} + export function isRuntimeEventWorkspaceFactEnvelope( value: unknown, ): value is RuntimeEventWorkspaceFactEnvelope { @@ -259,6 +371,9 @@ export function isRuntimeEventWorkspaceFactEnvelope( if (value.kind === WORKSPACE_BASELINE_ACCEPTED_FACT_KIND) { return isWorkspaceBaselineAcceptedV1(value.payload); } + if (value.kind === WORKSPACE_VERSION_ACCEPTED_FACT_KIND) { + return isWorkspaceVersionAcceptedV1(value.payload); + } return false; } @@ -300,7 +415,8 @@ export function scanWorkspaceBaselineAuthority( const issues: WorkspaceAuthorityIssue[] = []; const seenEventIds = new Set(); const epochRows = new Map(); - const versionRows = new Map(); + const baselineRows = new Map(); + const successorRows = new Map(); const versionIds = new Map(); for (const row of rows) { @@ -330,11 +446,13 @@ export function scanWorkspaceBaselineAuthority( epochRows.set(epochId, matches); continue; } - const matches = versionRows.get(epochId) ?? []; + const target = + fact.kind === WORKSPACE_BASELINE_ACCEPTED_FACT_KIND ? baselineRows : successorRows; + const matches = target.get(epochId) ?? []; matches.push(row); - versionRows.set(epochId, matches); + target.set(epochId, matches); const priorEpoch = versionIds.get(fact.payload.workspaceVersionId); - if (priorEpoch !== undefined && priorEpoch !== epochId) { + if (priorEpoch !== undefined) { issues.push({ code: 'duplicate_workspace_version', eventId: event.id, @@ -346,14 +464,20 @@ export function scanWorkspaceBaselineAuthority( } const baselines: ScannedWorkspaceBaselineAuthority[] = []; - const epochIds = new Set([...epochRows.keys(), ...versionRows.keys()]); + const successors: ScannedWorkspaceSuccessorAuthority[] = []; + const heads: WorkspaceHeadRecordV1[] = []; + const epochIds = new Set([...epochRows.keys(), ...baselineRows.keys(), ...successorRows.keys()]); for (const epochId of [...epochIds].sort()) { const opened = epochRows.get(epochId) ?? []; - const accepted = versionRows.get(epochId) ?? []; + const accepted = baselineRows.get(epochId) ?? []; + const pendingSuccessors = [...(successorRows.get(epochId) ?? [])].sort( + (left, right) => + left.eventSeq - right.eventSeq || left.event.id.localeCompare(right.event.id), + ); if (opened.length === 0) { issues.push({ code: 'orphan_baseline_version', - eventId: accepted[0]!.event.id, + eventId: (accepted[0] ?? pendingSuccessors[0])!.event.id, workspaceEpochId: epochId, }); continue; @@ -382,15 +506,15 @@ export function scanWorkspaceBaselineAuthority( }); continue; } + let baseline: ScannedWorkspaceBaselineAuthority; try { - baselines.push( - assertWorkspaceBaselineAuthorityPair({ - epochOpenedEvent: opened[0]!.event, - baselineAcceptedEvent: accepted[0]!.event, - epochEventSeq: opened[0]!.eventSeq, - baselineEventSeq: accepted[0]!.eventSeq, - }), - ); + baseline = assertWorkspaceBaselineAuthorityPair({ + epochOpenedEvent: opened[0]!.event, + baselineAcceptedEvent: accepted[0]!.event, + epochEventSeq: opened[0]!.eventSeq, + baselineEventSeq: accepted[0]!.eventSeq, + }); + baselines.push(baseline); } catch (error) { const code = error instanceof WorkspaceAuthorityContractError @@ -401,11 +525,57 @@ export function scanWorkspaceBaselineAuthority( eventId: accepted[0]!.event.id, workspaceEpochId: epochId, }); + continue; + } + + let head: WorkspaceHeadRecordV1 = { + repositoryId: baseline.epoch.repositoryId, + workspaceId: baseline.epoch.workspaceId, + workspaceEpochId: baseline.epoch.workspaceEpochId, + workspaceVersionId: baseline.baseline.workspaceVersionId, + acceptedEventId: baseline.baselineAcceptedEventId, + commitOid: baseline.baseline.commitOid, + treeOid: baseline.baseline.treeOid, + revision: 1, + }; + for (const row of pendingSuccessors) { + try { + const scanned = assertWorkspaceSuccessorAuthority({ + baseline, + currentHead: head, + event: row.event, + eventSeq: row.eventSeq, + }); + successors.push(scanned); + head = { + repositoryId: scanned.successor.repositoryId, + workspaceId: scanned.successor.workspaceId, + workspaceEpochId: scanned.successor.workspaceEpochId, + workspaceVersionId: scanned.successor.workspaceVersionId, + acceptedEventId: scanned.acceptedEventId, + commitOid: scanned.successor.commitOid, + treeOid: scanned.successor.treeOid, + revision: head.revision + 1, + }; + } catch (error) { + issues.push({ + code: + error instanceof WorkspaceAuthorityContractError + ? error.code + : 'successor_contract_conflict', + eventId: row.event.id, + workspaceEpochId: epochId, + }); + break; + } } + heads.push(head); } return { baselines: issues.length === 0 ? baselines : [], + successors: issues.length === 0 ? successors : [], + heads: issues.length === 0 ? heads : [], issues, hasCorruption: issues.length > 0, }; @@ -415,7 +585,10 @@ class WorkspaceAuthorityContractError extends Error { constructor( readonly code: Extract< WorkspaceAuthorityIssueCode, - 'event_order_conflict' | 'baseline_contract_conflict' + | 'event_order_conflict' + | 'baseline_contract_conflict' + | 'successor_contract_conflict' + | 'workspace_head_conflict' >, message: string, ) { @@ -484,6 +657,52 @@ function assertWorkspaceBaselineAuthorityPair(input: { }; } +function assertWorkspaceSuccessorAuthority(input: { + baseline: ScannedWorkspaceBaselineAuthority; + currentHead: WorkspaceHeadRecordV1; + event: RuntimeEvent; + eventSeq: number; +}): ScannedWorkspaceSuccessorAuthority { + const lane = validateWorkspaceFactEventLane(input.event); + const fact = input.event.actions?.workspaceFact; + if (!lane.ok || fact?.kind !== WORKSPACE_VERSION_ACCEPTED_FACT_KIND) { + throw new WorkspaceAuthorityContractError( + 'successor_contract_conflict', + 'Invalid workspace successor authority event lane', + ); + } + const successor = fact.payload; + const expectedSeq = input.currentHead.revision + 2; + if (input.eventSeq !== expectedSeq) { + throw new WorkspaceAuthorityContractError( + 'event_order_conflict', + 'Workspace successor facts must form one contiguous authority sequence', + ); + } + if ( + successor.repositoryId !== input.baseline.epoch.repositoryId || + successor.workspaceId !== input.baseline.epoch.workspaceId || + successor.workspaceEpochId !== input.baseline.epoch.workspaceEpochId || + successor.objectFormat !== input.baseline.epoch.objectFormat || + successor.policyHash !== input.baseline.epoch.policyHash || + successor.parents[0] !== input.currentHead.workspaceVersionId || + successor.baseAcceptedEventId !== input.currentHead.acceptedEventId || + successor.baseHeadRevision !== input.currentHead.revision + ) { + throw new WorkspaceAuthorityContractError( + 'workspace_head_conflict', + 'Workspace successor does not advance the current canonical head', + ); + } + return { + successor, + acceptedEventId: input.event.id, + acceptedAt: input.event.ts, + eventSeq: input.eventSeq, + authority: workspaceAuthorityIdentity(successor.workspaceEpochId), + }; +} + function assertWorkspaceBaselineAuthorityInput(input: WorkspaceBaselineAuthorityInput): void { if ( !EVENT_ID_PATTERN.test(input.epochOpenedEventId) || @@ -499,6 +718,18 @@ function assertWorkspaceBaselineAuthorityInput(input: WorkspaceBaselineAuthority } } +function assertWorkspaceSuccessorAuthorityInput(input: WorkspaceSuccessorAuthorityInput): void { + if ( + !EVENT_ID_PATTERN.test(input.acceptedEventId) || + !Number.isSafeInteger(input.committedAt) || + input.committedAt < 0 || + !isWorkspaceSuccessorDescriptor(input.successor) || + !isWorkspaceMutationOrigin(input.origin) + ) { + throw new Error('Invalid workspace successor authority input'); + } +} + function isWorkspaceEpochOpenedV1(value: unknown): value is WorkspaceEpochOpenedV1 { return ( hasExactKeys(value, [ @@ -575,6 +806,76 @@ function isWorkspaceBaselineAcceptedV1(value: unknown): value is WorkspaceBaseli return isWorkspaceBaselineDescriptor(value, value.objectFormat); } +function isWorkspaceVersionAcceptedV1(value: unknown): value is WorkspaceVersionAcceptedV1 { + if ( + !hasExactKeys(value, [ + 'protocol', + 'repositoryId', + 'workspaceId', + 'workspaceEpochId', + 'workspaceVersionId', + 'objectFormat', + 'parents', + 'origin', + 'baseAcceptedEventId', + 'baseHeadRevision', + 'commitOid', + 'treeOid', + 'policyHash', + 'treeDeltaDigest', + 'changedFileCount', + 'deletedFileCount', + 'executionProfileDigest', + ]) || + value.protocol !== 'workspace_version_accepted_v1' || + !isWorkspaceSuccessorDescriptor({ + ...value, + parentWorkspaceVersionId: Array.isArray(value.parents) ? value.parents[0] : undefined, + }) || + !Array.isArray(value.parents) || + value.parents.length !== 1 || + !hasExactKeys(value.origin, ['kind', 'operationId', 'dispatchEventId', 'outcomeEventId']) || + value.origin.kind !== 'tool_mutation' || + !isWorkspaceMutationOrigin(value.origin) + ) { + return false; + } + return true; +} + +function isWorkspaceSuccessorDescriptor(value: unknown): value is WorkspaceSuccessorDescriptorV1 { + if (!isRecord(value)) return false; + return ( + isIdentifier(value.repositoryId, 'repositoryId') && + isIdentifier(value.workspaceId, 'workspaceId') && + isIdentifier(value.workspaceEpochId, 'workspaceEpochId') && + isIdentifier(value.workspaceVersionId, 'workspaceVersionId') && + isIdentifier(value.parentWorkspaceVersionId, 'workspaceVersionId') && + EVENT_ID_PATTERN.test(String(value.baseAcceptedEventId)) && + value.workspaceVersionId !== value.parentWorkspaceVersionId && + isObjectFormat(value.objectFormat) && + Number.isSafeInteger(value.baseHeadRevision) && + Number(value.baseHeadRevision) > 0 && + isGitOid(value.commitOid, value.objectFormat) && + isGitOid(value.treeOid, value.objectFormat) && + isSha256Digest(value.policyHash) && + isSha256Digest(value.treeDeltaDigest) && + isNonNegativeSafeInteger(value.changedFileCount) && + isNonNegativeSafeInteger(value.deletedFileCount) && + isSha256Digest(value.executionProfileDigest) + ); +} + +function isWorkspaceMutationOrigin(value: unknown): value is WorkspaceMutationOriginV1 { + if (!isRecord(value)) return false; + return ( + EVENT_ID_PATTERN.test(String(value.operationId)) && + EVENT_ID_PATTERN.test(String(value.dispatchEventId)) && + EVENT_ID_PATTERN.test(String(value.outcomeEventId)) && + value.dispatchEventId !== value.outcomeEventId + ); +} + function isWorkspaceBaselineDescriptor( value: unknown, objectFormat: WorkspaceGitObjectFormat, diff --git a/packages/storage/src/__tests__/sqlite-runtime-crash.test.ts b/packages/storage/src/__tests__/sqlite-runtime-crash.test.ts index 8d8a7ce95c..b5c140a17f 100644 --- a/packages/storage/src/__tests__/sqlite-runtime-crash.test.ts +++ b/packages/storage/src/__tests__/sqlite-runtime-crash.test.ts @@ -16,6 +16,8 @@ import { import { bindWorkspaceBaselineAuthorityStoreRootInternal, commitWorkspaceBaselineInternal, + commitWorkspaceSuccessorInternal, + type WorkspaceSuccessorCommitInput, } from '../workspace-version-authority-internal.js'; const CRASH_READ_ARGS_HASH = canonicalToolArgsHash('Read', { @@ -135,6 +137,38 @@ if (childMode) { ); }); }); + + it('rolls back a workspace successor when killed inside its authority transaction', { + timeout: 30_000, + }, async () => { + await withKilledChild('inside_workspace_successor', async (store) => { + assert.equal( + (await store.readWorkspaceHead(`workspace_${'2'.repeat(32)}`, `epoch_${'3'.repeat(32)}`)) + ?.workspaceVersionId, + `version_${'5'.repeat(32)}`, + ); + assert.equal( + (await store.readToolOperation('workspace-successor-operation'))?.currentState, + 'prepared', + ); + assert.equal(await store.readWorkspaceVersion(`version_${'7'.repeat(32)}`), undefined); + }); + }); + + it('returns the accepted workspace successor after a process is killed post-commit', { + timeout: 30_000, + }, async () => { + await withKilledChild('after_workspace_successor_commit', async (store) => { + bindWorkspaceBaselineAuthorityStoreRootInternal(store, 'a'.repeat(64)); + const retry = await commitWorkspaceSuccessorInternal(store, workspaceSuccessorCommit()); + assert.equal(retry.created, false); + assert.equal(retry.head.workspaceVersionId, `version_${'7'.repeat(32)}`); + assert.equal( + (await store.readToolOperation('workspace-successor-operation'))?.currentState, + 'outcome_committed', + ); + }); + }); }); } @@ -216,12 +250,28 @@ async function runCrashChild(mode: string): Promise { if (point === 'after_workspace_version_event_insert' && mode === 'inside_workspace_baseline') { blockUntilKilled(); } + if ( + point === 'after_workspace_successor_event_insert' && + mode === 'inside_workspace_successor' + ) { + blockUntilKilled(); + } }; const store = createSqliteRuntimeStore(dbPath, { failpoint }); - if (mode === 'inside_workspace_baseline' || mode === 'after_workspace_baseline_commit') { + if ( + mode === 'inside_workspace_baseline' || + mode === 'after_workspace_baseline_commit' || + mode === 'inside_workspace_successor' || + mode === 'after_workspace_successor_commit' + ) { bindWorkspaceBaselineAuthorityStoreRootInternal(store, 'a'.repeat(64)); await commitWorkspaceBaselineInternal(store, workspaceBaselineInput()); if (mode === 'after_workspace_baseline_commit') blockUntilKilled(); + if (mode === 'inside_workspace_successor' || mode === 'after_workspace_successor_commit') { + await store.commitToolPrepared(workspaceSuccessorPreparedCommit()); + await commitWorkspaceSuccessorInternal(store, workspaceSuccessorCommit()); + if (mode === 'after_workspace_successor_commit') blockUntilKilled(); + } throw new Error(`Workspace baseline crash mode ${mode} missed its failpoint`); } await store.commitToolPrepared(preparedCommit()); @@ -294,6 +344,123 @@ function workspaceBaselineInput(): WorkspaceBaselineAuthorityInput { }; } +function workspaceSuccessorPreparedCommit() { + const args = { path: 'notes.txt', content: 'successor' }; + const canonicalArgsHash = canonicalToolArgsHash('Write', args); + return { + operationId: 'workspace-successor-operation', + journalEventId: 'workspace-successor-operation_prepared', + runtimeEvent: { + id: 'workspace-successor-call', + invocationId: 'workspace-successor-invocation', + runId: 'workspace-successor-run', + sessionId: 'workspace-successor-session', + turnId: 'workspace-successor-turn', + ts: 1_700_000_000_001, + partial: false, + role: 'model' as const, + author: 'agent' as const, + content: { + kind: 'function_call' as const, + id: 'workspace-successor-call-id', + name: 'Write', + args, + }, + refs: { + operationId: 'workspace-successor-operation', + toolCallId: 'workspace-successor-call-id', + }, + }, + dispatchRuntimeEvent: { + id: 'workspace-successor-dispatch', + invocationId: 'workspace-successor-invocation', + runId: 'workspace-successor-run', + sessionId: 'workspace-successor-session', + turnId: 'workspace-successor-turn', + ts: 1_700_000_000_001, + partial: false, + role: 'system' as const, + author: 'system' as const, + actions: { + toolDispatch: { + protocol: 't1_after_preflight_v1' as const, + operationId: 'workspace-successor-operation', + providerToolCallId: 'workspace-successor-call-id', + toolName: 'Write', + canonicalArgsHash, + recoveryMode: 'reconcile' as const, + }, + }, + refs: { + operationId: 'workspace-successor-operation', + toolCallId: 'workspace-successor-call-id', + }, + }, + providerToolCallId: 'workspace-successor-call-id', + toolName: 'Write', + canonicalArgsHash, + recoveryMode: 'reconcile' as const, + committedAt: 1_700_000_000_001, + }; +} + +function workspaceSuccessorCommit(): WorkspaceSuccessorCommitInput { + return { + successor: { + acceptedEventId: 'workspace-successor-accepted', + committedAt: 1_700_000_000_002, + successor: { + repositoryId: `repository_${'1'.repeat(32)}`, + workspaceId: `workspace_${'2'.repeat(32)}`, + workspaceEpochId: `epoch_${'3'.repeat(32)}`, + workspaceVersionId: `version_${'7'.repeat(32)}`, + objectFormat: 'sha1', + parentWorkspaceVersionId: `version_${'5'.repeat(32)}`, + baseAcceptedEventId: 'workspace-version-event-1', + baseHeadRevision: 1, + commitOid: '7'.repeat(40), + treeOid: '8'.repeat(40), + policyHash: `sha256:${'4'.repeat(64)}`, + treeDeltaDigest: `sha256:${'9'.repeat(64)}`, + changedFileCount: 1, + deletedFileCount: 0, + executionProfileDigest: `sha256:${'a'.repeat(64)}`, + }, + origin: { + operationId: 'workspace-successor-operation', + dispatchEventId: 'workspace-successor-dispatch', + outcomeEventId: 'workspace-successor-outcome', + }, + }, + toolOutcome: { + operationId: 'workspace-successor-operation', + journalEventId: 'workspace-successor-operation_outcome', + committedAt: 1_700_000_000_002, + runtimeEvent: { + id: 'workspace-successor-outcome', + invocationId: 'workspace-successor-invocation', + runId: 'workspace-successor-run', + sessionId: 'workspace-successor-session', + turnId: 'workspace-successor-turn', + ts: 1_700_000_000_002, + partial: false, + role: 'tool', + author: 'tool', + content: { + kind: 'function_response', + id: 'workspace-successor-call-id', + name: 'Write', + result: 'Wrote notes.txt', + }, + refs: { + operationId: 'workspace-successor-operation', + toolCallId: 'workspace-successor-call-id', + }, + }, + }, + }; +} + function outcomeCommit() { return { operationId: 'operation-1', diff --git a/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts b/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts index 282da536b3..fffb51f30e 100644 --- a/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts +++ b/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts @@ -20,6 +20,8 @@ import { import { bindWorkspaceBaselineAuthorityStoreRootInternal, commitWorkspaceBaselineInternal, + commitWorkspaceSuccessorInternal, + type WorkspaceSuccessorCommitInput, } from '../workspace-version-authority-internal.js'; const TEST_STORAGE_ROOT_ID = 'a'.repeat(64); @@ -136,6 +138,231 @@ describe('workspace version persistence authority', () => { }); }); + it('atomically commits one tool outcome with its successor workspace head', async () => { + await withDatabase(async ({ dbPath, store }) => { + const { baseline, input } = await prepareSuccessorCommit(store); + const result = await commitWorkspaceSuccessorInternal(store, input); + + assert.equal(result.created, true); + assert.equal(result.head.revision, 2); + assert.equal( + (await store.readToolOperation(input.toolOutcome.operationId))?.currentState, + 'outcome_committed', + ); + assert.deepEqual( + await store.readWorkspaceHead(baseline.epoch.workspaceId, baseline.epoch.workspaceEpochId), + result.head, + ); + assert.equal( + (await store.readWorkspaceVersion(result.head.workspaceVersionId))?.origin.kind, + 'tool_mutation', + ); + + const retry = await commitWorkspaceSuccessorInternal(store, input); + assert.deepEqual(retry, { ...result, created: false }); + const raw = new DatabaseSync(dbPath); + try { + assert.equal(count(raw, 'runtime_workspace_versions'), 2); + assert.equal(count(raw, 'runtime_workspace_heads'), 1); + assert.equal( + countWhere(raw, 'runtime_events', 'session_id = ?', WORKSPACE_AUTHORITY_SESSION_ID), + 3, + ); + } finally { + raw.close(); + } + + const corrupt = new DatabaseSync(dbPath); + try { + const row = corrupt + .prepare('SELECT payload_json FROM runtime_events WHERE event_id = ?') + .get(input.successor.acceptedEventId) as { payload_json: string }; + const event = JSON.parse(row.payload_json) as RuntimeEvent; + const fact = event.actions?.workspaceFact; + assert.equal(fact?.kind, 'maka.workspace.version_accepted'); + if (fact?.kind === 'maka.workspace.version_accepted') { + fact.payload.origin.outcomeEventId = 'other-outcome-event'; + } + corrupt + .prepare('UPDATE runtime_events SET payload_json = ? WHERE event_id = ?') + .run(JSON.stringify(event), input.successor.acceptedEventId); + } finally { + corrupt.close(); + } + await assert.rejects( + store.readWorkspaceHead(baseline.epoch.workspaceId, baseline.epoch.workspaceEpochId), + /workspace successor tool evidence: identity_conflict/i, + ); + }); + }); + + it('rejects a failed Write outcome without advancing the workspace head', async () => { + await withDatabase(async ({ store }) => { + const { baseline, input } = await prepareSuccessorCommit(store); + assert.equal(input.toolOutcome.runtimeEvent.content?.kind, 'function_response'); + if (input.toolOutcome.runtimeEvent.content?.kind !== 'function_response') { + throw new Error('Expected a function response fixture'); + } + input.toolOutcome.runtimeEvent.content.isError = true; + + await assert.rejects( + commitWorkspaceSuccessorInternal(store, input), + /workspace successor requires a successful tool outcome/i, + ); + assert.equal( + (await store.readToolOperation(input.toolOutcome.operationId))?.currentState, + 'prepared', + ); + assert.equal( + (await store.readWorkspaceHead(baseline.epoch.workspaceId, baseline.epoch.workspaceEpochId)) + ?.workspaceVersionId, + baseline.baseline.workspaceVersionId, + ); + }); + }); + + it('rolls back tool outcome, successor fact, projection, and head together', async () => { + await withDatabase(async ({ dbPath, store, setFailpoint }) => { + const { baseline, input } = await prepareSuccessorCommit(store); + setFailpoint('after_workspace_successor_event_insert'); + await assert.rejects(commitWorkspaceSuccessorInternal(store, input), /failpoint/); + setFailpoint(undefined); + store.close(); + + const reopened = createSqliteRuntimeStore(dbPath); + try { + assert.equal( + (await reopened.readToolOperation(input.toolOutcome.operationId))?.currentState, + 'prepared', + ); + assert.equal( + ( + await reopened.readWorkspaceHead( + baseline.epoch.workspaceId, + baseline.epoch.workspaceEpochId, + ) + )?.workspaceVersionId, + baseline.baseline.workspaceVersionId, + ); + assert.equal( + await reopened.readWorkspaceVersion(input.successor.successor.workspaceVersionId), + undefined, + ); + } finally { + reopened.close(); + } + }); + }); + + it('rejects a stale successor without settling its prepared tool operation', async () => { + await withDatabase(async ({ store }) => { + const first = await prepareSuccessorCommit(store, 1); + const stale = await prepareSuccessorCommit(store, 2); + + await commitWorkspaceSuccessorInternal(store, first.input); + await assert.rejects( + commitWorkspaceSuccessorInternal(store, stale.input), + /compare-and-set base head conflict/i, + ); + assert.equal( + (await store.readToolOperation(stale.input.toolOutcome.operationId))?.currentState, + 'prepared', + ); + }); + }); + + it('returns an earlier exact successor retry after the canonical head advances', async () => { + await withDatabase(async ({ store }) => { + const first = await prepareSuccessorCommit(store, 1); + const firstResult = await commitWorkspaceSuccessorInternal(store, first.input); + const second = await prepareSuccessorCommit(store, 2); + const secondResult = await commitWorkspaceSuccessorInternal(store, second.input); + assert.equal(secondResult.head.revision, firstResult.head.revision + 1); + + const retry = await commitWorkspaceSuccessorInternal(store, first.input); + assert.deepEqual(retry, { ...firstResult, created: false }); + assert.deepEqual( + await store.readWorkspaceHead( + first.baseline.epoch.workspaceId, + first.baseline.epoch.workspaceEpochId, + ), + secondResult.head, + ); + }); + }); + + it('rejects a failed outcome referenced by immutable successor authority', async () => { + await withDatabase(async ({ dbPath, store }) => { + const { input } = await prepareSuccessorCommit(store); + await commitWorkspaceSuccessorInternal(store, input); + + const raw = new DatabaseSync(dbPath); + try { + const row = raw + .prepare('SELECT payload_json FROM runtime_events WHERE event_id = ?') + .get(input.toolOutcome.runtimeEvent.id) as { payload_json: string }; + const outcome = JSON.parse(row.payload_json) as RuntimeEvent; + assert.equal(outcome.content?.kind, 'function_response'); + if (outcome.content?.kind !== 'function_response') { + throw new Error('Expected a function response fixture'); + } + outcome.content.isError = true; + raw + .prepare('UPDATE runtime_events SET payload_json = ? WHERE event_id = ?') + .run(JSON.stringify(outcome), input.toolOutcome.runtimeEvent.id); + } finally { + raw.close(); + } + + await assert.rejects( + store.rebuildWorkspaceVersionProjections(), + /workspace successor tool evidence: identity_conflict/i, + ); + }); + }); + + it('upgrades a populated schema 12 baseline before accepting its successor', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-workspace-schema-12-')); + const dbPath = join(root, 'runtime.sqlite'); + const baseline = baselineInput(); + const original = createSqliteRuntimeStore(dbPath); + bindWorkspaceBaselineAuthorityStoreRootInternal(original, TEST_STORAGE_ROOT_ID); + await commitWorkspaceBaselineInternal(original, baseline); + original.close(); + + try { + const legacy = new DatabaseSync(dbPath); + try { + recreateWorkspaceTablesAsSchema12(legacy); + } finally { + legacy.close(); + } + + const upgraded = createSqliteRuntimeStore(dbPath); + bindWorkspaceBaselineAuthorityStoreRootInternal(upgraded, TEST_STORAGE_ROOT_ID); + try { + assert.equal(upgraded.schemaVersion(), 13); + assert.equal( + ( + await upgraded.readWorkspaceHead( + baseline.epoch.workspaceId, + baseline.epoch.workspaceEpochId, + ) + )?.workspaceVersionId, + baseline.baseline.workspaceVersionId, + ); + + const prepared = await prepareSuccessorCommit(upgraded); + const accepted = await commitWorkspaceSuccessorInternal(upgraded, prepared.input); + assert.equal(accepted.head.revision, 2); + } finally { + upgraded.close(); + } + } finally { + await rm(root, { recursive: true, force: true }); + } + }); + it('refuses to silently claim unbound workspace authority facts for another root', async () => { await withDatabase(async ({ dbPath, store }) => { await commitWorkspaceBaselineInternal(store, baselineInput()); @@ -483,6 +710,197 @@ async function withDatabase( } } +async function prepareSuccessorCommit( + store: ReturnType, + variant = 1, +): Promise<{ + baseline: WorkspaceBaselineAuthorityInput; + input: WorkspaceSuccessorCommitInput; +}> { + const baseline = baselineInput(); + const opened = await commitWorkspaceBaselineInternal(store, baseline); + const args = { path: 'notes.txt', content: 'successor' }; + const argsHash = canonicalToolArgsHash('Write', args); + const operationId = `operation-successor-${variant}`; + const toolCallId = `call-successor-${variant}`; + const commitDigit = variant === 1 ? '7' : '6'; + const treeDigit = variant === 1 ? '8' : '5'; + await store.commitToolPrepared({ + operationId, + journalEventId: `${operationId}_prepared`, + runtimeEvent: { + id: `call-successor-event-${variant}`, + sessionId: 'session-successor', + invocationId: 'invocation-successor', + runId: 'run-successor', + turnId: 'turn-successor', + ts: baseline.committedAt + 1, + partial: false, + role: 'model', + author: 'agent', + content: { kind: 'function_call', id: toolCallId, name: 'Write', args }, + refs: { operationId, toolCallId }, + }, + dispatchRuntimeEvent: { + id: `dispatch-successor-event-${variant}`, + sessionId: 'session-successor', + invocationId: 'invocation-successor', + runId: 'run-successor', + turnId: 'turn-successor', + ts: baseline.committedAt + 1, + partial: false, + role: 'system', + author: 'system', + actions: { + toolDispatch: { + protocol: 't1_after_preflight_v1', + operationId, + providerToolCallId: toolCallId, + toolName: 'Write', + canonicalArgsHash: argsHash, + recoveryMode: 'reconcile', + }, + }, + refs: { operationId, toolCallId }, + }, + providerToolCallId: toolCallId, + toolName: 'Write', + canonicalArgsHash: argsHash, + recoveryMode: 'reconcile', + committedAt: baseline.committedAt + 1, + }); + + return { + baseline, + input: { + successor: { + acceptedEventId: `workspace-successor-event-${variant}`, + committedAt: baseline.committedAt + 2, + successor: { + repositoryId: baseline.epoch.repositoryId, + workspaceId: baseline.epoch.workspaceId, + workspaceEpochId: baseline.epoch.workspaceEpochId, + workspaceVersionId: `version_${commitDigit.repeat(32)}`, + objectFormat: baseline.epoch.objectFormat, + parentWorkspaceVersionId: opened.head.workspaceVersionId, + baseAcceptedEventId: opened.head.acceptedEventId, + baseHeadRevision: opened.head.revision, + commitOid: commitDigit.repeat(40), + treeOid: treeDigit.repeat(40), + policyHash: baseline.epoch.policyHash, + treeDeltaDigest: `sha256:${'9'.repeat(64)}`, + changedFileCount: 1, + deletedFileCount: 0, + executionProfileDigest: `sha256:${'a'.repeat(64)}`, + }, + origin: { + operationId, + dispatchEventId: `dispatch-successor-event-${variant}`, + outcomeEventId: `outcome-successor-event-${variant}`, + }, + }, + toolOutcome: { + operationId, + journalEventId: `${operationId}_outcome`, + committedAt: baseline.committedAt + 2, + runtimeEvent: { + id: `outcome-successor-event-${variant}`, + sessionId: 'session-successor', + invocationId: 'invocation-successor', + runId: 'run-successor', + turnId: 'turn-successor', + ts: baseline.committedAt + 2, + partial: false, + role: 'tool', + author: 'tool', + content: { + kind: 'function_response', + id: toolCallId, + name: 'Write', + result: 'Wrote notes.txt', + }, + refs: { operationId, toolCallId }, + }, + }, + }, + }; +} + +function recreateWorkspaceTablesAsSchema12(database: DatabaseSync): void { + database.exec(` + PRAGMA foreign_keys = OFF; + BEGIN IMMEDIATE; + + ALTER TABLE runtime_workspace_heads RENAME TO runtime_workspace_heads_schema_13; + ALTER TABLE runtime_workspace_versions RENAME TO runtime_workspace_versions_schema_13; + + CREATE TABLE runtime_workspace_versions ( + workspace_version_id TEXT PRIMARY KEY, + repository_id TEXT NOT NULL, + workspace_id TEXT NOT NULL, + workspace_epoch_id TEXT NOT NULL, + object_format TEXT NOT NULL CHECK (object_format IN ('sha1', 'sha256')), + origin_kind TEXT NOT NULL CHECK (origin_kind = 'baseline'), + origin_event_id TEXT NOT NULL, + parents_json TEXT NOT NULL CHECK (parents_json = '[]'), + commit_oid TEXT NOT NULL, + tree_oid TEXT NOT NULL, + policy_hash TEXT NOT NULL, + tree_delta_digest TEXT NOT NULL, + changed_file_count INTEGER NOT NULL CHECK (changed_file_count >= 0), + deleted_file_count INTEGER NOT NULL CHECK (deleted_file_count = 0), + accepted_event_id TEXT NOT NULL UNIQUE REFERENCES runtime_events(event_id), + protocol_version INTEGER NOT NULL CHECK (protocol_version = 1), + committed_at INTEGER NOT NULL, + FOREIGN KEY (workspace_id, workspace_epoch_id) + REFERENCES runtime_workspace_epochs(workspace_id, workspace_epoch_id), + UNIQUE (workspace_id, workspace_epoch_id, workspace_version_id, accepted_event_id) + ); + + INSERT INTO runtime_workspace_versions ( + workspace_version_id, repository_id, workspace_id, workspace_epoch_id, + object_format, origin_kind, origin_event_id, parents_json, + commit_oid, tree_oid, policy_hash, tree_delta_digest, + changed_file_count, deleted_file_count, accepted_event_id, + protocol_version, committed_at + ) + SELECT + workspace_version_id, repository_id, workspace_id, workspace_epoch_id, + object_format, origin_kind, origin_event_id, parents_json, + commit_oid, tree_oid, policy_hash, tree_delta_digest, + changed_file_count, deleted_file_count, accepted_event_id, + protocol_version, committed_at + FROM runtime_workspace_versions_schema_13; + + CREATE TABLE runtime_workspace_heads ( + workspace_id TEXT NOT NULL, + workspace_epoch_id TEXT NOT NULL, + repository_id TEXT NOT NULL, + workspace_version_id TEXT NOT NULL, + accepted_event_id TEXT NOT NULL, + commit_oid TEXT NOT NULL, + tree_oid TEXT NOT NULL, + revision INTEGER NOT NULL CHECK (revision > 0), + PRIMARY KEY (workspace_id, workspace_epoch_id), + FOREIGN KEY (workspace_id, workspace_epoch_id) + REFERENCES runtime_workspace_epochs(workspace_id, workspace_epoch_id), + FOREIGN KEY (workspace_id, workspace_epoch_id, workspace_version_id, accepted_event_id) + REFERENCES runtime_workspace_versions( + workspace_id, workspace_epoch_id, workspace_version_id, accepted_event_id + ) + ); + + INSERT INTO runtime_workspace_heads + SELECT * FROM runtime_workspace_heads_schema_13; + + DROP TABLE runtime_workspace_heads_schema_13; + DROP TABLE runtime_workspace_versions_schema_13; + PRAGMA user_version = 12; + COMMIT; + PRAGMA foreign_keys = ON; + `); +} + function baselineInput( overrides: Partial = {}, ): WorkspaceBaselineAuthorityInput { diff --git a/packages/storage/src/sqlite-runtime-schema.ts b/packages/storage/src/sqlite-runtime-schema.ts index bfbc391806..542cba1fc7 100644 --- a/packages/storage/src/sqlite-runtime-schema.ts +++ b/packages/storage/src/sqlite-runtime-schema.ts @@ -1,6 +1,6 @@ import type { DatabaseSync } from 'node:sqlite'; -export const SQLITE_RUNTIME_SCHEMA_VERSION = 12; +export const SQLITE_RUNTIME_SCHEMA_VERSION = 13; export const RUNTIME_RECOVERY_AUTHORITY_CAPABILITY = 'runtime_recovery_authority'; export const RUNTIME_RECOVERY_AUTHORITY_CAPABILITY_VERSION = 1; export const RUNTIME_CONTINUATION_AUTHORITY_CAPABILITY = 'runtime_continuation_authority'; @@ -301,6 +301,102 @@ const MIGRATIONS: ReadonlyMap = new Map([ DROP TABLE IF EXISTS headless_task_run_events; `, ], + [ + 13, + ` + ALTER TABLE runtime_workspace_heads RENAME TO runtime_workspace_heads_v12; + ALTER TABLE runtime_workspace_versions RENAME TO runtime_workspace_versions_v12; + + CREATE TABLE runtime_workspace_versions ( + workspace_version_id TEXT PRIMARY KEY, + repository_id TEXT NOT NULL, + workspace_id TEXT NOT NULL, + workspace_epoch_id TEXT NOT NULL, + object_format TEXT NOT NULL CHECK (object_format IN ('sha1', 'sha256')), + origin_kind TEXT NOT NULL CHECK (origin_kind IN ('baseline', 'tool_mutation')), + origin_event_id TEXT NOT NULL, + parents_json TEXT NOT NULL, + operation_id TEXT, + dispatch_event_id TEXT REFERENCES runtime_events(event_id), + outcome_event_id TEXT REFERENCES runtime_events(event_id), + base_head_revision INTEGER CHECK (base_head_revision IS NULL OR base_head_revision > 0), + execution_profile_digest TEXT, + commit_oid TEXT NOT NULL, + tree_oid TEXT NOT NULL, + policy_hash TEXT NOT NULL, + tree_delta_digest TEXT NOT NULL, + changed_file_count INTEGER NOT NULL CHECK (changed_file_count >= 0), + deleted_file_count INTEGER NOT NULL CHECK (deleted_file_count >= 0), + accepted_event_id TEXT NOT NULL UNIQUE REFERENCES runtime_events(event_id), + protocol_version INTEGER NOT NULL CHECK (protocol_version = 1), + committed_at INTEGER NOT NULL, + FOREIGN KEY (workspace_id, workspace_epoch_id) + REFERENCES runtime_workspace_epochs(workspace_id, workspace_epoch_id), + UNIQUE ( + workspace_id, + workspace_epoch_id, + workspace_version_id, + accepted_event_id + ), + CHECK ( + (origin_kind = 'baseline' AND parents_json = '[]' AND operation_id IS NULL + AND dispatch_event_id IS NULL AND outcome_event_id IS NULL + AND base_head_revision IS NULL AND execution_profile_digest IS NULL) + OR + (origin_kind = 'tool_mutation' AND parents_json <> '[]' AND operation_id IS NOT NULL + AND dispatch_event_id IS NOT NULL AND outcome_event_id IS NOT NULL + AND base_head_revision IS NOT NULL AND execution_profile_digest IS NOT NULL) + ) + ); + + INSERT INTO runtime_workspace_versions ( + workspace_version_id, repository_id, workspace_id, workspace_epoch_id, + object_format, origin_kind, origin_event_id, parents_json, + operation_id, dispatch_event_id, outcome_event_id, base_head_revision, + execution_profile_digest, commit_oid, tree_oid, policy_hash, + tree_delta_digest, changed_file_count, deleted_file_count, + accepted_event_id, protocol_version, committed_at + ) + SELECT + workspace_version_id, repository_id, workspace_id, workspace_epoch_id, + object_format, origin_kind, origin_event_id, parents_json, + NULL, NULL, NULL, NULL, NULL, commit_oid, tree_oid, policy_hash, + tree_delta_digest, changed_file_count, deleted_file_count, + accepted_event_id, protocol_version, committed_at + FROM runtime_workspace_versions_v12; + + CREATE TABLE runtime_workspace_heads ( + workspace_id TEXT NOT NULL, + workspace_epoch_id TEXT NOT NULL, + repository_id TEXT NOT NULL, + workspace_version_id TEXT NOT NULL, + accepted_event_id TEXT NOT NULL, + commit_oid TEXT NOT NULL, + tree_oid TEXT NOT NULL, + revision INTEGER NOT NULL CHECK (revision > 0), + PRIMARY KEY (workspace_id, workspace_epoch_id), + FOREIGN KEY (workspace_id, workspace_epoch_id) + REFERENCES runtime_workspace_epochs(workspace_id, workspace_epoch_id), + FOREIGN KEY ( + workspace_id, + workspace_epoch_id, + workspace_version_id, + accepted_event_id + ) REFERENCES runtime_workspace_versions( + workspace_id, + workspace_epoch_id, + workspace_version_id, + accepted_event_id + ) + ); + + INSERT INTO runtime_workspace_heads + SELECT * FROM runtime_workspace_heads_v12; + + DROP TABLE runtime_workspace_heads_v12; + DROP TABLE runtime_workspace_versions_v12; + `, + ], ]); export function configureSqliteRuntimeDatabase(db: DatabaseSync): void { diff --git a/packages/storage/src/sqlite-runtime-store.ts b/packages/storage/src/sqlite-runtime-store.ts index b94421ff1c..3c10b02f20 100644 --- a/packages/storage/src/sqlite-runtime-store.ts +++ b/packages/storage/src/sqlite-runtime-store.ts @@ -6,10 +6,12 @@ import type { DatabaseSync, SQLInputValue } from 'node:sqlite'; import { isDeepStrictEqual } from 'node:util'; import { buildWorkspaceBaselineAuthorityEvents, + buildWorkspaceSuccessorAuthorityEvent, scanWorkspaceBaselineAuthority, WORKSPACE_AUTHORITY_SESSION_ID, WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1, type ScannedWorkspaceBaselineAuthority, + type ScannedWorkspaceSuccessorAuthority, type WorkspaceAuthorityLedgerRow, type WorkspaceBaselineAuthorityInput, type WorkspaceBaselineCommitResult, @@ -71,7 +73,11 @@ import { RUNTIME_WORKSPACE_VERSION_AUTHORITY_CAPABILITY_VERSION, SQLITE_RUNTIME_SCHEMA_VERSION, } from './sqlite-runtime-schema.js'; -import { registerWorkspaceBaselineAuthorityWriterInternal } from './workspace-version-authority-internal.js'; +import { + registerWorkspaceBaselineAuthorityWriterInternal, + type WorkspaceSuccessorCommitInput, + type WorkspaceSuccessorCommitResult, +} from './workspace-version-authority-internal.js'; import type { ConversationCopyRuntimeEventBatch, ImmutableSteeringMessageProof, @@ -142,6 +148,9 @@ export type SqliteRuntimeStoreFailpoint = | 'after_workspace_epoch_projection_insert' | 'after_workspace_version_projection_insert' | 'after_workspace_head_projection_insert' + | 'after_workspace_successor_event_insert' + | 'after_workspace_successor_projection_insert' + | 'after_workspace_successor_head_update' | 'after_workspace_canonical_scan'; export interface SqliteRuntimeStoreOptions { @@ -1147,14 +1156,15 @@ export class SqliteRuntimeStore const events = buildWorkspaceBaselineAuthorityEvents(input); return this.transaction(() => { this.#assertWorkspaceStorageRootBinding(rootId); - const existingBaselines = this.readCanonicalWorkspaceBaselinesSync(); + const existingAuthority = this.readCanonicalWorkspaceAuthoritySync(); + const existingBaselines = existingAuthority.baselines; const existing = existingBaselines.find( (candidate) => candidate.epoch.workspaceId === input.epoch.workspaceId && candidate.epoch.workspaceEpochId === input.epoch.workspaceEpochId, ); if (existing) { - this.assertWorkspaceProjectionsMatchSync(existingBaselines); + this.assertWorkspaceProjectionsMatchSync(existingAuthority); if ( !isDeepStrictEqual( [ @@ -1166,11 +1176,15 @@ export class SqliteRuntimeStore ) { throw new Error('Workspace baseline authority conflict'); } - return { created: false, head: workspaceHeadRecord(existing) }; + const head = existingAuthority.heads.find( + (candidate) => candidate.workspaceEpochId === input.epoch.workspaceEpochId, + ); + if (!head) throw new Error('Workspace baseline authority head is unavailable'); + return { created: false, head }; } if (this.workspaceProjectionCountSync() !== 0 || existingBaselines.length !== 0) { - this.assertWorkspaceProjectionsMatchSync(existingBaselines); + this.assertWorkspaceProjectionsMatchSync(existingAuthority); } this.assertWorkspaceAuthorityStreamIsEmpty(events.epochOpenedEvent); this.assertInvocationIdentity([events.epochOpenedEvent, events.baselineAcceptedEvent]); @@ -1193,19 +1207,175 @@ export class SqliteRuntimeStore } this.options.failpoint?.('after_workspace_version_event_insert'); - const scanned = this.readCanonicalWorkspaceBaselinesSync(); - const accepted = scanned.find( + const scanned = this.readCanonicalWorkspaceAuthoritySync(); + const accepted = scanned.baselines.find( (candidate) => candidate.epoch.workspaceEpochId === input.epoch.workspaceEpochId, ); if (!accepted) throw new Error('Workspace baseline authority scan lost the committed epoch'); this.insertWorkspaceEpochProjection(accepted, input.committedAt); this.options.failpoint?.('after_workspace_epoch_projection_insert'); - this.insertWorkspaceVersionProjection(accepted, input.committedAt); + this.insertWorkspaceBaselineVersionProjection(accepted, input.committedAt); this.options.failpoint?.('after_workspace_version_projection_insert'); - this.insertWorkspaceHeadProjection(accepted); + const acceptedHead = scanned.heads.find( + (candidate) => candidate.workspaceEpochId === input.epoch.workspaceEpochId, + ); + if (!acceptedHead) throw new Error('Workspace baseline authority scan lost its head'); + this.insertWorkspaceHeadProjection(acceptedHead); this.options.failpoint?.('after_workspace_head_projection_insert'); this.assertWorkspaceProjectionsMatchSync(scanned); - return { created: true, head: workspaceHeadRecord(accepted) }; + const head = scanned.heads.find( + (candidate) => candidate.workspaceEpochId === input.epoch.workspaceEpochId, + ); + if (!head) throw new Error('Workspace baseline authority scan lost the committed head'); + return { created: true, head }; + }); + } + + async #commitWorkspaceSuccessor( + input: WorkspaceSuccessorCommitInput, + rootId: string, + ): Promise { + const toolOutcome: CommitToolOutcomeInput = { + ...input.toolOutcome, + runtimeEvent: canonicalizeRuntimeEventForStorage(input.toolOutcome.runtimeEvent), + }; + assertNoReservedWorkspaceAuthorityAppend(toolOutcome.runtimeEvent); + assertOutcomeInput(toolOutcome); + const successorEvent = buildWorkspaceSuccessorAuthorityEvent(input.successor); + if ( + input.successor.origin.operationId !== toolOutcome.operationId || + input.successor.origin.outcomeEventId !== toolOutcome.runtimeEvent.id + ) { + throw new Error('Workspace successor does not match its tool outcome identity'); + } + if ( + toolOutcome.runtimeEvent.content?.kind !== 'function_response' || + toolOutcome.runtimeEvent.content.isError === true + ) { + throw new Error('Workspace successor requires a successful tool outcome'); + } + + return this.transaction(() => { + this.#assertWorkspaceStorageRootBinding(rootId); + const before = this.readCanonicalWorkspaceAuthoritySync(); + this.assertWorkspaceProjectionsMatchSync(before); + const currentHead = before.heads.find( + (candidate) => + candidate.workspaceId === input.successor.successor.workspaceId && + candidate.workspaceEpochId === input.successor.successor.workspaceEpochId, + ); + if (!currentHead) throw new Error('Workspace successor base head is unavailable'); + + const existing = before.successors.find( + (candidate) => + candidate.acceptedEventId === input.successor.acceptedEventId || + candidate.successor.workspaceVersionId === input.successor.successor.workspaceVersionId, + ); + if (existing) { + assertStoredRuntimeEventEquals( + successorEvent, + this.readRuntimeEventJson(successorEvent.id), + ); + const operation = this.readToolOperationSync(toolOutcome.operationId); + if (!operation?.resultEventId) { + throw new Error('Workspace successor exists without its tool outcome'); + } + assertStoredRuntimeEventEquals( + toolOutcome.runtimeEvent, + this.readRuntimeEventJson(operation.resultEventId), + ); + return { + created: false, + head: { + repositoryId: existing.successor.repositoryId, + workspaceId: existing.successor.workspaceId, + workspaceEpochId: existing.successor.workspaceEpochId, + workspaceVersionId: existing.successor.workspaceVersionId, + acceptedEventId: existing.acceptedEventId, + commitOid: existing.successor.commitOid, + treeOid: existing.successor.treeOid, + revision: existing.successor.baseHeadRevision + 1, + }, + outcomeRuntimeEventSeq: this.runtimeEventSeq(operation.resultEventId), + }; + } + + const successor = input.successor.successor; + if ( + successor.repositoryId !== currentHead.repositoryId || + successor.parentWorkspaceVersionId !== currentHead.workspaceVersionId || + successor.baseAcceptedEventId !== currentHead.acceptedEventId || + successor.baseHeadRevision !== currentHead.revision + ) { + throw new Error('Workspace successor compare-and-set base head conflict'); + } + const operation = this.readToolOperationSync(toolOutcome.operationId); + if ( + !operation || + operation.currentState !== 'prepared' || + operation.resultEventId !== undefined || + operation.dispatchEventId !== input.successor.origin.dispatchEventId || + operation.recoveryMode !== 'reconcile' || + (operation.toolName !== 'Write' && operation.toolName !== 'Edit') + ) { + throw new Error('Workspace successor requires one prepared Write/Edit reconcile operation'); + } + + const outcomeResult = this.commitToolOutcomeSync(toolOutcome); + const successorSeq = this.insertRuntimeEvent( + successorEvent, + input.successor.committedAt, + false, + ); + if (successorSeq !== currentHead.revision + 2) { + throw new Error('Workspace successor fact is not the next authority event'); + } + this.options.failpoint?.('after_workspace_successor_event_insert'); + + const after = this.readCanonicalWorkspaceAuthoritySync(); + const accepted = after.successors.find( + (candidate) => candidate.acceptedEventId === input.successor.acceptedEventId, + ); + const nextHead = after.heads.find( + (candidate) => + candidate.workspaceId === successor.workspaceId && + candidate.workspaceEpochId === successor.workspaceEpochId, + ); + if (!accepted || !nextHead) { + throw new Error('Workspace successor authority scan lost the committed version'); + } + this.insertWorkspaceSuccessorVersionProjection(accepted, input.successor.committedAt); + this.options.failpoint?.('after_workspace_successor_projection_insert'); + const updated = this.db + .prepare(` + UPDATE runtime_workspace_heads + SET workspace_version_id = ?, accepted_event_id = ?, commit_oid = ?, tree_oid = ?, + revision = ? + WHERE workspace_id = ? AND workspace_epoch_id = ? + AND workspace_version_id = ? AND accepted_event_id = ? AND revision = ? + `) + .run( + nextHead.workspaceVersionId, + nextHead.acceptedEventId, + nextHead.commitOid, + nextHead.treeOid, + nextHead.revision, + currentHead.workspaceId, + currentHead.workspaceEpochId, + currentHead.workspaceVersionId, + currentHead.acceptedEventId, + currentHead.revision, + ); + if (updated.changes !== 1) { + throw new Error('Workspace successor head compare-and-set failed'); + } + this.options.failpoint?.('after_workspace_successor_head_update'); + this.assertWorkspaceProjectionsMatchSync(after); + return { + created: true, + head: nextHead, + outcomeRuntimeEventSeq: outcomeResult.runtimeEventSeq, + }; }); } @@ -1215,6 +1385,7 @@ export class SqliteRuntimeStore this, databasePath, (input, rootId) => this.#commitWorkspaceBaseline(input, rootId), + (input, rootId) => this.#commitWorkspaceSuccessor(input, rootId), (rootId) => this.#bindWorkspaceStorageRoot(rootId), readWorkspaceHead, ); @@ -1287,9 +1458,9 @@ export class SqliteRuntimeStore workspaceEpochId: string, ): Promise { return this.readTransaction(() => { - const baselines = this.readCanonicalWorkspaceBaselinesSync(); - this.assertWorkspaceProjectionsMatchSync(baselines); - const baseline = baselines.find( + const authority = this.readCanonicalWorkspaceAuthoritySync(); + this.assertWorkspaceProjectionsMatchSync(authority); + const baseline = authority.baselines.find( (candidate) => candidate.epoch.workspaceId === workspaceId && candidate.epoch.workspaceEpochId === workspaceEpochId, @@ -1302,12 +1473,16 @@ export class SqliteRuntimeStore workspaceVersionId: string, ): Promise { return this.readTransaction(() => { - const baselines = this.readCanonicalWorkspaceBaselinesSync(); - this.assertWorkspaceProjectionsMatchSync(baselines); - const baseline = baselines.find( + const authority = this.readCanonicalWorkspaceAuthoritySync(); + this.assertWorkspaceProjectionsMatchSync(authority); + const baseline = authority.baselines.find( (candidate) => candidate.baseline.workspaceVersionId === workspaceVersionId, ); - return baseline ? workspaceVersionRecord(baseline) : undefined; + if (baseline) return workspaceBaselineVersionRecord(baseline); + const successor = authority.successors.find( + (candidate) => candidate.successor.workspaceVersionId === workspaceVersionId, + ); + return successor ? workspaceSuccessorVersionRecord(successor) : undefined; }); } @@ -1316,42 +1491,46 @@ export class SqliteRuntimeStore workspaceEpochId: string, ): Promise { return this.readTransaction(() => { - const baselines = this.readCanonicalWorkspaceBaselinesSync(); - this.assertWorkspaceProjectionsMatchSync(baselines); - const baseline = baselines.find( + const authority = this.readCanonicalWorkspaceAuthoritySync(); + this.assertWorkspaceProjectionsMatchSync(authority); + return authority.heads.find( (candidate) => - candidate.epoch.workspaceId === workspaceId && - candidate.epoch.workspaceEpochId === workspaceEpochId, + candidate.workspaceId === workspaceId && candidate.workspaceEpochId === workspaceEpochId, ); - return baseline ? workspaceHeadRecord(baseline) : undefined; }); } async rebuildWorkspaceVersionProjections(): Promise { return this.transaction(() => { - const baselines = this.readCanonicalWorkspaceBaselinesSync(); + const authority = this.readCanonicalWorkspaceAuthoritySync(); this.db.prepare('DELETE FROM runtime_workspace_heads').run(); this.db.prepare('DELETE FROM runtime_workspace_versions').run(); this.db.prepare('DELETE FROM runtime_workspace_epochs').run(); - for (const baseline of baselines) { + for (const baseline of authority.baselines) { const committedAt = Math.max( this.runtimeEventCommittedAt(baseline.epochOpenedEventId), this.runtimeEventCommittedAt(baseline.baselineAcceptedEventId), ); this.insertWorkspaceEpochProjection(baseline, committedAt); - this.insertWorkspaceVersionProjection(baseline, committedAt); - this.insertWorkspaceHeadProjection(baseline); + this.insertWorkspaceBaselineVersionProjection(baseline, committedAt); + } + for (const successor of authority.successors) { + this.insertWorkspaceSuccessorVersionProjection( + successor, + this.runtimeEventCommittedAt(successor.acceptedEventId), + ); } - this.assertWorkspaceProjectionsMatchSync(baselines); + for (const head of authority.heads) this.insertWorkspaceHeadProjection(head); + this.assertWorkspaceProjectionsMatchSync(authority); return { - epochs: baselines.length, - versions: baselines.length, - heads: baselines.length, + epochs: authority.baselines.length, + versions: authority.baselines.length + authority.successors.length, + heads: authority.heads.length, }; }); } - private readCanonicalWorkspaceBaselinesSync() { + private readCanonicalWorkspaceAuthoritySync() { const partial = this.db .prepare(` SELECT stream_key FROM runtime_partial_snapshots @@ -1371,8 +1550,9 @@ export class SqliteRuntimeStore ORDER BY invocation_id ASC, event_seq ASC, event_id ASC `) .all() as unknown as RuntimeEventPrefixStorageRow[]; - const authorityRows: WorkspaceAuthorityLedgerRow[] = rows.map((row) => ({ - event: decodeRuntimeEventStorageRow(row), + const events = rows.map(decodeRuntimeEventStorageRow); + const authorityRows: WorkspaceAuthorityLedgerRow[] = rows.map((row, index) => ({ + event: events[index]!, eventSeq: row.event_seq, })); const scan = scanWorkspaceBaselineAuthority(authorityRows); @@ -1382,8 +1562,34 @@ export class SqliteRuntimeStore `Corrupt workspace RuntimeEvent authority: ${issue.code} at ${issue.eventId}`, ); } + const toolScan = scanToolLedger(events); + for (const accepted of scan.successors) { + const origin = accepted.successor.origin; + const operation = toolScan.operations.find( + (candidate) => candidate.operationId === origin.operationId, + ); + const dispatch = operation?.dispatchEvent?.actions?.toolDispatch; + const response = operation?.responseEvent; + if ( + !operation || + operation.issues.length > 0 || + operation.dispatchEvent?.id !== origin.dispatchEventId || + !dispatch || + dispatch.operationId !== origin.operationId || + dispatch.recoveryMode !== 'reconcile' || + (dispatch.toolName !== 'Write' && dispatch.toolName !== 'Edit') || + !response || + response.id !== origin.outcomeEventId || + response.content?.kind !== 'function_response' || + response.content.isError === true + ) { + throw new Error( + `Corrupt workspace successor tool evidence: identity_conflict at ${accepted.acceptedEventId}`, + ); + } + } this.options.failpoint?.('after_workspace_canonical_scan'); - return scan.baselines; + return scan; } private assertWorkspaceAuthorityStreamIsEmpty(event: RuntimeEvent): void { @@ -1452,7 +1658,7 @@ export class SqliteRuntimeStore ); } - private insertWorkspaceVersionProjection( + private insertWorkspaceBaselineVersionProjection( accepted: ReturnType['baselines'][number], committedAt: number, ): void { @@ -1468,6 +1674,11 @@ export class SqliteRuntimeStore origin_kind, origin_event_id, parents_json, + operation_id, + dispatch_event_id, + outcome_event_id, + base_head_revision, + execution_profile_digest, commit_oid, tree_oid, policy_hash, @@ -1477,7 +1688,8 @@ export class SqliteRuntimeStore accepted_event_id, protocol_version, committed_at - ) VALUES (?, ?, ?, ?, ?, 'baseline', ?, '[]', ?, ?, ?, ?, ?, ?, ?, 1, ?) + ) VALUES (?, ?, ?, ?, ?, 'baseline', ?, '[]', NULL, NULL, NULL, NULL, NULL, + ?, ?, ?, ?, ?, ?, ?, 1, ?) `) .run( baseline.workspaceVersionId, @@ -1497,10 +1709,63 @@ export class SqliteRuntimeStore ); } - private insertWorkspaceHeadProjection( - accepted: ReturnType['baselines'][number], + private insertWorkspaceSuccessorVersionProjection( + accepted: ScannedWorkspaceSuccessorAuthority, + committedAt: number, ): void { - const head = workspaceHeadRecord(accepted); + const { successor } = accepted; + this.db + .prepare(` + INSERT INTO runtime_workspace_versions ( + workspace_version_id, + repository_id, + workspace_id, + workspace_epoch_id, + object_format, + origin_kind, + origin_event_id, + parents_json, + operation_id, + dispatch_event_id, + outcome_event_id, + base_head_revision, + execution_profile_digest, + commit_oid, + tree_oid, + policy_hash, + tree_delta_digest, + changed_file_count, + deleted_file_count, + accepted_event_id, + protocol_version, + committed_at + ) VALUES (?, ?, ?, ?, ?, 'tool_mutation', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?) + `) + .run( + successor.workspaceVersionId, + successor.repositoryId, + successor.workspaceId, + successor.workspaceEpochId, + successor.objectFormat, + successor.origin.outcomeEventId, + JSON.stringify(successor.parents), + successor.origin.operationId, + successor.origin.dispatchEventId, + successor.origin.outcomeEventId, + successor.baseHeadRevision, + successor.executionProfileDigest, + successor.commitOid, + successor.treeOid, + successor.policyHash, + successor.treeDeltaDigest, + successor.changedFileCount, + successor.deletedFileCount, + accepted.acceptedEventId, + committedAt, + ); + } + + private insertWorkspaceHeadProjection(head: WorkspaceHeadRecordV1): void { this.db .prepare(` INSERT INTO runtime_workspace_heads ( @@ -1527,15 +1792,18 @@ export class SqliteRuntimeStore } private assertWorkspaceProjectionsMatchSync( - baselines: ReturnType['baselines'], + authority: ReturnType, ): void { - const expectedEpochs = baselines + const expectedEpochs = authority.baselines .map(workspaceEpochProjectionRow) .sort(compareWorkspaceEpochRow); - const expectedVersions = baselines - .map(workspaceVersionProjectionRow) - .sort(compareWorkspaceVersionRow); - const expectedHeads = baselines.map(workspaceHeadProjectionRow).sort(compareWorkspaceHeadRow); + const expectedVersions = [ + ...authority.baselines.map(workspaceBaselineVersionProjectionRow), + ...authority.successors.map(workspaceSuccessorVersionProjectionRow), + ].sort(compareWorkspaceVersionRow); + const expectedHeads = authority.heads + .map(workspaceHeadProjectionRow) + .sort(compareWorkspaceHeadRow); const epochs = ( this.db .prepare(` @@ -1578,6 +1846,11 @@ export class SqliteRuntimeStore origin_kind, origin_event_id, parents_json, + operation_id, + dispatch_event_id, + outcome_event_id, + base_head_revision, + execution_profile_digest, commit_oid, tree_oid, policy_hash, @@ -3341,6 +3614,11 @@ interface WorkspaceVersionProjectionRow { origin_kind: string; origin_event_id: string; parents_json: string; + operation_id: string | null; + dispatch_event_id: string | null; + outcome_event_id: string | null; + base_head_revision: number | null; + execution_profile_digest: string | null; commit_oid: string; tree_oid: string; policy_hash: string; @@ -3374,26 +3652,19 @@ function workspaceEpochRecord( }; } -function workspaceVersionRecord( - authority: ScannedWorkspaceBaselineAuthority, -): WorkspaceVersionRecordV1 { +function workspaceBaselineVersionRecord(authority: ScannedWorkspaceBaselineAuthority) { return { ...authority.baseline, - baselineAcceptedEventId: authority.baselineAcceptedEventId, + acceptedEventId: authority.baselineAcceptedEventId, committedAt: authority.baselineAcceptedAt, }; } -function workspaceHeadRecord(authority: ScannedWorkspaceBaselineAuthority): WorkspaceHeadRecordV1 { +function workspaceSuccessorVersionRecord(authority: ScannedWorkspaceSuccessorAuthority) { return { - repositoryId: authority.epoch.repositoryId, - workspaceId: authority.epoch.workspaceId, - workspaceEpochId: authority.epoch.workspaceEpochId, - workspaceVersionId: authority.baseline.workspaceVersionId, - acceptedEventId: authority.baselineAcceptedEventId, - commitOid: authority.baseline.commitOid, - treeOid: authority.baseline.treeOid, - revision: 1, + ...authority.successor, + acceptedEventId: authority.acceptedEventId, + committedAt: authority.acceptedAt, }; } @@ -3424,10 +3695,10 @@ function workspaceEpochProjectionRow( }; } -function workspaceVersionProjectionRow( +function workspaceBaselineVersionProjectionRow( authority: ScannedWorkspaceBaselineAuthority, ): WorkspaceVersionProjectionRow { - const record = workspaceVersionRecord(authority); + const record = workspaceBaselineVersionRecord(authority); return { workspace_version_id: record.workspaceVersionId, repository_id: record.repositoryId, @@ -3437,22 +3708,54 @@ function workspaceVersionProjectionRow( origin_kind: record.origin.kind, origin_event_id: record.origin.epochOpenedEventId, parents_json: '[]', + operation_id: null, + dispatch_event_id: null, + outcome_event_id: null, + base_head_revision: null, + execution_profile_digest: null, commit_oid: record.commitOid, tree_oid: record.treeOid, policy_hash: record.policyHash, tree_delta_digest: record.treeDeltaDigest, changed_file_count: record.changedFileCount, deleted_file_count: record.deletedFileCount, - accepted_event_id: record.baselineAcceptedEventId, + accepted_event_id: record.acceptedEventId, protocol_version: 1, committed_at: record.committedAt, }; } -function workspaceHeadProjectionRow( - authority: ScannedWorkspaceBaselineAuthority, -): WorkspaceHeadProjectionRow { - const record = workspaceHeadRecord(authority); +function workspaceSuccessorVersionProjectionRow( + authority: ScannedWorkspaceSuccessorAuthority, +): WorkspaceVersionProjectionRow { + const record = workspaceSuccessorVersionRecord(authority); + return { + workspace_version_id: record.workspaceVersionId, + repository_id: record.repositoryId, + workspace_id: record.workspaceId, + workspace_epoch_id: record.workspaceEpochId, + object_format: record.objectFormat, + origin_kind: record.origin.kind, + origin_event_id: record.origin.outcomeEventId, + parents_json: JSON.stringify(record.parents), + operation_id: record.origin.operationId, + dispatch_event_id: record.origin.dispatchEventId, + outcome_event_id: record.origin.outcomeEventId, + base_head_revision: record.baseHeadRevision, + execution_profile_digest: record.executionProfileDigest, + commit_oid: record.commitOid, + tree_oid: record.treeOid, + policy_hash: record.policyHash, + tree_delta_digest: record.treeDeltaDigest, + changed_file_count: record.changedFileCount, + deleted_file_count: record.deletedFileCount, + accepted_event_id: record.acceptedEventId, + protocol_version: 1, + committed_at: record.committedAt, + }; +} + +function workspaceHeadProjectionRow(record: WorkspaceHeadRecordV1): WorkspaceHeadProjectionRow { return { workspace_id: record.workspaceId, workspace_epoch_id: record.workspaceEpochId, diff --git a/packages/storage/src/workspace-version-authority-internal.ts b/packages/storage/src/workspace-version-authority-internal.ts index 5c2fc998a8..f9c7f93a93 100644 --- a/packages/storage/src/workspace-version-authority-internal.ts +++ b/packages/storage/src/workspace-version-authority-internal.ts @@ -2,7 +2,9 @@ import type { WorkspaceBaselineAuthorityInput, WorkspaceBaselineCommitResult, WorkspaceHeadRecordV1, + WorkspaceSuccessorAuthorityInput, } from '@maka/core/workspace-version-authority'; +import type { RuntimeEvent } from '@maka/core/runtime-event'; import { lstatSync } from 'node:fs'; import { realpath } from 'node:fs/promises'; import { basename, dirname, join, normalize, resolve } from 'node:path'; @@ -13,6 +15,24 @@ type WorkspaceBaselineAuthorityWriter = ( rootId: string, ) => Promise; type WorkspaceStorageRootBinder = (rootId: string) => void; +export interface WorkspaceSuccessorCommitInput { + successor: WorkspaceSuccessorAuthorityInput; + toolOutcome: { + operationId: string; + journalEventId: string; + runtimeEvent: RuntimeEvent; + committedAt: number; + }; +} +export interface WorkspaceSuccessorCommitResult { + created: boolean; + head: WorkspaceHeadRecordV1; + outcomeRuntimeEventSeq: number; +} +type WorkspaceSuccessorAuthorityWriter = ( + input: WorkspaceSuccessorCommitInput, + rootId: string, +) => Promise; type WorkspaceHeadReader = ( workspaceId: string, workspaceEpochId: string, @@ -20,6 +40,7 @@ type WorkspaceHeadReader = ( interface WorkspaceBaselineAuthorityRegistration { readonly writer: WorkspaceBaselineAuthorityWriter; + readonly successorWriter: WorkspaceSuccessorAuthorityWriter; readonly readHead: WorkspaceHeadReader; readonly bindStorageRoot: WorkspaceStorageRootBinder; readonly databasePath: string; @@ -36,6 +57,7 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( store: object, databasePath: string, writer: WorkspaceBaselineAuthorityWriter, + successorWriter: WorkspaceSuccessorAuthorityWriter, bindStorageRoot: WorkspaceStorageRootBinder, readHead: WorkspaceHeadReader, ): void { @@ -45,6 +67,7 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( const resolvedDatabasePath = resolve(databasePath); workspaceBaselineAuthorityWriters.set(store, { writer, + successorWriter, readHead, bindStorageRoot, databasePath: resolvedDatabasePath, @@ -80,6 +103,18 @@ export function commitWorkspaceBaselineInternal( return registration.writer(input, registration.boundRootId); } +export function commitWorkspaceSuccessorInternal( + store: object, + input: WorkspaceSuccessorCommitInput, +): Promise { + const registration = workspaceBaselineAuthorityWriters.get(store); + if (!registration) throw new Error('Workspace successor authority writer is unavailable'); + if (!registration.boundRootId) { + throw new Error('Workspace successor authority store has no durable storage-root binding'); + } + return registration.successorWriter(input, registration.boundRootId); +} + export function bindWorkspaceBaselineAuthorityStoreRootInternal( store: object, rootId: string,