diff --git a/.github/workflows/gitoxide-helper-admission.yml b/.github/workflows/gitoxide-helper-admission.yml index 065686a705..070a80368e 100644 --- a/.github/workflows/gitoxide-helper-admission.yml +++ b/.github/workflows/gitoxide-helper-admission.yml @@ -26,8 +26,14 @@ on: - 'packages/runtime-host/src/__tests__/gitoxide-helper-*.test.ts' - 'packages/runtime-host/src/server/gitoxide-repository-admission-authority-internal.ts' - 'packages/runtime-host/src/server/gitoxide-mutation-candidate-receipt-authority-internal.ts' + - 'packages/runtime-host/src/server/gitoxide-managed-write-edit-owner-internal.ts' - 'packages/runtime-host/src/__tests__/gitoxide-repository-admission-authority-internal.test.ts' + - 'packages/runtime-host/src/__tests__/gitoxide-managed-write-edit-owner-internal.test.ts' - 'packages/runtime-host/src/__tests__/fixtures/gitoxide-candidate-receipt-crash-child.ts' + - 'packages/runtime-host/src/__tests__/fixtures/gitoxide-managed-write-edit-owner-crash-child.ts' + - 'packages/storage/src/execution-stores-workspace-authority-internal.ts' + - 'packages/storage/src/workspace-version-authority-internal.ts' + - 'packages/storage/src/sqlite-runtime-store.ts' - 'packages/runtime/package.json' - 'docs/architecture/gitoxide-*.md' push: @@ -40,8 +46,14 @@ on: - 'packages/runtime-host/src/__tests__/gitoxide-helper-*.test.ts' - 'packages/runtime-host/src/server/gitoxide-repository-admission-authority-internal.ts' - 'packages/runtime-host/src/server/gitoxide-mutation-candidate-receipt-authority-internal.ts' + - 'packages/runtime-host/src/server/gitoxide-managed-write-edit-owner-internal.ts' - 'packages/runtime-host/src/__tests__/gitoxide-repository-admission-authority-internal.test.ts' + - 'packages/runtime-host/src/__tests__/gitoxide-managed-write-edit-owner-internal.test.ts' - 'packages/runtime-host/src/__tests__/fixtures/gitoxide-candidate-receipt-crash-child.ts' + - 'packages/runtime-host/src/__tests__/fixtures/gitoxide-managed-write-edit-owner-crash-child.ts' + - 'packages/storage/src/execution-stores-workspace-authority-internal.ts' + - 'packages/storage/src/workspace-version-authority-internal.ts' + - 'packages/storage/src/sqlite-runtime-store.ts' - 'packages/runtime/package.json' - 'docs/architecture/gitoxide-*.md' @@ -93,3 +105,4 @@ jobs: packages/runtime-host/dist/__tests__/gitoxide-helper-artifact-authority-internal.test.js packages/runtime-host/dist/__tests__/gitoxide-helper-invocation-internal.test.js packages/runtime-host/dist/__tests__/gitoxide-repository-admission-authority-internal.test.js + packages/runtime-host/dist/__tests__/gitoxide-managed-write-edit-owner-internal.test.js diff --git a/docs/architecture/gitoxide-managed-write-edit-owner-v1.zh-CN.md b/docs/architecture/gitoxide-managed-write-edit-owner-v1.zh-CN.md new file mode 100644 index 0000000000..46ff050fbc --- /dev/null +++ b/docs/architecture/gitoxide-managed-write-edit-owner-v1.zh-CN.md @@ -0,0 +1,72 @@ +# Gitoxide managed Write/Edit owner v1 + +## 主要不变量 + +一次 managed Write/Edit 在 T1 前绑定 durable workspace epoch、canonical head、单一路径和固定执行 +profile。T1 后只能收敛到以下三类 durable 结果之一: + +- 无变化或已证明无副作用的失败,由 SQLite 原子提交 terminal T2 并释放 reservation; +- exact result 被固化为 operation-bound Gitoxide candidate,由 SQLite 原子提交 T2、successor 和 + canonical head,再将 `refs/maka/accepted` CAS 到该 candidate; +- 证据不完整或状态无法确定时保持 unsettled,禁止 generic T2、fallback 和工具副作用重放。 + +## Owner 与权限 + +- Runtime 独占原始 Write/Edit 参数和 provider result。Host 只能读取 Runtime-issued immutable + operation proof,不能返回或替换 provider result。 +- Execution Stores authority 从 durable epoch 读取 `workspaceInstanceId`。调用者不能自报 reservation + identity。 +- Gitoxide repository authority 只从 exact accepted commit/tree 读取 base content;Write/Edit 是 + `F(immutable base, Runtime-owned args)` 的纯变换,不读写 live checkout。 +- Candidate receipt authority 拥有 candidate ref、receipt 和 exact retry。 +- SQLite workspace authority 是 accepted truth owner。只有成功的 successor transaction 才签发 + owner-bound projection capability。 +- Gitoxide helper 只消费该 capability 做 accepted-ref CAS;裸路径、OID、receipt 或 caller object 都不能 + 推进 accepted ref。 + +## 原子性与恢复边界 + +Git ref 与 SQLite 无法组成一个物理事务,因此 v1 采用两个有序线性化点: + +1. SQLite transaction 提交 tool outcome、workspace successor、head CAS,并释放 T1 reservation; +2. Gitoxide helper 将 accepted ref 从 exact base CAS 到 exact candidate。 + +进程若死在两者之间,新 owner 按 operation ID 从 `tool_operations` 主键和 RuntimeEvent event ID 读取 +exact call、dispatch、outcome。该读取是有界主键查询,不需要 schema 15 表达式索引,也不扫描完整账本。 +恢复随后: + +1. 以 durable parent head 重开 accepted repository; +2. 从 immutable base 和 durable call args 重新计算 pure transform; +3. exact-replay candidate receipt/ref; +4. exact-replay 已提交的 SQLite successor,以重新签发 process-local projection capability; +5. 重放 accepted-ref CAS。 + +恢复不会调用普通工具实现,也不会产生第二次文件系统副作用。 + +## 失败状态与回滚 + +| 状态 | 行为 | +| --- | --- | +| T1 前路径、epoch、head、version 或 helper admission 不匹配 | 拒绝进入 T1 | +| T1 后 Runtime 证明 no-change/no-effect failure | 原子 terminal T2,释放 reservation,不推进 head | +| candidate/receipt 与 base、path、content 或 profile 不匹配 | unsettled;保留 reservation 或 accepted truth,禁止覆盖 | +| SQLite successor 未提交 | 不签发 projection capability,candidate 只是未接受 artifact | +| SQLite successor 已提交、accepted ref 仍为 base | 重建 proof 并重放 CAS,不重跑工具 | +| accepted ref 已为 candidate | exact replay success | +| accepted ref 为第三值或 durable evidence 不一致 | fail closed,不 reset、不自动覆盖 | + +## 平台能力矩阵 + +| 平台 | v1 承诺 | +| --- | --- | +| Linux | short-lived Gitoxide helper、SQLite crash convergence、exact candidate/ref replay | +| macOS | 与 Linux 相同的进程崩溃合同;不承诺 `fsync` 之外的断电语义 | +| Windows | 相同的 SQLite/ref 协议;successor-commit 后的子进程 exit/reopen 已进入三平台 Gitoxide helper workflow | + +三平台均不依赖 system Git、bundled Git CLI、linked worktree rotation 或 live checkout 写权限。 + +## 当前交付边界 + +本切片提供 owner 与真实 helper/SQLite 合同测试,但尚不自行创建 Desktop/CLI managed coding session。 +产品组合与 continuation 分别属于后续交付;它们只能消费这里签发的窄 capability,不能重新开放裸路径或 +caller-provided execution profile。 diff --git a/docs/architecture/runtime-recovery-post-gitoxide-rebuild-ledger.zh-CN.md b/docs/architecture/runtime-recovery-post-gitoxide-rebuild-ledger.zh-CN.md index 59b42909dd..dd1c829ee9 100644 --- a/docs/architecture/runtime-recovery-post-gitoxide-rebuild-ledger.zh-CN.md +++ b/docs/architecture/runtime-recovery-post-gitoxide-rebuild-ledger.zh-CN.md @@ -119,3 +119,14 @@ Gitoxide helper 与 npm producer 分别保留自己的 trust root。packaged Hos - 延期 projection quarantine GC 到 projection owner 自己的 lifecycle 交付; - 延期 dependency environment 与 Write/Edit 恢复的耦合;纯文件 transform 不应被 npm producer 阻塞; - 禁止为使旧 PR 可编译而恢复 v1/v2 fallback。 + +## 7. 当前重建状态 + +- R0:完成。Runtime 不再接受 Host 提交的 `providerResult`;运行时夹带字段也会被忽略。 +- R1:完成。candidate receipt/ref 支持 exact replay 与 source import retry。 +- R2:owner 主链已完成:durable epoch/head、pure transform、candidate、SQLite successor 与 accepted-ref + projection 已串联;真实子进程会在 SQLite successor 提交后直接退出,新 owner 只从 durable evidence + 重建 candidate 并推进 accepted ref。该用例已进入三平台 Gitoxide helper workflow。 +- R3:accepted-ref projection 已实现为 accepted truth 的派生 CAS;filesystem checkout projection 仍保持 + 独立延期,不参与 canonical Write/Edit read/write。 +- R4、R5:尚未在新基线上重建。 diff --git a/packages/runtime-host/src/__tests__/fixtures/gitoxide-managed-write-edit-owner-crash-child.ts b/packages/runtime-host/src/__tests__/fixtures/gitoxide-managed-write-edit-owner-crash-child.ts new file mode 100644 index 0000000000..9cf512d4fd --- /dev/null +++ b/packages/runtime-host/src/__tests__/fixtures/gitoxide-managed-write-edit-owner-crash-child.ts @@ -0,0 +1,184 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import { createHash } from 'node:crypto'; +import { readFile, realpath, stat } from 'node:fs/promises'; +import { canonicalToolArgsHash } from '@maka/core/tool-args-identity'; +import { openInteractiveExecutionStoresForWrite } from '@maka/storage/execution-stores'; +import { + discoverMarkedStorageRoot, + tryAcquireInteractiveRootOwner, +} from '@maka/storage/root-authority'; +import { + admitGitoxideHelperArtifactInternal, + issueGitoxideHelperReleaseArtifactClaimInternal, +} from '../../server/gitoxide-helper-artifact-authority-internal.js'; +import { createGitoxideManagedWriteEditOwnerInternal } from '../../server/gitoxide-managed-write-edit-owner-internal.js'; + +interface Fixture { + readonly storageRoot: string; + readonly repositoryPath: string; + readonly helperPath: string; + readonly workspaceId: string; + readonly workspaceEpochId: string; + readonly operationId: string; + readonly toolCallId: string; + readonly args: { + readonly path: string; + readonly content: string; + }; +} + +const fixturePath = process.argv[2]; +if (!fixturePath) throw new Error('Missing managed Write crash fixture path'); +const fixture = JSON.parse(await readFile(fixturePath, 'utf8')) as Fixture; +const rootCapability = await discoverMarkedStorageRoot({ path: fixture.storageRoot }); +if (rootCapability.kind !== 'interactive') throw new Error('Crash fixture root kind is invalid'); +const rootOwner = await tryAcquireInteractiveRootOwner(rootCapability); +if (!rootOwner) throw new Error('Crash fixture could not acquire the storage root'); +const stores = await openInteractiveExecutionStoresForWrite(rootOwner.lease); + +const helperPath = await realpath(fixture.helperPath); +const helperBytes = await readFile(helperPath); +const helperInfo = await stat(helperPath); +const releaseOwnerToken = {}; +const invocationOwnerToken = {}; +const claim = issueGitoxideHelperReleaseArtifactClaimInternal(releaseOwnerToken, { + executablePath: helperPath, + expectedSha256: `sha256:${createHash('sha256').update(helperBytes).digest('hex')}`, + expectedBytes: helperInfo.size, + platform: process.platform, + arch: process.arch, + protocolVersion: 1, + supportedOperations: [ + 'inspect_repository', + 'import_source_head', + 'create_candidate', + 'promote_candidate', + 'observe_accepted_ref', + 'read_tree_file', + ], +}); +const helperCapability = await admitGitoxideHelperArtifactInternal({ + releaseOwnerToken, + invocationOwnerToken, + claim, +}); +const owner = createGitoxideManagedWriteEditOwnerInternal({ + storageRootLease: rootOwner.lease, + stores, + invocationOwnerToken, + helperCapability, + repositoryPath: fixture.repositoryPath, + workspaceId: fixture.workspaceId, + workspaceEpochId: fixture.workspaceEpochId, + failpoint(point) { + if (point === 'after_workspace_successor_commit') process.exit(73); + }, +}); +const admission = await owner.admitManagedMutation({ + operationId: fixture.operationId, + toolName: 'Write', + persistedArgs: fixture.args, + abortSignal: new AbortController().signal, +}); +if (admission.immutableBase?.content !== 'before\n') { + throw new Error('Crash fixture opened a different immutable base'); +} +const identity = { + sessionId: 'session-real-write', + invocationId: 'invocation-real-write', + runId: 'run-real-write', + turnId: 'turn-real-write', +}; +const canonicalArgsHash = canonicalToolArgsHash('Write', fixture.args); +await stores.runtimeEventStore.commitToolPrepared({ + operationId: fixture.operationId, + journalEventId: `${fixture.operationId}_prepared`, + runtimeEvent: { + id: 'call-event-real-write', + ...identity, + ts: 2, + partial: false, + role: 'model', + author: 'agent', + content: { + kind: 'function_call', + id: fixture.toolCallId, + name: 'Write', + args: fixture.args, + }, + refs: { operationId: fixture.operationId, toolCallId: fixture.toolCallId }, + }, + dispatchRuntimeEvent: { + id: 'dispatch-real-write', + ...identity, + ts: 2, + partial: false, + role: 'system', + author: 'system', + actions: { + toolDispatch: { + protocol: 't1_after_preflight_v1', + operationId: fixture.operationId, + providerToolCallId: fixture.toolCallId, + toolName: 'Write', + canonicalArgsHash, + recoveryMode: 'reconcile', + managedMutation: admission.durableDispatch, + }, + }, + refs: { operationId: fixture.operationId, toolCallId: fixture.toolCallId }, + }, + providerToolCallId: fixture.toolCallId, + toolName: 'Write', + canonicalArgsHash, + recoveryMode: 'reconcile', + committedAt: 2, +}); +const durableOutcome = { + id: 'outcome-event-real-write', + ...identity, + ts: 3, + partial: false, + role: 'tool' as const, + author: 'tool' as const, + content: { + kind: 'function_response' as const, + id: fixture.toolCallId, + name: 'Write', + result: { + kind: 'json' as const, + value: { kind: 'file_write', path: fixture.args.path }, + }, + }, + refs: { operationId: fixture.operationId, toolCallId: fixture.toolCallId }, +}; +await admission.execute(async () => ({ + content: durableOutcome.content.result, + isError: false, + durationMs: 1, + mutationResult: { + path: fixture.args.path, + content: fixture.args.content, + changed: true, + }, + durableOutcome, +})); +throw new Error('Crash fixture did not stop after the durable successor commit'); diff --git a/packages/runtime-host/src/__tests__/gitoxide-managed-write-edit-owner-internal.test.ts b/packages/runtime-host/src/__tests__/gitoxide-managed-write-edit-owner-internal.test.ts new file mode 100644 index 0000000000..e141df3864 --- /dev/null +++ b/packages/runtime-host/src/__tests__/gitoxide-managed-write-edit-owner-internal.test.ts @@ -0,0 +1,315 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import assert from 'node:assert/strict'; +import { execFileSync, spawn } from 'node:child_process'; +import { createHash } from 'node:crypto'; +import { once } from 'node:events'; +import { mkdtemp, readFile, realpath, rm, stat, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import test, { type TestContext } from 'node:test'; +import { WORKSPACE_MATERIALIZATION_SEMANTICS_V1 } from '@maka/core/workspace-version-authority'; +import { openInteractiveExecutionStoresForWrite } from '@maka/storage/execution-stores'; +import { resolveStorageRoot, tryAcquireInteractiveRootOwner } from '@maka/storage/root-authority'; +import { commitExecutionStoresWorkspaceBaselineForTestInternal } from '@maka/storage/test-only/execution-stores-workspace-authority'; +import { + admitGitoxideHelperArtifactInternal, + type GitoxideHelperInvocationCapability, + issueGitoxideHelperReleaseArtifactClaimInternal, +} from '../server/gitoxide-helper-artifact-authority-internal.js'; +import { createGitoxideManagedWriteEditOwnerInternal } from '../server/gitoxide-managed-write-edit-owner-internal.js'; +import { + admitGitoxideRepositoryInternal, + importAdmittedGitoxideRepositoryInternal, +} from '../server/gitoxide-repository-admission-authority-internal.js'; + +test('rejects a non-canonical managed path before consulting Gitoxide', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-gitoxide-write-edit-owner-')); + const rootCapability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const rootOwner = await tryAcquireInteractiveRootOwner(rootCapability); + assert.ok(rootOwner); + const stores = await openInteractiveExecutionStoresForWrite(rootOwner.lease); + try { + const owner = createGitoxideManagedWriteEditOwnerInternal({ + storageRootLease: rootOwner.lease, + stores, + invocationOwnerToken: {}, + helperCapability: {} as GitoxideHelperInvocationCapability, + repositoryPath: join(root, 'managed.git'), + workspaceId: `workspace_${'1'.repeat(32)}`, + workspaceEpochId: `epoch_${'2'.repeat(32)}`, + }); + + await assert.rejects( + owner.admitManagedMutation({ + operationId: 'operation-path-alias', + toolName: 'Write', + persistedArgs: { path: 'dir\\file.txt', content: 'after\n' }, + abortSignal: new AbortController().signal, + }), + /path must already be canonical/i, + ); + } finally { + await stores.sessionStore.close?.(); + await rootOwner.close(); + await rm(root, { recursive: true, force: true }); + } +}); + +test('requires the durable workspace epoch before consulting Gitoxide', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-gitoxide-write-edit-epoch-')); + const rootCapability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const rootOwner = await tryAcquireInteractiveRootOwner(rootCapability); + assert.ok(rootOwner); + const stores = await openInteractiveExecutionStoresForWrite(rootOwner.lease); + try { + const owner = createGitoxideManagedWriteEditOwnerInternal({ + storageRootLease: rootOwner.lease, + stores, + invocationOwnerToken: {}, + helperCapability: {} as GitoxideHelperInvocationCapability, + repositoryPath: join(root, 'managed.git'), + workspaceId: `workspace_${'1'.repeat(32)}`, + workspaceEpochId: `epoch_${'2'.repeat(32)}`, + }); + + await assert.rejects( + owner.admitManagedMutation({ + operationId: 'operation-missing-epoch', + toolName: 'Write', + persistedArgs: { path: 'file.txt', content: 'after\n' }, + abortSignal: new AbortController().signal, + }), + /durable workspace epoch is unavailable/i, + ); + } finally { + await stores.sessionStore.close?.(); + await rootOwner.close(); + await rm(root, { recursive: true, force: true }); + } +}); + +test('reopens after a process crash and promotes the exact durable Write successor', async (t) => { + const helper = await admittedHelper(); + if (!helper) { + t.skip('MAKA_GITOXIDE_HELPER_PATH is required for the real owner contract test'); + return; + } + const sourceRepositoryPath = await createRepository(t); + await writeFile(join(sourceRepositoryPath, 'notes.txt'), 'before\n', 'utf8'); + git(sourceRepositoryPath, ['add', 'notes.txt']); + git(sourceRepositoryPath, [ + '-c', + 'user.name=Maka Test', + '-c', + 'user.email=maka@example.invalid', + 'commit', + '--quiet', + '-m', + 'fixture', + ]); + + const root = await realpath(await mkdtemp(join(tmpdir(), 'maka-gitoxide-write-edit-full-'))); + t.after(() => rm(root, { recursive: true, force: true })); + const rootCapability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const rootOwner = await tryAcquireInteractiveRootOwner(rootCapability); + assert.ok(rootOwner); + const stores = await openInteractiveExecutionStoresForWrite(rootOwner.lease); + + const admissionOwnerToken = {}; + const importedRepositoryOwnerToken = {}; + const admitted = await admitGitoxideRepositoryInternal({ + ...helper, + admissionOwnerToken, + repositoryPath: sourceRepositoryPath, + }); + assert.equal(admitted.kind, 'accepted'); + if (admitted.kind !== 'accepted') return; + const repositoryPath = join(root, 'managed.git'); + const imported = await importAdmittedGitoxideRepositoryInternal({ + admissionOwnerToken, + repositoryCapability: admitted.capability, + acceptedRepositoryOwnerToken: importedRepositoryOwnerToken, + destinationRepositoryPath: repositoryPath, + }); + const ids = { + repositoryId: `repository_${'1'.repeat(32)}`, + workspaceId: `workspace_${'2'.repeat(32)}`, + workspaceEpochId: `epoch_${'3'.repeat(32)}`, + workspaceInstanceId: `instance_${'4'.repeat(32)}`, + workspaceVersionId: `version_${'5'.repeat(32)}`, + } as const; + const baseline = await commitExecutionStoresWorkspaceBaselineForTestInternal(stores, { + epochOpenedEventId: 'gitoxide-write-edit-epoch', + baselineAcceptedEventId: 'gitoxide-write-edit-baseline', + committedAt: 1, + epoch: { + ...ids, + mode: 'managed_worktree', + objectFormat: 'sha1', + sourceCommitOid: imported.sourceHeadCommitOid, + sourceTreeOid: imported.sourceTreeOid, + materializationProfileDigest: sha256('gitoxide-materialization-v1'), + materializationSemantics: WORKSPACE_MATERIALIZATION_SEMANTICS_V1, + policyHash: sha256('managed-tree-policy-v3'), + }, + baseline: { + workspaceVersionId: ids.workspaceVersionId, + commitOid: imported.baselineCommitOid, + treeOid: imported.baselineTreeOid, + treeDeltaDigest: sha256('gitoxide-baseline-delta-v1'), + changedFileCount: 1, + deletedFileCount: 0, + }, + }); + const operationId = 'operation-real-write'; + const toolCallId = 'call-real-write'; + const args = { path: 'notes.txt', content: 'after\n' }; + const fixturePath = join(root, 'managed-write-crash-fixture.json'); + await writeFile( + fixturePath, + `${JSON.stringify({ + storageRoot: root, + repositoryPath, + helperPath: helper.helperPath, + workspaceId: ids.workspaceId, + workspaceEpochId: ids.workspaceEpochId, + operationId, + toolCallId, + args, + })}\n`, + 'utf8', + ); + await stores.sessionStore.close?.(); + await rootOwner.close(); + + const child = spawn( + process.execPath, + [ + join(import.meta.dirname, 'fixtures', 'gitoxide-managed-write-edit-owner-crash-child.js'), + fixturePath, + ], + { stdio: ['ignore', 'pipe', 'pipe'] }, + ); + const stderr: Buffer[] = []; + child.stderr.on('data', (chunk: Buffer) => stderr.push(chunk)); + const [exitCode] = (await once(child, 'exit')) as [number | null]; + assert.equal(exitCode, 73, Buffer.concat(stderr).toString('utf8')); + assert.equal( + gitBare(repositoryPath, ['rev-parse', 'refs/maka/accepted']), + baseline.head.commitOid, + ); + + const reopenedRootCapability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const reopenedRootOwner = await tryAcquireInteractiveRootOwner(reopenedRootCapability); + assert.ok(reopenedRootOwner); + t.after(() => reopenedRootOwner.close()); + const reopenedStores = await openInteractiveExecutionStoresForWrite(reopenedRootOwner.lease); + t.after(() => reopenedStores.sessionStore.close?.()); + const owner = createGitoxideManagedWriteEditOwnerInternal({ + storageRootLease: reopenedRootOwner.lease, + stores: reopenedStores, + invocationOwnerToken: helper.invocationOwnerToken, + helperCapability: helper.helperCapability, + repositoryPath, + workspaceId: ids.workspaceId, + workspaceEpochId: ids.workspaceEpochId, + }); + assert.equal(await owner.reconcileAcceptedProjection(), 'promoted'); + + const reopened = await owner.admitManagedMutation({ + operationId: 'operation-read-promoted-head', + toolName: 'Write', + persistedArgs: { path: 'notes.txt', content: 'after\n' }, + abortSignal: new AbortController().signal, + }); + assert.equal(reopened.durableDispatch.baseHeadRevision, baseline.head.revision + 1); + assert.equal( + gitBare(repositoryPath, ['rev-parse', 'refs/maka/accepted']), + reopened.durableDispatch.baseCommitOid, + ); + assert.deepEqual(reopened.immutableBase, { content: 'after\n' }); + assert.equal(await owner.reconcileAcceptedProjection(), 'already_current'); +}); + +async function admittedHelper(): Promise< + | { + readonly invocationOwnerToken: object; + readonly helperCapability: GitoxideHelperInvocationCapability; + readonly helperPath: string; + } + | undefined +> { + const configuredHelperPath = process.env.MAKA_GITOXIDE_HELPER_PATH; + if (!configuredHelperPath) return undefined; + const helperPath = await realpath(configuredHelperPath); + const helperBytes = await readFile(helperPath); + const helperInfo = await stat(helperPath); + const expectedSha256 = sha256Bytes(helperBytes); + const releaseOwnerToken = {}; + const invocationOwnerToken = {}; + const claim = issueGitoxideHelperReleaseArtifactClaimInternal(releaseOwnerToken, { + executablePath: helperPath, + expectedSha256, + expectedBytes: helperInfo.size, + platform: process.platform, + arch: process.arch, + protocolVersion: 1, + supportedOperations: [ + 'inspect_repository', + 'import_source_head', + 'create_candidate', + 'promote_candidate', + 'observe_accepted_ref', + 'read_tree_file', + ], + }); + const helperCapability = await admitGitoxideHelperArtifactInternal({ + releaseOwnerToken, + invocationOwnerToken, + claim, + }); + return { invocationOwnerToken, helperCapability, helperPath }; +} + +async function createRepository(t: TestContext): Promise { + const path = await realpath(await mkdtemp(join(tmpdir(), 'maka-gitoxide-write-edit-source-'))); + t.after(() => rm(path, { recursive: true, force: true })); + git(path, ['init', '--quiet', '--object-format=sha1']); + return path; +} + +function git(cwd: string, args: readonly string[]): string { + return execFileSync('git', ['-C', cwd, ...args], { encoding: 'utf8' }).trim(); +} + +function gitBare(repositoryPath: string, args: readonly string[]): string { + return execFileSync('git', [`--git-dir=${repositoryPath}`, ...args], { + encoding: 'utf8', + }).trim(); +} + +function sha256(value: string): `sha256:${string}` { + return `sha256:${createHash('sha256').update(value, 'utf8').digest('hex')}`; +} + +function sha256Bytes(value: Buffer): `sha256:${string}` { + return `sha256:${createHash('sha256').update(value).digest('hex')}`; +} diff --git a/packages/runtime-host/src/server/gitoxide-managed-write-edit-owner-internal.ts b/packages/runtime-host/src/server/gitoxide-managed-write-edit-owner-internal.ts new file mode 100644 index 0000000000..4f4afcd7df --- /dev/null +++ b/packages/runtime-host/src/server/gitoxide-managed-write-edit-owner-internal.ts @@ -0,0 +1,782 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import { createHash } from 'node:crypto'; +import { isDeepStrictEqual } from 'node:util'; +import { + isCanonicalManagedMutationPathV1, + MANAGED_MUTATION_EXECUTION_PROFILE_V1, + type RuntimeEvent, + type RuntimeEventManagedWorkspaceMutationV2, +} from '@maka/core/runtime-event'; +import { canonicalToolArgsHash } from '@maka/core/tool-args-identity'; +import type { + WorkspaceEpochRecordV1, + WorkspaceHeadRecordV1, + WorkspaceSuccessorAuthorityInput, + WorkspaceVersionRecordV1, +} from '@maka/core/workspace-version-authority'; +import type { + RuntimeManagedMutationAdmission, + RuntimeManagedMutationOperationProof, + RuntimeManagedMutationSettlement, + ToolRuntimeInput, +} from '@maka/runtime/tool-runtime'; +import { transformManagedMutation } from '@maka/runtime/managed-mutation-transform'; +import type { InteractiveExecutionStoresWriter } from '@maka/storage/execution-stores'; +import { + issueExecutionStoresWorkspaceMutationAuthorityInternal, + requireExecutionStoresWorkspaceMutationAuthorityInternal, + type ExecutionStoresWorkspaceMutationAuthorityInternal, +} from '@maka/storage/execution-stores-workspace-authority-internal'; +import type { StorageRootLease } from '@maka/storage/root-authority'; +import type { GitoxideHelperInvocationCapability } from './gitoxide-helper-artifact-authority-internal.js'; +import { GitoxideHelperInvocationError } from './gitoxide-helper-invocation-internal.js'; +import { + createGitoxideMutationCandidateAuthorityInternal, + type GitoxideMutationCandidateProofV1, +} from './gitoxide-mutation-candidate-receipt-authority-internal.js'; +import { + readGitoxideTreeFileInternal, + reopenGitoxideAcceptedRepositoryInternal, + type GitoxideAcceptedRepositoryCapability, +} from './gitoxide-repository-admission-authority-internal.js'; + +const ACCEPTED_REF = 'refs/maka/accepted'; +const SHA1_PATTERN = /^[0-9a-f]{40}$/u; + +export type GitoxideManagedWriteEditOwnerFailpoint = 'after_workspace_successor_commit'; + +export interface GitoxideManagedWriteEditOwnerInternal { + readonly admitManagedMutation: NonNullable; + reconcileAcceptedProjection(abortSignal?: AbortSignal): Promise<'already_current' | 'promoted'>; +} + +export interface GitoxideManagedWriteEditOwnerInputInternal { + readonly storageRootLease: StorageRootLease<'interactive', 'write'>; + readonly stores: InteractiveExecutionStoresWriter; + readonly invocationOwnerToken: object; + readonly helperCapability: GitoxideHelperInvocationCapability; + readonly repositoryPath: string; + readonly workspaceId: string; + readonly workspaceEpochId: string; + readonly failpoint?: (point: GitoxideManagedWriteEditOwnerFailpoint) => void | Promise; +} + +export function createGitoxideManagedWriteEditOwnerInternal( + input: GitoxideManagedWriteEditOwnerInputInternal, +): GitoxideManagedWriteEditOwnerInternal { + const ownerToken = {}; + const issuedSuccessors = new WeakMap(); + const capability = issueExecutionStoresWorkspaceMutationAuthorityInternal({ + ownerToken, + stores: input.stores, + verifyCandidate(candidateOutcome) { + const successor = issuedSuccessors.get(candidateOutcome); + if (!successor) throw new Error('Managed Write/Edit candidate proof is invalid'); + return successor; + }, + }); + // Resolve the capability now so an invalid execution-stores owner cannot be + // published as a usable Write/Edit admission. + const persistence = requireExecutionStoresWorkspaceMutationAuthorityInternal( + ownerToken, + capability, + ); + + const admitManagedMutation: NonNullable = async ( + request, + ) => { + if (request.toolName !== 'Write' && request.toolName !== 'Edit') { + throw new Error('Gitoxide managed mutation admits only Write and Edit'); + } + const path = requireCanonicalPath(request.persistedArgs); + const epoch = await persistence.readEpoch(input.workspaceId, input.workspaceEpochId); + if (!epochMatchesOwner(epoch, input.workspaceId, input.workspaceEpochId)) { + throw new Error('Gitoxide managed mutation durable workspace epoch is unavailable'); + } + const head = await persistence.readHead(input.workspaceId, input.workspaceEpochId); + if (!head) throw new Error('Gitoxide managed mutation has no accepted workspace head'); + const version = await persistence.readVersion(head.workspaceVersionId); + if (!version || !versionMatchesHead(version, head, epoch)) { + throw new Error('Gitoxide managed mutation workspace version is unavailable'); + } + request.abortSignal.throwIfAborted(); + + const acceptedRepositoryOwnerToken = {}; + const accepted = await reopenGitoxideAcceptedRepositoryInternal({ + invocationOwnerToken: input.invocationOwnerToken, + helperCapability: input.helperCapability, + acceptedRepositoryOwnerToken, + repositoryPath: input.repositoryPath, + acceptedRef: ACCEPTED_REF, + expectedAcceptedCommitOid: head.commitOid, + expectedAcceptedTreeOid: head.treeOid, + managedTreePolicyVersion: 3, + abortSignal: request.abortSignal, + }); + const baseContent = await readAcceptedFile({ + acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: accepted.acceptedRepositoryCapability, + path, + abortSignal: request.abortSignal, + }); + const candidateAuthority = await createGitoxideMutationCandidateAuthorityInternal({ + storageRootLease: input.storageRootLease, + baseHead: head, + acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: accepted.acceptedRepositoryCapability, + projectionOwnerToken: ownerToken, + }); + const durableDispatch = freezeManagedDispatch({ + epoch, + head, + expectedPath: path, + }); + + const admission: RuntimeManagedMutationAdmission = Object.freeze({ + durableDispatch, + immutableBase: Object.freeze({ content: baseContent }), + execute: async (operation: () => Promise) => + settleManagedMutation({ + request, + operation, + path, + head, + version, + epoch, + persistence, + candidateAuthority, + issuedSuccessors, + ownerToken, + failpoint: input.failpoint, + }), + dispose: async () => undefined, + }); + return admission; + }; + + const reconcileAcceptedProjection = async ( + abortSignal?: AbortSignal, + ): Promise<'already_current' | 'promoted'> => { + abortSignal?.throwIfAborted(); + const epoch = await persistence.readEpoch(input.workspaceId, input.workspaceEpochId); + if (!epochMatchesOwner(epoch, input.workspaceId, input.workspaceEpochId)) { + throw new Error('Gitoxide projection recovery durable workspace epoch is unavailable'); + } + const head = await persistence.readHead(input.workspaceId, input.workspaceEpochId); + if (!head) throw new Error('Gitoxide projection recovery has no accepted workspace head'); + const version = await persistence.readVersion(head.workspaceVersionId); + if (!version || !versionMatchesHead(version, head, epoch)) { + throw new Error('Gitoxide projection recovery workspace version is unavailable'); + } + try { + await reopenExactAcceptedRepository({ + input, + head, + ...(abortSignal ? { abortSignal } : {}), + }); + return 'already_current'; + } catch (currentProjectionError) { + if (version.protocol !== 'workspace_version_accepted_v1' || version.parents.length !== 1) { + throw new Error('Gitoxide accepted projection does not match its baseline head', { + cause: currentProjectionError, + }); + } + const parent = await persistence.readVersion(version.parents[0]); + if (!parent || !parentMatchesSuccessor(parent, version)) { + throw new Error('Gitoxide projection recovery parent version is unavailable', { + cause: currentProjectionError, + }); + } + const parentHead: WorkspaceHeadRecordV1 = Object.freeze({ + repositoryId: parent.repositoryId, + workspaceId: parent.workspaceId, + workspaceEpochId: parent.workspaceEpochId, + workspaceVersionId: parent.workspaceVersionId, + acceptedEventId: parent.acceptedEventId, + commitOid: parent.commitOid, + treeOid: parent.treeOid, + revision: version.baseHeadRevision, + }); + const acceptedRepositoryOwnerToken = {}; + const accepted = await reopenExactAcceptedRepository({ + input, + head: parentHead, + acceptedRepositoryOwnerToken, + ...(abortSignal ? { abortSignal } : {}), + }).catch((parentProjectionError) => { + throw new Error('Gitoxide accepted projection matches neither durable head nor parent', { + cause: new AggregateError([currentProjectionError, parentProjectionError]), + }); + }); + const evidence = await persistence.readMutationEvidence(version.origin.operationId); + if (!evidence) throw new Error('Gitoxide projection recovery operation evidence is missing'); + const recovered = await reconstructAcceptedSuccessor({ + epoch, + parentHead, + parentVersion: parent, + successorVersion: version, + evidence, + acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: accepted.acceptedRepositoryCapability, + ...(abortSignal ? { abortSignal } : {}), + }); + const candidateAuthority = await createGitoxideMutationCandidateAuthorityInternal({ + storageRootLease: input.storageRootLease, + baseHead: parentHead, + acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: accepted.acceptedRepositoryCapability, + projectionOwnerToken: ownerToken, + }); + const candidate = await candidateAuthority.capture({ + operationId: version.origin.operationId, + path: recovered.path, + content: recovered.content, + executionProfileDigest: MANAGED_MUTATION_EXECUTION_PROFILE_V1, + ...(abortSignal ? { abortSignal } : {}), + }); + assertCandidateProof(candidate, parentHead, recovered.path, recovered.content); + const successor = buildSuccessor({ + operationId: version.origin.operationId, + dispatchEventId: evidence.dispatchEvent.id, + outcome: evidence.outcomeEvent, + head: parentHead, + version: parent, + candidate, + }); + if (!successorMatchesVersion(successor, version)) { + throw new Error('Gitoxide projection recovery candidate conflicts with durable successor'); + } + issuedSuccessors.set(candidate, successor); + const replay = await persistence.commitSuccessor({ + candidateOutcome: candidate, + toolOutcome: toolOutcome(version.origin.operationId, evidence.outcomeEvent), + }); + if (replay.created || !headRecordsEqual(replay.head, head)) { + throw new Error('Gitoxide projection recovery did not exact-replay the durable successor'); + } + const promoted = await candidateAuthority.promote({ + proof: candidate, + projectionCapability: replay.projectionCapability, + nextAcceptedRepositoryOwnerToken: {}, + ...(abortSignal ? { abortSignal } : {}), + }); + if ( + promoted.acceptedCommitOid !== head.commitOid || + promoted.acceptedTreeOid !== head.treeOid + ) { + throw new Error('Recovered Gitoxide projection conflicts with the durable workspace head'); + } + return 'promoted'; + } + }; + + return Object.freeze({ admitManagedMutation, reconcileAcceptedProjection }); +} + +async function settleManagedMutation(input: { + readonly request: Parameters>[0]; + readonly operation: () => Promise; + readonly path: string; + readonly head: WorkspaceHeadRecordV1; + readonly version: WorkspaceVersionRecordV1; + readonly epoch: WorkspaceEpochRecordV1; + readonly persistence: ExecutionStoresWorkspaceMutationAuthorityInternal; + readonly candidateAuthority: Awaited< + ReturnType + >; + readonly issuedSuccessors: WeakMap; + readonly ownerToken: object; + readonly failpoint?: (point: GitoxideManagedWriteEditOwnerFailpoint) => void | Promise; +}): Promise { + const proof = await input.operation(); + const reservation = await input.persistence.readActiveMutation(input.epoch.workspaceInstanceId); + if (!reservationMatchesAdmission(reservation, input)) { + return unsettled('Managed Write/Edit operation has no exact durable reservation'); + } + + if (proof.isError) { + return commitTerminal(input, proof, 'operation_failed_no_effect'); + } + const mutation = proof.mutationResult; + if (!mutation || mutation.path !== input.path) { + return unsettled('Managed Write/Edit operation has no exact mutation result'); + } + if (!mutation.changed) { + return commitTerminal(input, proof, 'no_workspace_change'); + } + + let candidate: GitoxideMutationCandidateProofV1; + try { + candidate = await input.candidateAuthority.capture({ + operationId: input.request.operationId, + path: input.path, + content: mutation.content, + executionProfileDigest: MANAGED_MUTATION_EXECUTION_PROFILE_V1, + abortSignal: input.request.abortSignal, + }); + assertCandidateProof(candidate, input.head, input.path, mutation.content); + } catch (error) { + return Object.freeze({ kind: 'unsettled' as const, error }); + } + + const successor = buildSuccessor({ + operationId: input.request.operationId, + dispatchEventId: reservation!.dispatchEventId, + outcome: proof.durableOutcome, + head: input.head, + version: input.version, + candidate, + }); + input.issuedSuccessors.set(candidate, successor); + try { + const committed = await input.persistence.commitSuccessor({ + candidateOutcome: candidate, + toolOutcome: toolOutcome(input.request.operationId, proof.durableOutcome), + }); + await input.failpoint?.('after_workspace_successor_commit'); + const promoted = await input.candidateAuthority.promote({ + proof: candidate, + projectionCapability: committed.projectionCapability, + nextAcceptedRepositoryOwnerToken: {}, + abortSignal: input.request.abortSignal, + }); + if ( + promoted.acceptedCommitOid !== committed.head.commitOid || + promoted.acceptedTreeOid !== committed.head.treeOid + ) { + return unsettled('Accepted Gitoxide projection conflicts with the durable workspace head'); + } + return Object.freeze({ + kind: 'workspace_successor_committed' as const, + durableOutcome: proof.durableOutcome, + }); + } catch (error) { + return Object.freeze({ kind: 'unsettled' as const, error }); + } +} + +async function commitTerminal( + input: Parameters[0], + proof: RuntimeManagedMutationOperationProof, + terminalKind: 'no_workspace_change' | 'operation_failed_no_effect', +): Promise { + if (proof.terminalOutcome?.kind !== terminalKind) { + return unsettled('Managed Write/Edit terminal proof is unavailable'); + } + try { + await input.persistence.commitTerminal({ + toolOutcome: toolOutcome(input.request.operationId, proof.terminalOutcome.durableOutcome), + }); + return Object.freeze({ + kind: + terminalKind === 'no_workspace_change' + ? ('no_workspace_change_committed' as const) + : ('operation_failed_no_effect_committed' as const), + durableOutcome: proof.terminalOutcome.durableOutcome, + }); + } catch (error) { + return Object.freeze({ kind: 'unsettled' as const, error }); + } +} + +async function readAcceptedFile(input: { + readonly acceptedRepositoryOwnerToken: object; + readonly acceptedRepositoryCapability: Parameters< + typeof readGitoxideTreeFileInternal + >[0]['acceptedRepositoryCapability']; + readonly path: string; + readonly abortSignal: AbortSignal; +}): Promise { + try { + return ( + await readGitoxideTreeFileInternal({ + acceptedRepositoryOwnerToken: input.acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: input.acceptedRepositoryCapability, + path: input.path, + abortSignal: input.abortSignal, + }) + ).content; + } catch (error) { + if ( + error instanceof GitoxideHelperInvocationError && + error.code === 'gitoxide_helper_operation_failed' && + error.helperReason === 'tree_file_unavailable' + ) { + return null; + } + throw error; + } +} + +function freezeManagedDispatch(input: { + readonly epoch: WorkspaceEpochRecordV1; + readonly head: WorkspaceHeadRecordV1; + readonly expectedPath: string; +}): Readonly { + return Object.freeze({ + protocol: 'managed_mutation_v2' as const, + repositoryId: input.head.repositoryId, + workspaceId: input.head.workspaceId, + workspaceEpochId: input.head.workspaceEpochId, + workspaceInstanceId: input.epoch.workspaceInstanceId, + objectFormat: 'sha1' as const, + baseWorkspaceVersionId: input.head.workspaceVersionId, + baseAcceptedEventId: input.head.acceptedEventId, + baseHeadRevision: input.head.revision, + baseCommitOid: input.head.commitOid, + baseTreeOid: input.head.treeOid, + expectedPath: input.expectedPath, + pathPolicyVersion: 3 as const, + executionProfileDigest: MANAGED_MUTATION_EXECUTION_PROFILE_V1, + }); +} + +function epochMatchesOwner( + epoch: WorkspaceEpochRecordV1 | undefined, + workspaceId: string, + workspaceEpochId: string, +): epoch is WorkspaceEpochRecordV1 { + return Boolean( + epoch && + epoch.workspaceId === workspaceId && + epoch.workspaceEpochId === workspaceEpochId && + epoch.mode === 'managed_worktree' && + epoch.objectFormat === 'sha1' && + /^instance_[0-9a-f]{32}$/u.test(epoch.workspaceInstanceId), + ); +} + +function versionMatchesHead( + version: WorkspaceVersionRecordV1, + head: WorkspaceHeadRecordV1, + epoch: WorkspaceEpochRecordV1, +): boolean { + return ( + version.repositoryId === epoch.repositoryId && + version.repositoryId === head.repositoryId && + version.workspaceId === head.workspaceId && + version.workspaceEpochId === head.workspaceEpochId && + version.workspaceVersionId === head.workspaceVersionId && + version.acceptedEventId === head.acceptedEventId && + version.commitOid === head.commitOid && + version.treeOid === head.treeOid && + version.objectFormat === 'sha1' && + head.revision >= 1 + ); +} + +function reservationMatchesAdmission( + reservation: Awaited< + ReturnType + >, + input: Pick[0], 'request' | 'path' | 'head' | 'epoch'>, +): boolean { + return Boolean( + reservation && + reservation.workspaceInstanceId === input.epoch.workspaceInstanceId && + reservation.repositoryId === input.head.repositoryId && + reservation.workspaceId === input.head.workspaceId && + reservation.workspaceEpochId === input.head.workspaceEpochId && + reservation.operationId === input.request.operationId && + reservation.baseWorkspaceVersionId === input.head.workspaceVersionId && + reservation.baseAcceptedEventId === input.head.acceptedEventId && + reservation.baseHeadRevision === input.head.revision && + reservation.baseCommitOid === input.head.commitOid && + reservation.baseTreeOid === input.head.treeOid && + reservation.expectedPath === input.path && + reservation.executionProfileDigest === MANAGED_MUTATION_EXECUTION_PROFILE_V1, + ); +} + +function assertCandidateProof( + proof: GitoxideMutationCandidateProofV1, + head: WorkspaceHeadRecordV1, + path: string, + content: string, +): void { + const receipt = proof.receipt; + if ( + receipt.disposition !== 'published' || + receipt.repositoryId !== head.repositoryId || + receipt.workspaceId !== head.workspaceId || + receipt.workspaceEpochId !== head.workspaceEpochId || + receipt.workspaceVersionId !== head.workspaceVersionId || + receipt.baseAcceptedEventId !== head.acceptedEventId || + receipt.baseHeadRevision !== head.revision || + receipt.baseCommitOid !== head.commitOid || + receipt.baseTreeOid !== head.treeOid || + receipt.path !== path || + receipt.contentSha256 !== sha256(content) || + receipt.executionProfileDigest !== MANAGED_MUTATION_EXECUTION_PROFILE_V1 || + !SHA1_PATTERN.test(receipt.candidateCommitOid) || + !SHA1_PATTERN.test(receipt.candidateTreeOid) || + !SHA1_PATTERN.test(receipt.resultBlobOid) + ) { + throw new Error('Gitoxide candidate proof conflicts with the admitted Write/Edit operation'); + } +} + +function buildSuccessor(input: { + readonly operationId: string; + readonly dispatchEventId: string; + readonly outcome: RuntimeEvent; + readonly head: WorkspaceHeadRecordV1; + readonly version: WorkspaceVersionRecordV1; + readonly candidate: GitoxideMutationCandidateProofV1; +}): WorkspaceSuccessorAuthorityInput { + const receipt = input.candidate.receipt; + const identity = digest('accepted-successor', input.operationId, receipt.candidateCommitOid); + return Object.freeze({ + acceptedEventId: `workspace-successor-${identity}`, + committedAt: input.outcome.ts, + successor: Object.freeze({ + repositoryId: input.head.repositoryId, + workspaceId: input.head.workspaceId, + workspaceEpochId: input.head.workspaceEpochId, + workspaceVersionId: `version_${identity}`, + objectFormat: 'sha1' as const, + parentWorkspaceVersionId: input.head.workspaceVersionId, + baseAcceptedEventId: input.head.acceptedEventId, + baseHeadRevision: input.head.revision, + commitOid: receipt.candidateCommitOid, + treeOid: receipt.candidateTreeOid, + policyHash: input.version.policyHash, + treeDeltaDigest: sha256( + `gitoxide-tree-delta-v1\0${input.head.treeOid}\0${receipt.candidateTreeOid}\0${receipt.path}\0${receipt.resultBlobOid}`, + ), + changedPaths: Object.freeze([receipt.path]), + changedFileCount: 1, + deletedFileCount: 0, + executionProfileDigest: MANAGED_MUTATION_EXECUTION_PROFILE_V1, + }), + origin: Object.freeze({ + operationId: input.operationId, + dispatchEventId: input.dispatchEventId, + outcomeEventId: input.outcome.id, + }), + }); +} + +function toolOutcome(operationId: string, runtimeEvent: RuntimeEvent) { + return Object.freeze({ + operationId, + journalEventId: `${operationId}_outcome`, + runtimeEvent, + committedAt: runtimeEvent.ts, + }); +} + +function unsettled(message: string): RuntimeManagedMutationSettlement { + return Object.freeze({ kind: 'unsettled' as const, error: new Error(message) }); +} + +function digest(domain: string, ...values: readonly string[]): string { + const hash = createHash('sha256').update(`maka-${domain}-v1\0`, 'utf8'); + for (const value of values) hash.update(value).update('\0'); + return hash.digest('hex').slice(0, 32); +} + +function sha256(value: string): `sha256:${string}` { + return `sha256:${createHash('sha256').update(value, 'utf8').digest('hex')}`; +} + +function requireCanonicalPath(args: unknown): string { + if (!args || typeof args !== 'object' || Array.isArray(args)) { + throw new Error('Gitoxide managed mutation arguments are invalid'); + } + const path = (args as Record).path; + if (!isCanonicalManagedMutationPathV1(path)) { + throw new Error('Gitoxide managed mutation path must already be canonical'); + } + return path; +} + +async function reopenExactAcceptedRepository(input: { + readonly input: GitoxideManagedWriteEditOwnerInputInternal; + readonly head: WorkspaceHeadRecordV1; + readonly acceptedRepositoryOwnerToken?: object; + readonly abortSignal?: AbortSignal; +}) { + return reopenGitoxideAcceptedRepositoryInternal({ + invocationOwnerToken: input.input.invocationOwnerToken, + helperCapability: input.input.helperCapability, + acceptedRepositoryOwnerToken: input.acceptedRepositoryOwnerToken ?? {}, + repositoryPath: input.input.repositoryPath, + acceptedRef: ACCEPTED_REF, + expectedAcceptedCommitOid: input.head.commitOid, + expectedAcceptedTreeOid: input.head.treeOid, + managedTreePolicyVersion: 3, + ...(input.abortSignal ? { abortSignal: input.abortSignal } : {}), + }); +} + +async function reconstructAcceptedSuccessor(input: { + readonly epoch: WorkspaceEpochRecordV1; + readonly parentHead: WorkspaceHeadRecordV1; + readonly parentVersion: WorkspaceVersionRecordV1; + readonly successorVersion: Extract< + WorkspaceVersionRecordV1, + { protocol: 'workspace_version_accepted_v1' } + >; + readonly evidence: NonNullable< + Awaited> + >; + readonly acceptedRepositoryOwnerToken: object; + readonly acceptedRepositoryCapability: GitoxideAcceptedRepositoryCapability; + readonly abortSignal?: AbortSignal; +}): Promise<{ readonly path: string; readonly content: string }> { + const call = input.evidence.callEvent; + const dispatch = input.evidence.dispatchEvent; + const outcome = input.evidence.outcomeEvent; + const callContent = call.content; + const dispatchFact = dispatch.actions?.toolDispatch; + const managed = dispatchFact?.managedMutation; + const outcomeContent = outcome.content; + if ( + callContent?.kind !== 'function_call' || + (callContent.name !== 'Write' && callContent.name !== 'Edit') || + call.refs?.operationId !== input.successorVersion.origin.operationId || + dispatch.id !== input.successorVersion.origin.dispatchEventId || + dispatch.refs?.operationId !== input.successorVersion.origin.operationId || + dispatchFact?.operationId !== input.successorVersion.origin.operationId || + dispatchFact.providerToolCallId !== callContent.id || + dispatchFact.toolName !== callContent.name || + dispatchFact.recoveryMode !== 'reconcile' || + dispatchFact.canonicalArgsHash !== canonicalToolArgsHash(callContent.name, callContent.args) || + outcome.id !== input.successorVersion.origin.outcomeEventId || + outcome.refs?.operationId !== input.successorVersion.origin.operationId || + outcomeContent?.kind !== 'function_response' || + outcomeContent.id !== callContent.id || + outcomeContent.name !== callContent.name || + outcomeContent.isError === true || + !managedMutationMatchesParent(managed, input.epoch, input.parentHead) + ) { + throw new Error('Gitoxide projection recovery operation evidence is invalid'); + } + const path = requireCanonicalPath(callContent.args); + if ( + path !== managed.expectedPath || + input.successorVersion.changedPaths.length !== 1 || + input.successorVersion.changedPaths[0] !== path || + input.successorVersion.executionProfileDigest !== MANAGED_MUTATION_EXECUTION_PROFILE_V1 + ) { + throw new Error('Gitoxide projection recovery path authority is invalid'); + } + const baseContent = await readAcceptedFile({ + acceptedRepositoryOwnerToken: input.acceptedRepositoryOwnerToken, + acceptedRepositoryCapability: input.acceptedRepositoryCapability, + path, + abortSignal: input.abortSignal ?? new AbortController().signal, + }); + const transformed = transformManagedMutation({ + toolName: callContent.name, + canonicalPath: path, + baseContent, + args: structuredClone(callContent.args), + }); + if (!transformed.changed) { + throw new Error('Gitoxide projection recovery successor has no content change'); + } + return Object.freeze({ path, content: transformed.content }); +} + +function managedMutationMatchesParent( + managed: RuntimeEventManagedWorkspaceMutationV2 | undefined, + epoch: WorkspaceEpochRecordV1, + parent: WorkspaceHeadRecordV1, +): managed is RuntimeEventManagedWorkspaceMutationV2 { + return Boolean( + managed && + managed.protocol === 'managed_mutation_v2' && + managed.repositoryId === parent.repositoryId && + managed.workspaceId === parent.workspaceId && + managed.workspaceEpochId === parent.workspaceEpochId && + managed.workspaceInstanceId === epoch.workspaceInstanceId && + managed.objectFormat === 'sha1' && + managed.baseWorkspaceVersionId === parent.workspaceVersionId && + managed.baseAcceptedEventId === parent.acceptedEventId && + managed.baseHeadRevision === parent.revision && + managed.baseCommitOid === parent.commitOid && + managed.baseTreeOid === parent.treeOid && + managed.pathPolicyVersion === 3 && + managed.executionProfileDigest === MANAGED_MUTATION_EXECUTION_PROFILE_V1, + ); +} + +function parentMatchesSuccessor( + parent: WorkspaceVersionRecordV1, + successor: Extract, +): boolean { + return ( + parent.repositoryId === successor.repositoryId && + parent.workspaceId === successor.workspaceId && + parent.workspaceEpochId === successor.workspaceEpochId && + parent.workspaceVersionId === successor.parents[0] && + parent.acceptedEventId === successor.baseAcceptedEventId && + parent.commitOid !== successor.commitOid && + successor.baseHeadRevision >= 1 + ); +} + +function successorMatchesVersion( + expected: WorkspaceSuccessorAuthorityInput, + actual: Extract, +): boolean { + const successor = expected.successor; + return ( + expected.acceptedEventId === actual.acceptedEventId && + expected.committedAt === actual.committedAt && + isDeepStrictEqual(expected.origin, { + operationId: actual.origin.operationId, + dispatchEventId: actual.origin.dispatchEventId, + outcomeEventId: actual.origin.outcomeEventId, + }) && + successor.repositoryId === actual.repositoryId && + successor.workspaceId === actual.workspaceId && + successor.workspaceEpochId === actual.workspaceEpochId && + successor.workspaceVersionId === actual.workspaceVersionId && + successor.objectFormat === actual.objectFormat && + successor.parentWorkspaceVersionId === actual.parents[0] && + successor.baseAcceptedEventId === actual.baseAcceptedEventId && + successor.baseHeadRevision === actual.baseHeadRevision && + successor.commitOid === actual.commitOid && + successor.treeOid === actual.treeOid && + successor.policyHash === actual.policyHash && + successor.treeDeltaDigest === actual.treeDeltaDigest && + isDeepStrictEqual(successor.changedPaths, actual.changedPaths) && + successor.changedFileCount === actual.changedFileCount && + successor.deletedFileCount === actual.deletedFileCount && + successor.executionProfileDigest === actual.executionProfileDigest + ); +} + +function headRecordsEqual(left: WorkspaceHeadRecordV1, right: WorkspaceHeadRecordV1): boolean { + return ( + left.repositoryId === right.repositoryId && + left.workspaceId === right.workspaceId && + left.workspaceEpochId === right.workspaceEpochId && + left.workspaceVersionId === right.workspaceVersionId && + left.acceptedEventId === right.acceptedEventId && + left.commitOid === right.commitOid && + left.treeOid === right.treeOid && + left.revision === right.revision + ); +} diff --git a/packages/runtime/package.json b/packages/runtime/package.json index 53d05f563c..ca99d97abf 100644 --- a/packages/runtime/package.json +++ b/packages/runtime/package.json @@ -113,6 +113,7 @@ "./tool-result-archive-capability": "./dist/tool-result-archive-capability.js", "./tool-result-archive-resource": "./dist/tool-result-archive-resource.js", "./tool-runtime": "./dist/tool-runtime.js", + "./managed-mutation-transform": "./dist/managed-mutation-transform.js", "./web-fetch-tool": "./dist/web-fetch-tool.js", "./web-search-tool": "./dist/web-search-tool.js", "./xai-oauth-enrollment": "./dist/xai-oauth-enrollment.js" diff --git a/packages/runtime/src/__tests__/tool-runtime-durable-boundary.test.ts b/packages/runtime/src/__tests__/tool-runtime-durable-boundary.test.ts index 19fcd830e5..2a44d5fc76 100644 --- a/packages/runtime/src/__tests__/tool-runtime-durable-boundary.test.ts +++ b/packages/runtime/src/__tests__/tool-runtime-durable-boundary.test.ts @@ -33,6 +33,7 @@ import { ToolRuntime, type MakaTool, type RuntimeManagedMutationAdmission, + type RuntimeManagedMutationSettlement, type ToolRuntimeInput, } from '../tool-runtime.js'; @@ -335,7 +336,6 @@ describe('ToolRuntime durable boundary', () => { }); return { kind: 'no_workspace_change_committed', - providerResult: proof.content, durableOutcome: proof.terminalOutcome!.durableOutcome, }; }, @@ -733,23 +733,17 @@ describe('ToolRuntime durable boundary', () => { admitManagedMutation: async (input) => { operationId = input.operationId; return managedAdmission(async (operation) => { - await operation(); - const result = { error: 'candidate was safely discarded' }; + const proof = await operation(); + assert.equal(proof.terminalOutcome?.kind, 'operation_failed_no_effect'); return { kind: 'operation_failed_no_effect_committed', - providerResult: result, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: result }, - true, - { terminalKind: 'operation_failed_no_effect' }, - ), + durableOutcome: proof.terminalOutcome!.durableOutcome, }; }); }, }, ); - const managedTool = tool(() => ({ ok: true })); + const managedTool = tool(() => ({ error: 'candidate was safely discarded' })); managedTool.name = 'Write'; managedTool.recoveryMode = 'reconcile'; managedTool.durableExecutionProfile = 'managed_mutation_v1'; @@ -777,20 +771,17 @@ describe('ToolRuntime durable boundary', () => { { admitManagedMutation: async (input) => { operationId = input.operationId; - return managedAdmission(async (operation) => { - await operation(); - const result = { ok: true, changed: false }; - return { - kind: 'no_workspace_change_committed', - providerResult: result, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: result }, - false, - { terminalKind: 'no_workspace_change' }, - ), - }; - }); + return { + ...managedAdmission(async (operation) => { + const proof = await operation(); + assert.equal(proof.terminalOutcome?.kind, 'no_workspace_change'); + return { + kind: 'no_workspace_change_committed', + durableOutcome: proof.terminalOutcome!.durableOutcome, + }; + }), + immutableBase: Object.freeze({ content: 'same' }), + }; }, }, ); @@ -799,13 +790,16 @@ describe('ToolRuntime durable boundary', () => { managedTool.recoveryMode = 'reconcile'; managedTool.durableExecutionProfile = 'managed_mutation_v1'; - assert.deepEqual(await harness.execute(managedTool), { ok: true, changed: false }); + assert.deepEqual( + await harness.executeWithInput(managedTool, { path: 'notes.txt', content: 'same' }), + { kind: 'file_write', path: 'notes.txt', bytes: 4 }, + ); const published = harness.events.at(-1); assert.equal(published?.type, 'tool_result'); assert.equal(published?.type === 'tool_result' && published.isError, false); }); - it('snapshots a safe-discard result before its owner can mutate it', async () => { + it('ignores a mutable provider result smuggled across the owner boundary', async () => { let operationId = ''; const ownerResult = { error: 'discarded-A' }; const appendedMessages: StoredMessage[] = []; @@ -826,22 +820,18 @@ describe('ToolRuntime durable boundary', () => { admitManagedMutation: async (input) => { operationId = input.operationId; return managedAdmission(async (operation) => { - await operation(); + const proof = await operation(); + assert.equal(proof.terminalOutcome?.kind, 'operation_failed_no_effect'); return { kind: 'operation_failed_no_effect_committed', providerResult: ownerResult, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: { error: 'discarded-A' } }, - true, - { terminalKind: 'operation_failed_no_effect' }, - ), - }; + durableOutcome: proof.terminalOutcome!.durableOutcome, + } as unknown as RuntimeManagedMutationSettlement; }); }, }, ); - const managedTool = tool(() => ({ ok: true })); + const managedTool = tool(() => ({ error: 'runtime-owned-A' })); managedTool.name = 'Write'; managedTool.recoveryMode = 'reconcile'; managedTool.durableExecutionProfile = 'managed_mutation_v1'; @@ -850,11 +840,11 @@ describe('ToolRuntime durable boundary', () => { const storedResult = appendedMessages.find((message) => message.type === 'tool_result'); assert.equal(ownerResult.error, 'mutated-B'); - assert.deepEqual(result, { error: 'discarded-A' }); + assert.deepEqual(result, { error: 'runtime-owned-A' }); assert.equal(Object.isFrozen(result), true); assert.deepEqual(storedResult?.type === 'tool_result' ? storedResult.content : undefined, { kind: 'json', - value: { error: 'discarded-A' }, + value: { error: 'runtime-owned-A' }, }); }); @@ -876,16 +866,11 @@ describe('ToolRuntime durable boundary', () => { operationId = input.operationId; return managedAdmission(async (operation) => { retainedOperation = operation; - const result = { error: 'candidate was safely discarded' }; + const proof = await operation(); + assert.equal(proof.terminalOutcome?.kind, 'operation_failed_no_effect'); return { kind: 'operation_failed_no_effect_committed', - providerResult: result, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: result }, - true, - { terminalKind: 'operation_failed_no_effect' }, - ), + durableOutcome: proof.terminalOutcome!.durableOutcome, }; }); }, @@ -893,7 +878,7 @@ describe('ToolRuntime durable boundary', () => { ); const managedTool = tool(() => { implementationCalls += 1; - return { ok: true }; + return { error: 'candidate was safely discarded' }; }); managedTool.name = 'Write'; managedTool.recoveryMode = 'reconcile'; @@ -904,7 +889,7 @@ describe('ToolRuntime durable boundary', () => { }); assert.ok(retainedOperation); await assert.rejects(retainedOperation(), /operation capability is closed/i); - assert.equal(implementationCalls, 0); + assert.equal(implementationCalls, 1); }); it('does not accept terminal settlement while a detached operation is running', async () => { @@ -930,7 +915,6 @@ describe('ToolRuntime durable boundary', () => { const result = { error: 'candidate was safely discarded' }; return { kind: 'operation_failed_no_effect_committed', - providerResult: result, durableOutcome: managedOutcomeEvent( operationId, { kind: 'json', value: result }, @@ -967,7 +951,7 @@ describe('ToolRuntime durable boundary', () => { ); }); - it('rejects a safe discard whose live error differs from its durable result', async () => { + it('rejects a terminal proof whose durable result differs from the Runtime result', async () => { let operationId = ''; const harness = makeHarness( { @@ -985,7 +969,6 @@ describe('ToolRuntime durable boundary', () => { await operation(); return { kind: 'operation_failed_no_effect_committed', - providerResult: { error: 'live provider error A' }, durableOutcome: managedOutcomeEvent( operationId, { kind: 'json', value: { error: 'durable replay error B' } }, @@ -1008,108 +991,6 @@ describe('ToolRuntime durable boundary', () => { ); }); - it('fail-stops safe-discard canonicalization without writing generic T2', async () => { - let genericOutcomeCalls = 0; - let operationId = ''; - const providerResult = Object.defineProperty({}, 'kind', { - enumerable: true, - get: () => { - throw new Error('provider result getter exploded'); - }, - }); - const harness = makeHarness( - { - commitToolPrepared: async () => ({ created: true, runtimeEventSeq: 1 }), - commitToolOutcome: async () => { - genericOutcomeCalls += 1; - return { created: true, runtimeEventSeq: 2 }; - }, - }, - undefined, - 'run-1', - { - admitManagedMutation: async (input) => { - operationId = input.operationId; - return managedAdmission(async (operation) => { - await operation(); - return { - kind: 'operation_failed_no_effect_committed', - providerResult, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: { error: 'discarded' } }, - true, - ), - }; - }); - }, - }, - ); - const managedTool = tool(() => ({ ok: true })); - managedTool.name = 'Write'; - managedTool.recoveryMode = 'reconcile'; - managedTool.durableExecutionProfile = 'managed_mutation_v1'; - - await assert.rejects( - harness.execute(managedTool), - /strict JSON.*accessor|provider result getter exploded|byte limit exceeded/i, - ); - assert.equal(genericOutcomeCalls, 0); - assert.equal( - harness.events.some((event) => event.type === 'tool_result'), - false, - ); - }); - - it('fail-stops an oversized safe discard before durable publication', async () => { - let genericOutcomeCalls = 0; - let operationId = ''; - const oversized = { error: 'x'.repeat(128) }; - const harness = makeHarness( - { - commitToolPrepared: async () => ({ created: true, runtimeEventSeq: 1 }), - commitToolOutcome: async () => { - genericOutcomeCalls += 1; - return { created: true, runtimeEventSeq: 2 }; - }, - }, - undefined, - 'run-1', - { - admitManagedMutation: async (input) => { - operationId = input.operationId; - return managedAdmission(async (operation) => { - await operation(); - return { - kind: 'operation_failed_no_effect_committed', - providerResult: oversized, - durableOutcome: managedOutcomeEvent( - operationId, - { kind: 'json', value: oversized }, - true, - { - origin: 'code_mode', - modelVisibility: 'hidden', - toolCallId: 'nested-call-1', - parentToolCallId: 'exec-1', - parentOperationId: 'exec-op-1', - }, - ), - }; - }); - }, - }, - ); - const managedTool = tool(() => ({ ok: true })); - managedTool.name = 'Write'; - managedTool.recoveryMode = 'reconcile'; - managedTool.durableExecutionProfile = 'managed_mutation_v1'; - - await assert.rejects(harness.executeNested(managedTool, 32), /byte limit exceeded/i); - assert.equal(genericOutcomeCalls, 0); - assert.equal(JSON.stringify(harness.events).includes(oversized.error), false); - }); - it('stops snapshot traversal as soon as a managed result exceeds its byte budget', async () => { let genericOutcomeCalls = 0; let lateGetterReads = 0; diff --git a/packages/runtime/src/tool-runtime.ts b/packages/runtime/src/tool-runtime.ts index 2c69b41419..ca7cb30462 100644 --- a/packages/runtime/src/tool-runtime.ts +++ b/packages/runtime/src/tool-runtime.ts @@ -535,8 +535,6 @@ export type RuntimeManagedMutationSettlement = } | { readonly kind: 'no_workspace_change_committed' | 'operation_failed_no_effect_committed'; - /** Exact value returned to the provider and canonicalized for durable replay. */ - readonly providerResult: unknown; readonly durableOutcome: RuntimeEvent; } | { readonly kind: 'unsettled'; readonly error: unknown }; @@ -1830,7 +1828,7 @@ export class ToolRuntime { // discarded; every other failure remains unsettled for recovery. throw new RuntimeManagedMutationUnsettledError(ownerError); } - const normalized = normalizeManagedMutationSettlement(settlement, ctx.maxResultBytes); + const normalized = normalizeManagedMutationSettlement(settlement); if (normalized.kind === 'workspace_successor_committed') { if (!runtimeOwnedValue) { throw new RuntimeManagedMutationUnsettledError( @@ -1845,9 +1843,14 @@ export class ToolRuntime { durableOutcome: normalized.durableOutcome, }; } else { + if (!runtimeOwnedValue) { + throw new RuntimeManagedMutationUnsettledError( + new Error('Managed mutation owner committed a terminal state without execution'), + ); + } settledExecution = { kind: 'managed', - value: normalized.value, + value: runtimeOwnedValue, durableOutcome: normalized.durableOutcome, terminalKind: normalized.kind === 'no_workspace_change_committed' @@ -3516,17 +3519,13 @@ function uncertainOutcomeSignalFromError(error: unknown): ToolUncertainOutcomeSi }; } -function normalizeManagedMutationSettlement( - settlement: unknown, - maxResultBytes: number | undefined, -): +function normalizeManagedMutationSettlement(settlement: unknown): | { kind: 'workspace_successor_committed'; durableOutcome: RuntimeEvent; } | { kind: 'no_workspace_change_committed' | 'operation_failed_no_effect_committed'; - value: RuntimeManagedMutationOperationValue; durableOutcome: RuntimeEvent; } { if (!settlement || typeof settlement !== 'object' || Array.isArray(settlement)) { @@ -3563,10 +3562,6 @@ function normalizeManagedMutationSettlement( }; } - if (!Object.hasOwn(record, 'providerResult')) { - throw new Error('Managed no-effect settlement has no provider result'); - } - const providerResult = snapshotManagedToolResult(record.providerResult, maxResultBytes); const response = durableOutcome.content; const expectedError = kind === 'operation_failed_no_effect_committed'; if ( @@ -3575,21 +3570,8 @@ function normalizeManagedMutationSettlement( ) { throw new Error('Managed no-effect settlement has the wrong durable outcome state'); } - const content = Object.freeze(coerceResultContent(providerResult)); - const outcome = Object.freeze({ - content, - isError: expectedError, - durationMs: - typeof durableOutcome.actions?.stateDelta?.durationMs === 'number' - ? durableOutcome.actions.stateDelta.durationMs - : 0, - }); return { kind, - value: Object.freeze({ - result: providerResult, - outcome, - }), durableOutcome, }; } diff --git a/packages/storage/package.json b/packages/storage/package.json index 15ee369c7c..2a48686d2d 100644 --- a/packages/storage/package.json +++ b/packages/storage/package.json @@ -17,6 +17,7 @@ "./encrypted-file-managed-secret-store": "./dist/encrypted-file-managed-secret-store.js", "./execution-stores": "./dist/execution-stores.js", "./execution-stores-workspace-authority-internal": "./dist/execution-stores-workspace-authority-internal.js", + "./test-only/execution-stores-workspace-authority": "./dist/test-only/execution-stores-workspace-authority.js", "./external-sessions": "./dist/external-sessions.js", "./file-update-lock": "./dist/file-update-lock.js", "./foreign-session-store": "./dist/foreign-session-store.js", diff --git a/packages/storage/src/__tests__/execution-stores-workspace-authority-internal.test.ts b/packages/storage/src/__tests__/execution-stores-workspace-authority-internal.test.ts index 63857d688f..de663b2b41 100644 --- a/packages/storage/src/__tests__/execution-stores-workspace-authority-internal.test.ts +++ b/packages/storage/src/__tests__/execution-stores-workspace-authority-internal.test.ts @@ -29,7 +29,7 @@ import { } from '../execution-stores-workspace-authority-internal.js'; import { resolveStorageRoot, tryAcquireInteractiveRootOwner } from '../root-authority.js'; -test('binds workspace mutation persistence to one execution-stores owner capability', async () => { +test('binds each workspace mutation verifier to its execution-stores owner capability', async () => { const root = await mkdtemp(join(tmpdir(), 'maka-execution-workspace-authority-')); const capability = await resolveStorageRoot({ path: root, kind: 'interactive' }); const rootOwner = await tryAcquireInteractiveRootOwner(capability); @@ -45,6 +45,20 @@ test('binds workspace mutation persistence to one execution-stores owner capabil throw new Error('not used'); }, }); + const secondOwnerToken = {}; + const secondAuthorityCapability = issueExecutionStoresWorkspaceMutationAuthorityInternal({ + ownerToken: secondOwnerToken, + stores, + verifyCandidate: () => { + throw new Error('not used'); + }, + }); + assert.ok( + requireExecutionStoresWorkspaceMutationAuthorityInternal( + secondOwnerToken, + secondAuthorityCapability, + ), + ); assert.throws( () => requireExecutionStoresWorkspaceMutationAuthorityInternal({}, authorityCapability), /capability is invalid/i, @@ -53,6 +67,14 @@ test('binds workspace mutation persistence to one execution-stores owner capabil ownerToken, authorityCapability, ); + assert.equal( + await authority.readEpoch( + 'workspace_'.concat('1'.repeat(32)), + 'epoch_'.concat('2'.repeat(32)), + ), + undefined, + ); + assert.equal(await authority.readMutationEvidence('operation-missing'), undefined); assert.equal( await authority.readHead( 'workspace_'.concat('1'.repeat(32)), diff --git a/packages/storage/src/execution-stores-workspace-authority-internal.ts b/packages/storage/src/execution-stores-workspace-authority-internal.ts index 47e97c895b..6a71dbb6c9 100644 --- a/packages/storage/src/execution-stores-workspace-authority-internal.ts +++ b/packages/storage/src/execution-stores-workspace-authority-internal.ts @@ -18,19 +18,25 @@ */ import type { + WorkspaceBaselineAuthorityInput, + WorkspaceBaselineCommitResult, WorkspaceHeadRecordV1, + WorkspaceEpochRecordV1, WorkspaceSuccessorAuthorityInput, WorkspaceVersionRecordV1, } from '@maka/core/workspace-version-authority'; import { adoptWorkspaceBaselineAuthorityStoreRootInternal, + commitWorkspaceBaselineInternal, commitManagedMutationTerminalInternal, - commitWorkspaceSuccessorInternal, + commitVerifiedWorkspaceSuccessorInternal, readActiveManagedMutationInternal, + readManagedMutationEvidenceInternal, + readWorkspaceEpochInternal, readWorkspaceHeadInternal, readWorkspaceVersionInternal, - registerWorkspaceSuccessorCandidateVerifierInternal, type ManagedMutationReservationRecordV1, + type ManagedMutationEvidenceRecordV1, type ManagedMutationTerminalCommitInput, type ManagedMutationTerminalCommitResult, type WorkspaceSuccessorCommitInput, @@ -51,6 +57,10 @@ export interface ExecutionStoresWorkspaceSuccessorCommitResult } export interface ExecutionStoresWorkspaceMutationAuthorityInternal { + readEpoch( + workspaceId: string, + workspaceEpochId: string, + ): Promise; readHead( workspaceId: string, workspaceEpochId: string, @@ -59,6 +69,7 @@ export interface ExecutionStoresWorkspaceMutationAuthorityInternal { readActiveMutation( workspaceInstanceId: string, ): Promise; + readMutationEvidence(operationId: string): Promise; commitSuccessor( input: WorkspaceSuccessorCommitInput, ): Promise; @@ -74,6 +85,7 @@ interface AuthoritySource { interface AuthorityCapabilityRecord extends AuthoritySource { readonly ownerToken: object; + readonly verifyCandidate: (candidateOutcome: object) => WorkspaceSuccessorAuthorityInput; } const sources = new WeakMap(); @@ -98,6 +110,17 @@ export function registerExecutionStoresWorkspaceMutationSourceInternal( sources.set(stores, Object.freeze({ store, rootId })); } +/** Test-only bridge used by cross-package production-shaped owner tests. */ +export function commitExecutionStoresWorkspaceBaselineForTestInternal( + stores: object, + input: WorkspaceBaselineAuthorityInput, +): Promise { + const source = sources.get(stores); + if (!source) throw new Error('Execution stores workspace baseline test source is unavailable'); + adoptWorkspaceBaselineAuthorityStoreRootInternal(source.store, source.rootId); + return commitWorkspaceBaselineInternal(source.store, input); +} + export function issueExecutionStoresWorkspaceMutationAuthorityInternal(input: { readonly ownerToken: object; readonly stores: object; @@ -106,13 +129,17 @@ export function issueExecutionStoresWorkspaceMutationAuthorityInternal(input: { const source = sources.get(input.stores); if (!source) throw new Error('Execution stores workspace mutation source is unavailable'); adoptWorkspaceBaselineAuthorityStoreRootInternal(source.store, source.rootId); - registerWorkspaceSuccessorCandidateVerifierInternal(source.store, input.verifyCandidate); const capability = Object.freeze({ kind: 'execution_stores_workspace_mutation_authority_v1' as const, }); capabilities.set( capability, - Object.freeze({ ownerToken: input.ownerToken, store: source.store, rootId: source.rootId }), + Object.freeze({ + ownerToken: input.ownerToken, + store: source.store, + rootId: source.rootId, + verifyCandidate: input.verifyCandidate, + }), ); return capability; } @@ -127,14 +154,22 @@ export function requireExecutionStoresWorkspaceMutationAuthorityInternal( } const store = record.store; return Object.freeze({ + readEpoch: (workspaceId: string, workspaceEpochId: string) => + readWorkspaceEpochInternal(store, workspaceId, workspaceEpochId), readHead: (workspaceId: string, workspaceEpochId: string) => readWorkspaceHeadInternal(store, workspaceId, workspaceEpochId), readVersion: (workspaceVersionId: string) => readWorkspaceVersionInternal(store, workspaceVersionId), readActiveMutation: (workspaceInstanceId: string) => readActiveManagedMutationInternal(store, workspaceInstanceId), + readMutationEvidence: (operationId: string) => + readManagedMutationEvidenceInternal(store, operationId), commitSuccessor: async (input: WorkspaceSuccessorCommitInput) => { - const result = await commitWorkspaceSuccessorInternal(store, input); + const successor = record.verifyCandidate(input.candidateOutcome); + const result = await commitVerifiedWorkspaceSuccessorInternal(store, { + successor, + toolOutcome: input.toolOutcome, + }); const projectionCapability = Object.freeze({ kind: 'workspace_successor_projection_capability_v1' as const, }); diff --git a/packages/storage/src/sqlite-runtime-store.ts b/packages/storage/src/sqlite-runtime-store.ts index 31931c221f..6fb9634982 100644 --- a/packages/storage/src/sqlite-runtime-store.ts +++ b/packages/storage/src/sqlite-runtime-store.ts @@ -1534,6 +1534,7 @@ export class SqliteRuntimeStore } private registerWorkspaceBaselineAuthorityWriter(databasePath: string): void { + const readWorkspaceEpoch = this.readWorkspaceEpoch.bind(this); const readWorkspaceHead = this.readWorkspaceHead.bind(this); const readWorkspaceVersion = this.readWorkspaceVersion.bind(this); registerWorkspaceBaselineAuthorityWriterInternal( @@ -1544,12 +1545,31 @@ export class SqliteRuntimeStore (input, rootId) => this.#commitManagedMutationTerminal(input, rootId), (rootId) => this.#bindWorkspaceStorageRoot(rootId), (rootId) => this.#adoptWorkspaceStorageRoot(rootId), + readWorkspaceEpoch, readWorkspaceHead, readWorkspaceVersion, (workspaceInstanceId) => this.#readActiveManagedMutation(workspaceInstanceId), + (operationId) => this.#readManagedMutationEvidence(operationId), ); } + async #readManagedMutationEvidence( + operationId: string, + ): Promise< + import('./workspace-version-authority-internal.js').ManagedMutationEvidenceRecordV1 | undefined + > { + return this.readTransaction(() => { + const operation = this.readToolOperationSync(operationId); + if (!operation?.dispatchEventId || !operation.resultEventId) return undefined; + return Object.freeze({ + operationId, + callEvent: this.readRequiredRuntimeEvent(operation.callEventId), + dispatchEvent: this.readRequiredRuntimeEvent(operation.dispatchEventId), + outcomeEvent: this.readRequiredRuntimeEvent(operation.resultEventId), + }); + }); + } + async #readActiveManagedMutation( workspaceInstanceId: string, ): Promise< diff --git a/packages/storage/src/test-only/execution-stores-workspace-authority.ts b/packages/storage/src/test-only/execution-stores-workspace-authority.ts new file mode 100644 index 0000000000..5c0b483e62 --- /dev/null +++ b/packages/storage/src/test-only/execution-stores-workspace-authority.ts @@ -0,0 +1,20 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +export { commitExecutionStoresWorkspaceBaselineForTestInternal } from '../execution-stores-workspace-authority-internal.js'; diff --git a/packages/storage/src/workspace-version-authority-internal.ts b/packages/storage/src/workspace-version-authority-internal.ts index e90b302e51..4df187b76d 100644 --- a/packages/storage/src/workspace-version-authority-internal.ts +++ b/packages/storage/src/workspace-version-authority-internal.ts @@ -20,6 +20,7 @@ import type { WorkspaceBaselineAuthorityInput, WorkspaceBaselineCommitResult, + WorkspaceEpochRecordV1, WorkspaceHeadRecordV1, WorkspaceSuccessorAuthorityInput, WorkspaceVersionRecordV1, @@ -46,7 +47,7 @@ export interface WorkspaceSuccessorCommitInput { committedAt: number; }; } -interface VerifiedWorkspaceSuccessorCommitInput { +export interface VerifiedWorkspaceSuccessorCommitInput { successor: WorkspaceSuccessorAuthorityInput; toolOutcome: WorkspaceSuccessorCommitInput['toolOutcome']; } @@ -77,6 +78,10 @@ type WorkspaceHeadReader = ( workspaceId: string, workspaceEpochId: string, ) => Promise; +type WorkspaceEpochReader = ( + workspaceId: string, + workspaceEpochId: string, +) => Promise; type WorkspaceVersionReader = ( workspaceVersionId: string, ) => Promise; @@ -96,18 +101,29 @@ export interface ManagedMutationReservationRecordV1 { readonly executionProfileDigest: string; readonly reservedAt: number; } +export interface ManagedMutationEvidenceRecordV1 { + readonly operationId: string; + readonly callEvent: RuntimeEvent; + readonly dispatchEvent: RuntimeEvent; + readonly outcomeEvent: RuntimeEvent; +} type ManagedMutationReservationReader = ( workspaceInstanceId: string, ) => Promise; +type ManagedMutationEvidenceReader = ( + operationId: string, +) => Promise; interface WorkspaceBaselineAuthorityRegistration { readonly writer: WorkspaceBaselineAuthorityWriter; readonly successorWriter: WorkspaceSuccessorAuthorityWriter; candidateVerifier?: WorkspaceSuccessorCandidateVerifier; readonly terminalWriter: ManagedMutationTerminalAuthorityWriter; + readonly readEpoch: WorkspaceEpochReader; readonly readHead: WorkspaceHeadReader; readonly readVersion: WorkspaceVersionReader; readonly readActiveManagedMutation: ManagedMutationReservationReader; + readonly readManagedMutationEvidence: ManagedMutationEvidenceReader; readonly bindStorageRoot: WorkspaceStorageRootBinder; readonly adoptStorageRoot: WorkspaceStorageRootAdopter; readonly databasePath: string; @@ -128,9 +144,11 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( terminalWriter: ManagedMutationTerminalAuthorityWriter, bindStorageRoot: WorkspaceStorageRootBinder, adoptStorageRoot: WorkspaceStorageRootAdopter, + readEpoch: WorkspaceEpochReader, readHead: WorkspaceHeadReader, readVersion: WorkspaceVersionReader, readActiveManagedMutation: ManagedMutationReservationReader, + readManagedMutationEvidence: ManagedMutationEvidenceReader, ): void { if (workspaceBaselineAuthorityWriters.has(store)) { throw new Error('Workspace baseline authority writer is already registered'); @@ -140,9 +158,11 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( writer, successorWriter, terminalWriter, + readEpoch, readHead, readVersion, readActiveManagedMutation, + readManagedMutationEvidence, bindStorageRoot, adoptStorageRoot, databasePath: resolvedDatabasePath, @@ -159,6 +179,25 @@ export function readActiveManagedMutationInternal( return registration.readActiveManagedMutation(workspaceInstanceId); } +export function readManagedMutationEvidenceInternal( + store: object, + operationId: string, +): Promise { + const registration = workspaceBaselineAuthorityWriters.get(store); + if (!registration) throw new Error('Managed mutation evidence reader is unavailable'); + return registration.readManagedMutationEvidence(operationId); +} + +export function readWorkspaceEpochInternal( + store: object, + workspaceId: string, + workspaceEpochId: string, +): Promise { + const registration = workspaceBaselineAuthorityWriters.get(store); + if (!registration) throw new Error('Workspace baseline authority reader is unavailable'); + return registration.readEpoch(workspaceId, workspaceEpochId); +} + export function readWorkspaceHeadInternal( store: object, workspaceId: string, @@ -209,10 +248,27 @@ export function commitWorkspaceSuccessorInternal( throw new Error('Workspace successor candidate verifier is unavailable'); } const successor = registration.candidateVerifier(input.candidateOutcome); - return registration.successorWriter( - { successor, toolOutcome: input.toolOutcome }, - registration.boundRootId, - ); + return commitVerifiedWorkspaceSuccessorInternal(store, { + successor, + toolOutcome: input.toolOutcome, + }); +} + +/** + * Storage-internal seam for an owner-bound candidate verifier. This module is + * not a package export; production callers can reach it only through the + * execution-stores capability that owns the verifier. + */ +export function commitVerifiedWorkspaceSuccessorInternal( + store: object, + input: VerifiedWorkspaceSuccessorCommitInput, +): 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 registerWorkspaceSuccessorCandidateVerifierInternal(