diff --git a/.github/workflows/gitoxide-helper-admission.yml b/.github/workflows/gitoxide-helper-admission.yml index 1bf0dc3662..c224e8ee90 100644 --- a/.github/workflows/gitoxide-helper-admission.yml +++ b/.github/workflows/gitoxide-helper-admission.yml @@ -34,6 +34,8 @@ on: - '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__/gitoxide-managed-session-owner-internal.test.ts' + - 'packages/runtime-host/src/__tests__/gitoxide-managed-continuation-crash.test.ts' + - 'packages/runtime-host/src/__tests__/fixtures/execution-host*.ts' - 'packages/runtime-host/src/__tests__/packaged-gitoxide-helper.test.ts' - 'packages/runtime-host/src/__tests__/hosted-execution-tool-profile.test.ts' - 'scripts/prepare-gitoxide-helper*' @@ -63,6 +65,8 @@ on: - '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__/gitoxide-managed-session-owner-internal.test.ts' + - 'packages/runtime-host/src/__tests__/gitoxide-managed-continuation-crash.test.ts' + - 'packages/runtime-host/src/__tests__/fixtures/execution-host*.ts' - 'packages/runtime-host/src/__tests__/packaged-gitoxide-helper.test.ts' - 'packages/runtime-host/src/__tests__/hosted-execution-tool-profile.test.ts' - 'scripts/prepare-gitoxide-helper*' @@ -126,6 +130,7 @@ jobs: 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 packages/runtime-host/dist/__tests__/gitoxide-managed-session-owner-internal.test.js + packages/runtime-host/dist/__tests__/gitoxide-managed-continuation-crash.test.js packages/runtime-host/dist/__tests__/packaged-gitoxide-helper.test.js packages/runtime-host/dist/__tests__/hosted-execution-tool-profile.test.js - name: Test packaged helper preparation diff --git a/docs/architecture/gitoxide-workspace-bound-continuation-v1.zh-CN.md b/docs/architecture/gitoxide-workspace-bound-continuation-v1.zh-CN.md new file mode 100644 index 0000000000..cd02369aa5 --- /dev/null +++ b/docs/architecture/gitoxide-workspace-bound-continuation-v1.zh-CN.md @@ -0,0 +1,73 @@ +# Gitoxide workspace-bound continuation v1 + +## 目标 + +本切片只证明一个主要不变量: + +> managed-coding-v1 的 continuation 必须同时绑定不可变 RuntimeEvent 前缀与同一 SQLite 事务读取到的 accepted Git workspace head;任一侧漂移都不得继续 provider dispatch。 + +普通 continuation 仍使用 `continuation_claim_v1`。Managed continuation 使用 +`continuation_claim_v2` 和 `continuation_source_v3`,不存在从 v2 静默降级为 v1 的路径。 + +## Owner 与权限 + +- Runtime owns:RuntimeEvent high-water、provider replay manifest、target Run identity。 +- SQLite owns:workspace epoch、accepted version/head、claim/start 的原子写入与重验。 +- Gitoxide helper owns:source HEAD 与 `refs/maka/accepted` 的只读复核。 +- Runtime Host owns:把 SQLite observation 与 Gitoxide observation 组合成 safety observation。 + +Host 只能通过 `execution_stores_workspace_continuation_authority_v1` 读取 continuation +boundary。该 opaque capability 不暴露 candidate、successor、T2 或 head 推进权限。 + +## 绑定内容 + +`managed_workspace_continuation_boundary_v1` 包含: + +- storage root、repository、workspace、epoch、instance identity; +- source commit/tree; +- accepted workspace version、accepted event、revision; +- accepted commit/tree; +- materialization、policy、execution profile digest。 + +最终 `boundaryDigest` 对 Runtime boundary 与上述 workspace boundary 一起做 canonical +SHA-256 commitment。`replayManifestDigest` 仍只表示 RuntimeEvent lineage,不能与 composite +boundary digest 混用。 + +## 原子性边界 + +1. continuation plan 从 immutable RuntimeEvents 建立 Runtime boundary; +2. SQLite 在单个 read transaction 中读取并校验 epoch/head/version/projection; +3. Host 用短生命周期 Gitoxide helper 重验 source HEAD 与 accepted ref; +4. claim transaction 再次重验 workspace boundary 后写入 claim; +5. continuation-start transaction 再次重验 workspace boundary后写入首个 RuntimeEvent; +6. 只有以上步骤完成后允许 provider dispatch。 + +Git 与 SQLite 不组成跨系统事务。Git observation 只作为 admission proof;accepted truth 仍是 +RuntimeEvents + 可重建 workspace authority projection。 + +## 失败状态与恢复 + +- boundary 缺失、helper 缺失、source/accepted-ref 漂移:park/fail closed; +- claim 已提交但 start 未提交:使用专用 repair writer,不调用 provider; +- start 已提交而 Host 退出:标记 `continuation_started_indeterminate`,禁止 provider replay; +- provider 已被调用但响应未知:保持 indeterminate,后续显式重试也不得再次调用 provider; +- malformed claim、row/payload mismatch、projection corruption:拒绝 session continuation authority。 + +回滚仅删除尚未提交的事务内状态。已经持久化的 claim/start 不回滚,通过重验与专用 repair +路径收敛。 + +## 平台矩阵 + +| 平台 | v1 承诺 | +| --- | --- | +| Linux | SQLite claim/start crash convergence;Gitoxide source/ref 重验;Host kill 后 provider-at-most-once | +| macOS | 与 Linux 相同;路径先 realpath canonicalize | +| Windows | 与 Linux 相同;进程终止证据由专用 Gitoxide workflow 执行 | + +三平台共用 `.github/workflows/gitoxide-helper-admission.yml` 中的真实 helper/Host crash test。 +没有系统 Git 或 PATH fallback。 + +## 产品范围 + +该能力只对显式 `managed-coding-v1` 生效。普通会话行为不变。当前 UI 是否开放该 profile 是独立 +产品决策;协议与 Host consumer 已真实接通,不依赖 Desktop 控件存在。 diff --git a/packages/core/src/__tests__/agent-run-continuation-source.test.ts b/packages/core/src/__tests__/agent-run-continuation-source.test.ts index b8c1eba5c2..e489a96c4e 100644 --- a/packages/core/src/__tests__/agent-run-continuation-source.test.ts +++ b/packages/core/src/__tests__/agent-run-continuation-source.test.ts @@ -75,6 +75,19 @@ describe('AgentRun continuation source decoding', () => { /Invalid AgentRun header schema/, ); }); + + it('accepts a V3 workspace-bound source with a distinct replay manifest digest', () => { + const source = { + ...validV2ContinuationSource(), + protocol: 'continuation_source_v3' as const, + replayManifestDigest: `sha256:${'c'.repeat(64)}` as const, + }; + + assert.deepEqual( + decodeAgentRunHeader(headerWithContinuation(source)).continuationSource, + source, + ); + }); }); function headerWithContinuation( diff --git a/packages/core/src/__tests__/runtime-boundary.test.ts b/packages/core/src/__tests__/runtime-boundary.test.ts index 5ec71103a3..a3b92f2a4c 100644 --- a/packages/core/src/__tests__/runtime-boundary.test.ts +++ b/packages/core/src/__tests__/runtime-boundary.test.ts @@ -24,6 +24,7 @@ import { buildImmutableRuntimePrefix, createRuntimeBoundaryCursor, decodeContinuationClaim, + digestWorkspaceBoundContinuationBoundary, runtimePrefixSegment, type RuntimeBoundaryCursorV1, type RuntimePrefixIdentityV1, @@ -298,8 +299,85 @@ describe('immutable RuntimeEvent boundary', () => { /target turnId reuses source identity/, ); }); + + it('binds a v2 continuation claim to one exact accepted workspace head', () => { + const boundary = boundaryForRuns('run-source'); + const workspaceBoundary = workspaceBoundaryV1(); + const claim = claimForBoundary(boundary); + const boundaryDigest = digestWorkspaceBoundContinuationBoundary(boundary, workspaceBoundary); + + const decoded = decodeContinuationClaim({ + ...claim, + protocol: 'continuation_claim_v2', + boundaryDigest, + workspaceBoundary, + targetRunHeader: { + ...claim.targetRunHeader, + continuationSource: { + ...claim.targetRunHeader.continuationSource, + protocol: 'continuation_source_v3', + boundaryDigest, + }, + }, + }); + + assert.equal(decoded.protocol, 'continuation_claim_v2'); + if (decoded.protocol !== 'continuation_claim_v2') return; + assert.deepEqual(decoded.workspaceBoundary, workspaceBoundary); + }); + + it('rejects a v2 claim when the accepted workspace version changes', () => { + const boundary = boundaryForRuns('run-source'); + const workspaceBoundary = workspaceBoundaryV1(); + const claim = claimForBoundary(boundary); + const boundaryDigest = digestWorkspaceBoundContinuationBoundary(boundary, workspaceBoundary); + + assert.throws( + () => + decodeContinuationClaim({ + ...claim, + protocol: 'continuation_claim_v2', + boundaryDigest, + workspaceBoundary: { + ...workspaceBoundary, + workspaceVersionId: `version_${'9'.repeat(32)}`, + }, + targetRunHeader: { + ...claim.targetRunHeader, + continuationSource: { + ...claim.targetRunHeader.continuationSource, + protocol: 'continuation_source_v3', + boundaryDigest, + }, + }, + }), + /workspace boundary digest mismatch/, + ); + }); }); +function workspaceBoundaryV1() { + return { + protocol: 'managed_workspace_continuation_boundary_v1' as const, + storageRootId: '1'.repeat(64), + repositoryId: `repository_${'2'.repeat(32)}`, + workspaceId: `workspace_${'3'.repeat(32)}`, + workspaceEpochId: `epoch_${'4'.repeat(32)}`, + workspaceInstanceId: `instance_${'5'.repeat(32)}`, + workspaceVersionId: `version_${'6'.repeat(32)}`, + acceptedEventId: 'workspace-accepted-event-1', + revision: 2, + objectFormat: 'sha1' as const, + sourceCommitOid: '5'.repeat(40), + sourceTreeOid: '6'.repeat(40), + commitOid: '7'.repeat(40), + treeOid: '8'.repeat(40), + materializationProfileDigest: `sha256:${'9'.repeat(64)}` as const, + policyHash: `sha256:${'a'.repeat(64)}` as const, + executionProfileDigest: `sha256:${'b'.repeat(64)}` as const, + }; +} + function runtimeIdentity(runId: string): RuntimePrefixIdentityV1 { return { sessionId: 'session-1', diff --git a/packages/core/src/agent-run.ts b/packages/core/src/agent-run.ts index 5ef38f8e25..fdf50e14ba 100644 --- a/packages/core/src/agent-run.ts +++ b/packages/core/src/agent-run.ts @@ -70,9 +70,20 @@ export interface AgentRunContinuationSourceV2 extends AgentRunContinuationSource replayManifestDigest: `sha256:${string}`; } +export interface AgentRunContinuationSourceV3 extends AgentRunContinuationSourceV1 { + protocol: 'continuation_source_v3'; + claimId: string; + /** Composite RuntimeEvent + accepted workspace boundary identity. */ + boundaryDigest: `sha256:${string}`; + sourcePrefixDigest: `sha256:${string}`; + /** RuntimeEvent-only replay lineage identity. */ + replayManifestDigest: `sha256:${string}`; +} + export type AgentRunContinuationSource = | AgentRunContinuationSourceV1 - | AgentRunContinuationSourceV2; + | AgentRunContinuationSourceV2 + | AgentRunContinuationSourceV3; export type RootExecutionDescriptor = | { @@ -104,6 +115,7 @@ export type RootExecutionDescriptor = sourceRuntimeEventHighWater: number; claimId: string; boundaryDigest: `sha256:${string}`; + replayManifestDigest?: `sha256:${string}`; providerReplayDigest: `sha256:${string}`; safetyDigest: `sha256:${string}`; targetInvocationId: string; @@ -150,6 +162,20 @@ const AGENT_RUN_CONTINUATION_SOURCE_V2_SHAPE = defineObjectShape()( + [ + 'protocol', + 'claimId', + 'boundaryDigest', + 'sourceInvocationId', + 'sourceRunId', + 'sourceTurnId', + 'sourceRuntimeEventHighWater', + 'sourcePrefixDigest', + 'replayManifestDigest', + ], + [], +); export interface AgentRunHeader { runId: string; @@ -282,14 +308,16 @@ export function agentRunMatchesHostedRootExecution( run.parentTurnId === execution.sourceTurnId && source !== undefined && 'protocol' in source && - source.protocol === 'continuation_source_v2' && + (source.protocol === 'continuation_source_v2' || + source.protocol === 'continuation_source_v3') && source.sourceInvocationId === execution.sourceInvocationId && source.sourceRunId === execution.sourceRunId && source.sourceTurnId === execution.sourceTurnId && source.sourceRuntimeEventHighWater === execution.sourceRuntimeEventHighWater && source.claimId === execution.claimId && source.boundaryDigest === execution.boundaryDigest && - source.replayManifestDigest === execution.boundaryDigest && + source.replayManifestDigest === + (execution.replayManifestDigest ?? execution.boundaryDigest) && run.resumedFromRunId === undefined && run.retriedFromRunId === undefined && run.agentId === undefined && @@ -718,9 +746,14 @@ function isAgentRunContinuationSource(value: unknown): value is AgentRunContinua value.sourceRuntimeEventHighWater >= 0; if (!common) return false; if (hasExactShape(value, AGENT_RUN_CONTINUATION_SOURCE_V1_SHAPE)) return true; - return ( + const isV2 = hasExactShape(value, AGENT_RUN_CONTINUATION_SOURCE_V2_SHAPE) && - value.protocol === 'continuation_source_v2' && + value.protocol === 'continuation_source_v2'; + const isV3 = + hasExactShape(value, AGENT_RUN_CONTINUATION_SOURCE_V3_SHAPE) && + value.protocol === 'continuation_source_v3'; + return ( + (isV2 || isV3) && typeof value.claimId === 'string' && value.claimId.length > 0 && typeof value.sourceInvocationId === 'string' && @@ -734,7 +767,7 @@ function isAgentRunContinuationSource(value: unknown): value is AgentRunContinua isSha256Digest(value.boundaryDigest) && isSha256Digest(value.sourcePrefixDigest) && isSha256Digest(value.replayManifestDigest) && - value.replayManifestDigest === value.boundaryDigest + (!isV2 || value.replayManifestDigest === value.boundaryDigest) ); } diff --git a/packages/core/src/runtime-boundary.ts b/packages/core/src/runtime-boundary.ts index e444a8862a..3f49be96be 100644 --- a/packages/core/src/runtime-boundary.ts +++ b/packages/core/src/runtime-boundary.ts @@ -84,6 +84,33 @@ export interface ContinuationClaimV1 { claimedAt: number; } +export interface ManagedWorkspaceContinuationBoundaryV1 { + protocol: 'managed_workspace_continuation_boundary_v1'; + storageRootId: string; + repositoryId: string; + workspaceId: string; + workspaceEpochId: string; + workspaceInstanceId: string; + workspaceVersionId: string; + acceptedEventId: string; + revision: number; + objectFormat: 'sha1' | 'sha256'; + sourceCommitOid: string; + sourceTreeOid: string; + commitOid: string; + treeOid: string; + materializationProfileDigest: RuntimeBoundaryDigest; + policyHash: RuntimeBoundaryDigest; + executionProfileDigest: RuntimeBoundaryDigest; +} + +export interface ContinuationClaimV2 extends Omit { + protocol: 'continuation_claim_v2'; + workspaceBoundary: ManagedWorkspaceContinuationBoundaryV1; +} + +export type ContinuationClaim = ContinuationClaimV1 | ContinuationClaimV2; + export function buildImmutableRuntimePrefix( identity: RuntimePrefixIdentityV1, rows: readonly RuntimePrefixRowV1[], @@ -182,6 +209,28 @@ export function digestRuntimeBoundaryManifest( return `sha256:${hash.digest('hex')}`; } +export function digestWorkspaceBoundContinuationBoundary( + boundary: RuntimeBoundaryCursorV1, + workspaceBoundary: ManagedWorkspaceContinuationBoundaryV1, +): RuntimeBoundaryDigest { + const runtime = decodeRuntimeBoundaryCursor(boundary); + const workspace = decodeManagedWorkspaceContinuationBoundary(workspaceBoundary); + const hash = nodeCrypto.createHash('sha256'); + updateLengthPrefixed(hash, Buffer.from('maka.workspace-bound-continuation.v1', 'utf8')); + updateLengthPrefixed( + hash, + Buffer.from( + stableJsonStringify({ + protocol: 'workspace_bound_continuation_boundary_v1', + runtime, + workspace, + }), + 'utf8', + ), + ); + return `sha256:${hash.digest('hex')}`; +} + export function decodeRuntimePrefixSegment(value: unknown): RuntimePrefixSegmentV1 { if ( !isRecord(value) || @@ -220,21 +269,24 @@ export function decodeRuntimeBoundaryCursor(value: unknown): RuntimeBoundaryCurs return cursor; } -export function decodeContinuationClaim(value: unknown): ContinuationClaimV1 { +export function decodeContinuationClaim(value: unknown): ContinuationClaim { + const isV2 = isRecord(value) && value.protocol === 'continuation_claim_v2'; + const expectedKeys = [ + 'protocol', + 'claimId', + 'boundaryDigest', + 'boundary', + ...(isV2 ? ['workspaceBoundary'] : []), + 'providerProjectionVersion', + 'providerReplayDigest', + 'target', + 'targetRunHeader', + 'claimedAt', + ]; if ( !isRecord(value) || - !hasExactKeys(value, [ - 'protocol', - 'claimId', - 'boundaryDigest', - 'boundary', - 'providerProjectionVersion', - 'providerReplayDigest', - 'target', - 'targetRunHeader', - 'claimedAt', - ]) || - value.protocol !== 'continuation_claim_v1' || + !hasExactKeys(value, expectedKeys) || + (value.protocol !== 'continuation_claim_v1' && value.protocol !== 'continuation_claim_v2') || !isNonEmptyString(value.claimId) || !isRecord(value.target) || !hasExactKeys(value.target, ['sessionId', 'invocationId', 'runId', 'turnId']) || @@ -251,8 +303,18 @@ export function decodeContinuationClaim(value: unknown): ContinuationClaimV1 { const boundary = decodeRuntimeBoundaryCursor(value.boundary); const boundaryDigest = decodeBoundaryDigest(value.boundaryDigest); const providerReplayDigest = decodeBoundaryDigest(value.providerReplayDigest); - if (boundaryDigest !== boundary.manifestDigest) { - throw new Error('Continuation claim boundary digest mismatch'); + const workspaceBoundary = isV2 + ? decodeManagedWorkspaceContinuationBoundary(value.workspaceBoundary) + : undefined; + const expectedBoundaryDigest = workspaceBoundary + ? digestWorkspaceBoundContinuationBoundary(boundary, workspaceBoundary) + : boundary.manifestDigest; + if (boundaryDigest !== expectedBoundaryDigest) { + throw new Error( + workspaceBoundary + ? 'Continuation claim workspace boundary digest mismatch' + : 'Continuation claim boundary digest mismatch', + ); } const source = boundary.segments.at(-1)!; if (value.target.sessionId !== source.identity.sessionId) { @@ -285,7 +347,8 @@ export function decodeContinuationClaim(value: unknown): ContinuationClaimV1 { targetRunHeader.failureMessage !== undefined || !continuationSource || !('protocol' in continuationSource) || - continuationSource.protocol !== 'continuation_source_v2' || + continuationSource.protocol !== + (workspaceBoundary ? 'continuation_source_v3' : 'continuation_source_v2') || continuationSource.claimId !== value.claimId || continuationSource.boundaryDigest !== boundaryDigest || continuationSource.sourceInvocationId !== source.identity.invocationId || @@ -297,11 +360,14 @@ export function decodeContinuationClaim(value: unknown): ContinuationClaimV1 { ) { throw new Error('Continuation claim target Run header mismatch'); } - return { - protocol: 'continuation_claim_v1', + const claim = { + protocol: workspaceBoundary + ? ('continuation_claim_v2' as const) + : ('continuation_claim_v1' as const), claimId: value.claimId, boundaryDigest, boundary, + ...(workspaceBoundary ? { workspaceBoundary } : {}), providerProjectionVersion: 1, providerReplayDigest, target: { @@ -313,6 +379,92 @@ export function decodeContinuationClaim(value: unknown): ContinuationClaimV1 { targetRunHeader, claimedAt: value.claimedAt as number, }; + return claim as ContinuationClaim; +} + +export function decodeManagedWorkspaceContinuationBoundary( + value: unknown, +): ManagedWorkspaceContinuationBoundaryV1 { + const keys = [ + 'protocol', + 'storageRootId', + 'repositoryId', + 'workspaceId', + 'workspaceEpochId', + 'workspaceInstanceId', + 'workspaceVersionId', + 'acceptedEventId', + 'revision', + 'objectFormat', + 'sourceCommitOid', + 'sourceTreeOid', + 'commitOid', + 'treeOid', + 'materializationProfileDigest', + 'policyHash', + 'executionProfileDigest', + ]; + if ( + !isRecord(value) || + !hasExactKeys(value, keys) || + value.protocol !== 'managed_workspace_continuation_boundary_v1' || + typeof value.storageRootId !== 'string' || + !/^[0-9a-f]{64}$/u.test(value.storageRootId) || + typeof value.repositoryId !== 'string' || + !/^repository_[0-9a-f]{32}$/u.test(value.repositoryId) || + typeof value.workspaceId !== 'string' || + !/^workspace_[0-9a-f]{32}$/u.test(value.workspaceId) || + typeof value.workspaceEpochId !== 'string' || + !/^epoch_[0-9a-f]{32}$/u.test(value.workspaceEpochId) || + typeof value.workspaceInstanceId !== 'string' || + !/^instance_[0-9a-f]{32}$/u.test(value.workspaceInstanceId) || + typeof value.workspaceVersionId !== 'string' || + !/^version_[0-9a-f]{32}$/u.test(value.workspaceVersionId) || + !isNonEmptyString(value.acceptedEventId) || + !Number.isSafeInteger(value.revision) || + (value.revision as number) <= 0 || + (value.objectFormat !== 'sha1' && value.objectFormat !== 'sha256') || + typeof value.sourceCommitOid !== 'string' || + typeof value.sourceTreeOid !== 'string' || + !oidMatchesFormat(value.sourceCommitOid, value.objectFormat) || + !oidMatchesFormat(value.sourceTreeOid, value.objectFormat) || + typeof value.commitOid !== 'string' || + typeof value.treeOid !== 'string' || + !oidMatchesFormat(value.commitOid, value.objectFormat) || + !oidMatchesFormat(value.treeOid, value.objectFormat) || + !isBoundaryDigest(value.materializationProfileDigest) || + !isBoundaryDigest(value.policyHash) || + !isBoundaryDigest(value.executionProfileDigest) + ) { + throw new Error('Invalid managed workspace continuation boundary'); + } + return { + protocol: 'managed_workspace_continuation_boundary_v1', + storageRootId: value.storageRootId, + repositoryId: value.repositoryId, + workspaceId: value.workspaceId, + workspaceEpochId: value.workspaceEpochId, + workspaceInstanceId: value.workspaceInstanceId, + workspaceVersionId: value.workspaceVersionId, + acceptedEventId: value.acceptedEventId, + revision: value.revision as number, + objectFormat: value.objectFormat, + sourceCommitOid: value.sourceCommitOid, + sourceTreeOid: value.sourceTreeOid, + commitOid: value.commitOid, + treeOid: value.treeOid, + materializationProfileDigest: value.materializationProfileDigest, + policyHash: value.policyHash, + executionProfileDigest: value.executionProfileDigest, + }; +} + +function oidMatchesFormat(value: string, format: 'sha1' | 'sha256'): boolean { + return new RegExp(`^[0-9a-f]{${format === 'sha1' ? 40 : 64}}$`, 'u').test(value); +} + +function isBoundaryDigest(value: unknown): value is RuntimeBoundaryDigest { + return typeof value === 'string' && /^sha256:[0-9a-f]{64}$/u.test(value); } function canonicalizePrefixRows( diff --git a/packages/core/src/runtime-event-store.ts b/packages/core/src/runtime-event-store.ts index 9bb29085ba..bf3a48557f 100644 --- a/packages/core/src/runtime-event-store.ts +++ b/packages/core/src/runtime-event-store.ts @@ -20,6 +20,7 @@ import type { RuntimeEvent } from './runtime-event.js'; import type { ContinuationClaimV1, + ContinuationClaimV2, ImmutableRuntimePrefixV1, RuntimeBoundaryDigest, } from './runtime-boundary.js'; @@ -33,6 +34,8 @@ import { WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1 } from './workspace-version-a export const TOOL_RECOVERY_BUNDLE_CAPABILITY_V1 = 'tool_recovery_bundle_v1' as const; export const RUNTIME_CONTINUATION_AUTHORITY_V1 = 'runtime_continuation_authority_v1' as const; +export const RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_V1 = + 'runtime_workspace_bound_continuation_authority_v1' as const; export interface RuntimeRecoveryBundleCommit { operationId: string; @@ -150,6 +153,41 @@ export interface RuntimeContinuationAuthorityStore extends RuntimeEventStore { }): Promise<{ created: boolean; runtimeEventSeq: number }>; } +export type WorkspaceBoundContinuationClaimResult = + | { kind: 'acquired'; claim: ContinuationClaimV2 } + | { kind: 'existing'; claim: ContinuationClaimV2 } + | { kind: 'conflict'; claim: ContinuationClaimV2 }; + +export interface RuntimeWorkspaceBoundContinuationAuthorityStore extends RuntimeEventStore { + readonly workspaceBoundContinuationAuthorityCapability: typeof RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_V1; + claimWorkspaceBoundContinuation(input: { + claim: ContinuationClaimV2; + }): Promise; + readWorkspaceBoundContinuationClaimByBoundary( + boundaryDigest: RuntimeBoundaryDigest, + ): Promise; + readWorkspaceBoundContinuationClaimStateByBoundary( + boundaryDigest: RuntimeBoundaryDigest, + ): Promise; + listWorkspaceBoundContinuationClaimsForRecovery( + sessionId: string, + ): Promise; + commitWorkspaceBoundContinuationStart(input: { + claim: ContinuationClaimV2; + event: RuntimeEvent; + }): Promise<{ created: boolean; runtimeEventSeq: number }>; + commitWorkspaceBoundContinuationRepairStart(input: { + claim: ContinuationClaimV2; + event: RuntimeEvent; + }): Promise<{ created: boolean; runtimeEventSeq: number }>; +} + +export interface ContinuationClaimStateV2 { + claim: ContinuationClaimV2; + startEventId?: string; + startKind?: 'runtime_admission' | 'claim_repair'; +} + export interface RuntimeWorkspaceVersionAuthorityStore extends RuntimeEventStore { readonly workspaceVersionAuthorityCapability: typeof WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1; readWorkspaceEpoch( diff --git a/packages/core/src/workspace-version-authority.ts b/packages/core/src/workspace-version-authority.ts index 386f0dff09..8eb3e7060c 100644 --- a/packages/core/src/workspace-version-authority.ts +++ b/packages/core/src/workspace-version-authority.ts @@ -17,6 +17,7 @@ * under the License. */ +import { createHash } from 'node:crypto'; import { isCanonicalManagedMutationPathV1, type RuntimeEvent } from './runtime-event.js'; export const WORKSPACE_EPOCH_OPENED_FACT_KIND = 'maka.workspace.epoch_opened' as const; @@ -31,6 +32,25 @@ export const WORKSPACE_MATERIALIZATION_SEMANTICS_V1 = export type WorkspaceGitObjectFormat = 'sha1' | 'sha256'; +/** Durable commitment to both Git materialization and Runtime transform semantics. */ +export function workspaceMutationPolicyHashV1( + materializationProfileDigest: `sha256:${string}`, + mutationExecutionProfileDigest: `sha256:${string}`, +): `sha256:${string}` { + if ( + !isSha256Digest(materializationProfileDigest) || + !isSha256Digest(mutationExecutionProfileDigest) + ) { + throw new Error('Invalid workspace mutation policy profile digest'); + } + const hash = createHash('sha256'); + hash.update('maka.workspace-mutation-policy.v1\0', 'utf8'); + hash.update(materializationProfileDigest, 'utf8'); + hash.update('\0', 'utf8'); + hash.update(mutationExecutionProfileDigest, 'utf8'); + return `sha256:${hash.digest('hex')}`; +} + export interface WorkspaceEpochDescriptorV1 { repositoryId: string; workspaceId: string; diff --git a/packages/runtime-host/src/__tests__/fixtures/execution-host-suite.ts b/packages/runtime-host/src/__tests__/fixtures/execution-host-suite.ts index 32f8ea3ffb..8e4953bc95 100644 --- a/packages/runtime-host/src/__tests__/fixtures/execution-host-suite.ts +++ b/packages/runtime-host/src/__tests__/fixtures/execution-host-suite.ts @@ -126,6 +126,16 @@ export interface ExecutionHostHandle { recoveryOutcome?: RuntimeEvent; } +export interface ExecutionHostTestOptions { + readonly packagedResourcesRoot?: string; + readonly providerCallLogPath?: string; + readonly continuationFailpoint?: + | 'after_continuation_claim_committed' + | 'after_run_created' + | 'after_continuation_start_committed'; + readonly providerFailpointAfterSend?: boolean; +} + export interface TurnLedger { runs: AgentRunHeader[]; userMessages: Array>; @@ -1045,8 +1055,9 @@ export class ExecutionFixture { runId: string; }, safeBoundaryResumeEnabled = true, + testOptions: ExecutionHostTestOptions = {}, ): Promise { - const child = this.spawnHost('inherit', recoveryProbe, safeBoundaryResumeEnabled); + const child = this.spawnHost('inherit', recoveryProbe, safeBoundaryResumeEnabled, testOptions); const ready = await waitForHostReady(child); return { child, ...ready }; } @@ -1233,10 +1244,31 @@ export class ExecutionFixture { stderr: 'inherit' | 'ignore', recoveryProbe?: { sessionId: string; runId: string }, safeBoundaryResumeEnabled = true, + testOptions: ExecutionHostTestOptions = {}, ): ChildProcess { const env = { ...process.env }; if (safeBoundaryResumeEnabled) env.MAKA_RUNTIME_SAFE_BOUNDARY_RESUME = '1'; else delete env.MAKA_RUNTIME_SAFE_BOUNDARY_RESUME; + if (testOptions.packagedResourcesRoot) { + env.MAKA_TEST_PACKAGED_RESOURCES_ROOT = testOptions.packagedResourcesRoot; + } else { + delete env.MAKA_TEST_PACKAGED_RESOURCES_ROOT; + } + if (testOptions.providerCallLogPath) { + env.MAKA_TEST_PROVIDER_CALL_LOG = testOptions.providerCallLogPath; + } else { + delete env.MAKA_TEST_PROVIDER_CALL_LOG; + } + if (testOptions.continuationFailpoint) { + env.MAKA_TEST_CONTINUATION_FAILPOINT = testOptions.continuationFailpoint; + } else { + delete env.MAKA_TEST_CONTINUATION_FAILPOINT; + } + if (testOptions.providerFailpointAfterSend) { + env.MAKA_TEST_PROVIDER_FAILPOINT_AFTER_SEND = '1'; + } else { + delete env.MAKA_TEST_PROVIDER_FAILPOINT_AFTER_SEND; + } const child = fork( new URL('./execution-host.js', import.meta.url), [ diff --git a/packages/runtime-host/src/__tests__/fixtures/execution-host.ts b/packages/runtime-host/src/__tests__/fixtures/execution-host.ts index b18f6d7ce5..70133e86c8 100644 --- a/packages/runtime-host/src/__tests__/fixtures/execution-host.ts +++ b/packages/runtime-host/src/__tests__/fixtures/execution-host.ts @@ -17,7 +17,8 @@ * under the License. */ -import { join } from 'node:path'; +import { appendFile } from 'node:fs/promises'; +import { isAbsolute, join } from 'node:path'; import { inspect } from 'node:util'; import { FakeBackend } from '@maka/runtime/test-only/fake-backend'; import { createSqliteRuntimeStore } from '@maka/storage/sqlite-runtime-store'; @@ -40,6 +41,38 @@ if (!Number.isSafeInteger(idleGraceMs) || idleGraceMs < 0) { throw new Error('execution-host requires a non-negative idle grace'); } +const packagedResourcesRoot = process.env.MAKA_TEST_PACKAGED_RESOURCES_ROOT; +if (packagedResourcesRoot) { + if (!isAbsolute(packagedResourcesRoot)) { + throw new Error('MAKA_TEST_PACKAGED_RESOURCES_ROOT must be absolute'); + } + Object.defineProperty(process.versions, 'electron', { + configurable: true, + value: 'test-runtime-host', + }); + Object.defineProperty(process, 'resourcesPath', { + configurable: true, + value: packagedResourcesRoot, + }); +} + +const providerCallLogPath = process.env.MAKA_TEST_PROVIDER_CALL_LOG; +const continuationFailpoint = process.env.MAKA_TEST_CONTINUATION_FAILPOINT; +const providerFailpointAfterSend = process.env.MAKA_TEST_PROVIDER_FAILPOINT_AFTER_SEND === '1'; + +class ObservedFakeBackend extends FakeBackend { + override async *send(input: Parameters[0]) { + if (providerCallLogPath) { + await appendFile(providerCallLogPath, `${input.turnId}\n`, 'utf8'); + } + if (providerFailpointAfterSend) { + process.send?.({ type: 'test.provider_failpoint', point: 'after_send_called' }); + await new Promise(() => undefined); + } + yield* super.send(input); + } +} + // The production composition registers no test backend. This fixture is a // candidate host in its own right, so it supplies the deterministic one through // the composition's `primaryBackendFactory` seam — the same path Desktop E2E @@ -55,10 +88,17 @@ const result = await startExecutionRuntimeHostCandidate( createComposition: (context, compositionOptions) => createExecutionRuntimeHostComposition(context, compositionOptions, { primaryBackendFactory: (backendContext) => { - const backend = new FakeBackend(backendContext); + const backend = new ObservedFakeBackend(backendContext); fakeBackends.add(backend); return backend; }, + continuationFailpoint: continuationFailpoint + ? async (point) => { + if (point !== continuationFailpoint) return; + process.send?.({ type: 'test.continuation_failpoint', point }); + await new Promise(() => undefined); + } + : undefined, }), }, ); diff --git a/packages/runtime-host/src/__tests__/gitoxide-managed-continuation-crash.test.ts b/packages/runtime-host/src/__tests__/gitoxide-managed-continuation-crash.test.ts new file mode 100644 index 0000000000..077e794f8b --- /dev/null +++ b/packages/runtime-host/src/__tests__/gitoxide-managed-continuation-crash.test.ts @@ -0,0 +1,463 @@ +/* + * 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 type { ChildProcess } from 'node:child_process'; +import { execFileSync } from 'node:child_process'; +import { createHash } from 'node:crypto'; +import { + chmod, + copyFile, + mkdir, + mkdtemp, + readFile, + realpath, + rm, + stat, + writeFile, +} from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { test } from 'node:test'; +import { createSqliteRuntimeStore } from '@maka/storage/sqlite-runtime-store'; +import { openInteractiveExecutionStoresForWrite } from '@maka/storage/execution-stores'; +import { resolveStorageRoot, tryAcquireInteractiveRootOwner } from '@maka/storage/root-authority'; +import { + admitGitoxideHelperArtifactInternal, + issueGitoxideHelperReleaseArtifactClaimInternal, +} from '../server/gitoxide-helper-artifact-authority-internal.js'; +import { + inspectGitoxideManagedContinuationBoundaryInternal, + openGitoxideManagedSessionOwnerInternal, +} from '../server/gitoxide-managed-session-owner-internal.js'; +import { + connectClient, + ExecutionFixture, + PROCESS_TIMEOUT_MS, + withTimeout, +} from './fixtures/execution-host-suite.js'; + +test('a started workspace-bound continuation survives Host death without provider replay', async (t) => { + const helperPath = process.env.MAKA_GITOXIDE_HELPER_PATH; + if (!helperPath) { + t.skip( + 'MAKA_GITOXIDE_HELPER_PATH is required for the real helper continuation test', + ); + return; + } + await withManagedContinuationFixture( + helperPath, + async ({ fixture, resourcesRoot, callLog, boundary }) => { + const source = await fixture.seedSafeBoundaryContinuationSource(); + const crashHost = await fixture.startHost(undefined, true, { + packagedResourcesRoot: resourcesRoot, + providerCallLogPath: callLog, + continuationFailpoint: 'after_continuation_start_committed', + }); + const crashClient = await connectClient(fixture.root); + const targetTurnId = 'turn-workspace-bound-host-crash'; + try { + const initialPlan = await crashClient.request('turn.resume.query', { + sessionId: fixture.sessionId, + }); + assert.equal(initialPlan.disposition, 'ready', JSON.stringify(initialPlan)); + const failpoint = waitForContinuationFailpoint(crashHost.child); + const start = crashClient + .request('turn.resume.start', { + sessionId: fixture.sessionId, + turnId: targetTurnId, + sourceRunId: source.sourceRunId, + sourceRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }) + .then( + () => undefined, + () => undefined, + ); + await failpoint; + await fixture.killHost(crashHost); + await withTimeout(start, PROCESS_TIMEOUT_MS, 'crashed continuation request did not close'); + } finally { + await crashClient.close().catch(() => undefined); + } + + assert.equal(await readFile(callLog, 'utf8'), ''); + const admission = (await fixture.readAdmissionChain()).find( + (candidate) => candidate.turnId === targetTurnId, + ); + assert.equal(admission?.execution.kind, 'safe_boundary_continuation'); + if (admission?.execution.kind !== 'safe_boundary_continuation') { + assert.fail('Workspace-bound continuation admission is missing'); + } + + const successorHost = await fixture.startHost(undefined, true, { + packagedResourcesRoot: resourcesRoot, + providerCallLogPath: callLog, + }); + const successorClient = await connectClient(fixture.root); + try { + const target = await successorClient.request('turn.query', { + sessionId: fixture.sessionId, + turnId: targetTurnId, + }); + assert.equal(target.runId, admission.runId); + assert.equal(target.status, 'failed'); + if (target.status !== 'failed') assert.fail('Crashed continuation Run was not closed'); + assert.equal(target.failureClass, 'app_restarted'); + const plan = { + sessionId: fixture.sessionId, + disposition: 'parked' as const, + reason: 'continuation_started_indeterminate' as const, + }; + assert.deepEqual( + await successorClient.request('turn.resume.query', { + sessionId: fixture.sessionId, + sourceRunId: source.sourceRunId, + expectedRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }), + plan, + ); + assert.deepEqual( + await successorClient.request('turn.resume.start', { + sessionId: fixture.sessionId, + turnId: `${targetTurnId}-retry`, + sourceRunId: source.sourceRunId, + sourceRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }), + { kind: 'parked', plan }, + ); + assert.equal(await readFile(callLog, 'utf8'), ''); + } finally { + await successorClient.close(); + await fixture.stopHost(successorHost); + } + + const store = createSqliteRuntimeStore(join(fixture.root, 'runtime.sqlite'), { + readOnly: true, + }); + try { + const state = await store.readWorkspaceBoundContinuationClaimStateByBoundary( + admission.execution.boundaryDigest, + ); + assert.equal(state?.claim.protocol, 'continuation_claim_v2'); + assert.equal(state?.startKind, 'runtime_admission'); + assert.ok(state?.startEventId); + assert.equal(state?.claim.workspaceBoundary.commitOid, boundary.commitOid); + assert.equal(state?.claim.workspaceBoundary.treeOid, boundary.treeOid); + assert.equal(state?.claim.workspaceBoundary.revision, boundary.revision); + } finally { + store.close(); + } + }, + ); +}); + +test('an accepted-head continuation never calls the provider twice after Host death', async (t) => { + const helperPath = process.env.MAKA_GITOXIDE_HELPER_PATH; + if (!helperPath) { + t.skip( + 'MAKA_GITOXIDE_HELPER_PATH is required for the real provider crash test', + ); + return; + } + await withManagedContinuationFixture( + helperPath, + async ({ fixture, resourcesRoot, callLog, boundary }) => { + assert.equal(boundary.revision, 1); + const source = await fixture.seedSafeBoundaryContinuationSource(); + const crashHost = await fixture.startHost(undefined, true, { + packagedResourcesRoot: resourcesRoot, + providerCallLogPath: callLog, + providerFailpointAfterSend: true, + }); + const crashClient = await connectClient(fixture.root); + const targetTurnId = 'turn-workspace-bound-provider-crash'; + try { + const initialPlan = await crashClient.request('turn.resume.query', { + sessionId: fixture.sessionId, + }); + assert.equal(initialPlan.disposition, 'ready', JSON.stringify(initialPlan)); + const failpoint = waitForProviderFailpoint(crashHost.child); + const start = crashClient + .request('turn.resume.start', { + sessionId: fixture.sessionId, + turnId: targetTurnId, + sourceRunId: source.sourceRunId, + sourceRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }) + .then( + () => undefined, + () => undefined, + ); + await failpoint; + assert.equal(await providerCallCount(callLog), 1); + await fixture.killHost(crashHost); + await withTimeout(start, PROCESS_TIMEOUT_MS, 'provider-crashed continuation did not close'); + } finally { + await crashClient.close().catch(() => undefined); + } + + const successorHost = await fixture.startHost(undefined, true, { + packagedResourcesRoot: resourcesRoot, + providerCallLogPath: callLog, + }); + const successorClient = await connectClient(fixture.root); + try { + const plan = await successorClient.request('turn.resume.query', { + sessionId: fixture.sessionId, + sourceRunId: source.sourceRunId, + expectedRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }); + assert.equal(plan.disposition, 'parked'); + assert.equal(plan.reason, 'continuation_started_indeterminate'); + const retry = await successorClient.request('turn.resume.start', { + sessionId: fixture.sessionId, + turnId: `${targetTurnId}-retry`, + sourceRunId: source.sourceRunId, + sourceRuntimeEventHighWater: source.sourceRuntimeEventHighWater, + }); + assert.equal(retry.kind, 'parked'); + assert.equal(await providerCallCount(callLog), 1); + } finally { + await successorClient.close(); + await fixture.stopHost(successorHost); + } + }, + ); +}); + +async function withManagedContinuationFixture( + helperInputPath: string, + run: (input: { + fixture: ExecutionFixture; + resourcesRoot: string; + callLog: string; + boundary: NonNullable< + Awaited> + >; + }) => Promise, +): Promise { + const base = await realpath(await mkdtemp(join(tmpdir(), 'maka-gitoxide-continuation-'))); + const root = join(base, 'root'); + const callLog = join(base, 'provider-calls.log'); + await mkdir(root); + await writeFile(callLog, '', 'utf8'); + git(root, ['init', '--quiet', '--object-format=sha1']); + await writeFile(join(root, 'notes.txt'), 'baseline\n', 'utf8'); + git(root, ['add', 'notes.txt']); + git(root, [ + '-c', + 'user.name=Maka Test', + '-c', + 'user.email=maka@example.invalid', + 'commit', + '--quiet', + '-m', + 'baseline', + ]); + + const resourcesRoot = await preparePackagedResources( + base, + helperInputPath, + ); + const capability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const owner = await tryAcquireInteractiveRootOwner(capability); + assert.ok(owner); + if (!owner) throw new Error('Unable to own managed continuation fixture root'); + const stores = await openInteractiveExecutionStoresForWrite(owner.lease); + let sessionId: string; + let boundary: NonNullable< + Awaited> + >; + try { + const session = await stores.sessionStore.create({ + cwd: root, + llmConnectionSlug: 'fake', + model: 'fake-model', + permissionMode: 'ask', + toolProfile: 'managed-coding-v1', + }); + sessionId = session.id; + const helper = await admitRealHelper(helperInputPath); + const sessionInput = { + storageRootLease: owner.lease, + stores, + sourceRoot: root, + sessionId, + ...helper, + }; + await openGitoxideManagedSessionOwnerInternal(sessionInput); + const observedBoundary = + await inspectGitoxideManagedContinuationBoundaryInternal(sessionInput); + assert.ok(observedBoundary); + boundary = observedBoundary; + } finally { + await stores.sessionStore.close?.(); + await owner.close(); + } + + const fixture = new ExecutionFixture(base, root, capability, sessionId); + try { + await run({ fixture, resourcesRoot, callLog, boundary: boundary! }); + } finally { + await fixture.close(); + } +} + +async function preparePackagedResources( + base: string, + helperInputPath: string, +): Promise { + const resourcesRoot = join(base, 'resources'); + const helperDirectory = join(resourcesRoot, 'gitoxide'); + const executableName = + process.platform === 'win32' ? 'maka-gitoxide-helper.exe' : 'maka-gitoxide-helper'; + const executablePath = join(helperDirectory, executableName); + await mkdir(helperDirectory, { recursive: true }); + await copyFile(await realpath(helperInputPath), executablePath); + if (process.platform !== 'win32') await chmod(executablePath, 0o755); + const [bytes, info] = await Promise.all([readFile(executablePath), stat(executablePath)]); + await writeFile( + join(resourcesRoot, 'gitoxide-helper.json'), + `${JSON.stringify({ + schemaVersion: 1, + protocol: 'maka_gitoxide_helper_release_v1', + provider: 'maka/gitoxide-helper', + platform: process.platform, + arch: process.arch, + protocolVersion: 1, + executableRelativePath: `gitoxide/${executableName}`, + bytes: info.size, + sha256: `sha256:${createHash('sha256').update(bytes).digest('hex')}`, + supportedOperations: [ + 'inspect_repository', + 'import_source_head', + 'create_candidate', + 'promote_candidate', + 'observe_accepted_ref', + 'read_tree_file', + ], + distributionReady: true, + })}\n`, + 'utf8', + ); + return resourcesRoot; +} + +async function admitRealHelper(helperInputPath: string) { + const executablePath = await realpath(helperInputPath); + const [bytes, info] = await Promise.all([readFile(executablePath), stat(executablePath)]); + const releaseOwnerToken = {}; + const invocationOwnerToken = {}; + const claim = issueGitoxideHelperReleaseArtifactClaimInternal(releaseOwnerToken, { + executablePath, + expectedSha256: `sha256:${createHash('sha256').update(bytes).digest('hex')}`, + expectedBytes: info.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 }; +} + +function waitForContinuationFailpoint(child: ChildProcess): Promise { + return withTimeout( + new Promise((resolve, reject) => { + const onMessage = (message: unknown): void => { + if ( + message && + typeof message === 'object' && + (message as { type?: unknown }).type === 'test.continuation_failpoint' + ) { + cleanup(); + resolve(); + } + }; + const onExit = (code: number | null, signal: NodeJS.Signals | null): void => { + cleanup(); + reject(new Error(`Runtime Host exited before continuation failpoint: ${code ?? signal}`)); + }; + const cleanup = (): void => { + child.off('message', onMessage); + child.off('exit', onExit); + }; + child.on('message', onMessage); + child.on('exit', onExit); + }), + PROCESS_TIMEOUT_MS * 3, + 'Runtime Host did not reach the continuation failpoint', + ); +} + +function waitForProviderFailpoint(child: ChildProcess): Promise { + return waitForChildMessage(child, 'test.provider_failpoint', 'provider'); +} + +function waitForChildMessage( + child: ChildProcess, + expectedType: string, + label: string, +): Promise { + return withTimeout( + new Promise((resolve, reject) => { + const onMessage = (message: unknown): void => { + if ( + message && + typeof message === 'object' && + (message as { type?: unknown }).type === expectedType + ) { + cleanup(); + resolve(); + } + }; + const onExit = (code: number | null, signal: NodeJS.Signals | null): void => { + cleanup(); + reject(new Error(`Runtime Host exited before ${label} failpoint: ${code ?? signal}`)); + }; + const cleanup = (): void => { + child.off('message', onMessage); + child.off('exit', onExit); + }; + child.on('message', onMessage); + child.on('exit', onExit); + }), + PROCESS_TIMEOUT_MS * 3, + `Runtime Host did not reach the ${label} failpoint`, + ); +} + +async function providerCallCount(path: string): Promise { + return (await readFile(path, 'utf8')).split(/\r?\n/u).filter(Boolean).length; +} + +function git(cwd: string, args: readonly string[]): string { + return execFileSync('git', args, { cwd, encoding: 'utf8' }).trim(); +} diff --git a/packages/runtime-host/src/__tests__/gitoxide-managed-session-owner-internal.test.ts b/packages/runtime-host/src/__tests__/gitoxide-managed-session-owner-internal.test.ts index 7f908852da..43c30c71ea 100644 --- a/packages/runtime-host/src/__tests__/gitoxide-managed-session-owner-internal.test.ts +++ b/packages/runtime-host/src/__tests__/gitoxide-managed-session-owner-internal.test.ts @@ -32,7 +32,10 @@ import { type GitoxideHelperInvocationCapability, issueGitoxideHelperReleaseArtifactClaimInternal, } from '../server/gitoxide-helper-artifact-authority-internal.js'; -import { openGitoxideManagedSessionOwnerInternal } from '../server/gitoxide-managed-session-owner-internal.js'; +import { + inspectGitoxideManagedContinuationBoundaryInternal, + openGitoxideManagedSessionOwnerInternal, +} from '../server/gitoxide-managed-session-owner-internal.js'; test('opens one durable Gitoxide baseline and reuses it for the same session', async (t) => { const helper = await admittedHelper(); @@ -113,6 +116,48 @@ test('fails closed when the source advances after its managed epoch opens', asyn ); }); +test('issues a continuation boundary only for the exact source and accepted Gitoxide head', async (t) => { + const helper = await admittedHelper(); + if (!helper) { + t.skip('MAKA_GITOXIDE_HELPER_PATH is required for the continuation boundary test'); + return; + } + const sourceRoot = await createRepository(t); + await writeFile(join(sourceRoot, 'notes.txt'), 'before\n', 'utf8'); + git(sourceRoot, ['add', 'notes.txt']); + commit(sourceRoot, 'baseline'); + const root = await realpath(await mkdtemp(join(tmpdir(), 'maka-gitoxide-continuation-'))); + t.after(() => rm(root, { recursive: true, force: true })); + const rootCapability = await resolveStorageRoot({ path: root, kind: 'interactive' }); + const rootOwner = await tryAcquireInteractiveRootOwner(rootCapability); + assert.ok(rootOwner); + t.after(() => rootOwner.close()); + const stores = await openInteractiveExecutionStoresForWrite(rootOwner.lease); + t.after(() => stores.sessionStore.close?.()); + const common = { + storageRootLease: rootOwner.lease, + stores, + ...helper, + sourceRoot, + sessionId: 'session-managed-continuation', + } as const; + const opened = await openGitoxideManagedSessionOwnerInternal(common); + const boundary = await inspectGitoxideManagedContinuationBoundaryInternal(common); + assert.ok(boundary); + assert.equal(boundary.repositoryId, opened.repositoryId); + assert.equal(boundary.workspaceEpochId, opened.workspaceEpochId); + assert.equal(boundary.sourceCommitOid, git(sourceRoot, ['rev-parse', 'HEAD'])); + assert.equal(boundary.commitOid, git(opened.repositoryPath, ['rev-parse', 'refs/maka/accepted'])); + + await writeFile(join(sourceRoot, 'notes.txt'), 'source advanced\n', 'utf8'); + git(sourceRoot, ['add', 'notes.txt']); + commit(sourceRoot, 'advance after boundary'); + await assert.rejects( + inspectGitoxideManagedContinuationBoundaryInternal(common), + /source has drifted/i, + ); +}); + test('retries an exact import after the publishing process exits before baseline commit', async (t) => { const helper = await admittedHelper(); if (!helper) { diff --git a/packages/runtime-host/src/server/execution-composition.ts b/packages/runtime-host/src/server/execution-composition.ts index 94757ce4d6..602134f4bf 100644 --- a/packages/runtime-host/src/server/execution-composition.ts +++ b/packages/runtime-host/src/server/execution-composition.ts @@ -41,6 +41,7 @@ import { BackendRegistry, SessionManager, type BackendFactory, + type RuntimeContinuationLifecycleEvent, } from '@maka/runtime/session-manager'; import { buildToolsForAgentDefinition } from '@maka/runtime/agent-catalog'; import { buildHistoryTools } from '@maka/runtime/history-tools'; @@ -152,7 +153,10 @@ import { HostNetworkProxyCoordinator } from './network-proxy-coordinator.js'; import { HostOAuthExecutionAuthority } from './oauth-execution-authority.js'; import { HostOAuthCoordinator, type HostOAuthCoordinatorInput } from './oauth-coordinator.js'; import { HostPlanCoordinator } from './plan-coordinator.js'; -import { openGitoxideManagedSessionOwnerInternal } from './gitoxide-managed-session-owner-internal.js'; +import { + inspectGitoxideManagedContinuationBoundaryInternal, + openGitoxideManagedSessionOwnerInternal, +} from './gitoxide-managed-session-owner-internal.js'; import { PackagedGitoxideHelperError, resolvePackagedGitoxideHelperInternal, @@ -227,6 +231,21 @@ export interface CreateExecutionRuntimeHostCompositionOptions { export interface ExecutionRuntimeHostCompositionDependencies { readonly primaryBackendFactory?: BackendFactory; + /** Production-shaped crash-test seam; product composition never supplies it. */ + readonly continuationFailpoint?: ( + point: + | 'after_continuation_claim_committed' + | 'after_run_created' + | 'after_continuation_start_committed' + | 'after_terminal_event_committed' + | 'after_terminal_header_committed', + ) => Promise; + /** Test/telemetry seam; it observes decisions but owns no durable state. */ + readonly onContinuationLifecycleEvent?: ( + event: RuntimeContinuationLifecycleEvent, + ) => void | Promise; + /** Production-shaped test diagnostic; correctness still fails closed. */ + readonly onContinuationSafetyError?: (error: unknown) => void; readonly oauthAuthorization?: Pick< HostOAuthCoordinatorInput, 'startCodexAuthorization' | 'pollCodexAuthorization' | 'exchangeCodexCode' @@ -1002,6 +1021,8 @@ export async function createExecutionRuntimeHostComposition( newId: randomUUID, now: Date.now, safeBoundaryResumeEnabled: process.env.MAKA_RUNTIME_SAFE_BOUNDARY_RESUME === '1', + continuationFailpoint: dependencies.continuationFailpoint, + onContinuationLifecycleEvent: dependencies.onContinuationLifecycleEvent, inspectContinuationSafety: createLocalContinuationSafetyInspector({ readSessionCwd: async (sessionId) => (await stores.sessionStore.readHeaderSnapshot(sessionId)).cwd, @@ -1027,6 +1048,23 @@ export async function createExecutionRuntimeHostComposition( resourcesLive || graphLive || graphWake.hasLiveSessionState(sessionId) || descendantLive ); }, + readManagedWorkspaceBoundary: async (sessionId) => { + const header = await stores.sessionStore.readHeaderSnapshot(sessionId); + if (header.toolProfile !== 'managed-coding-v1') return undefined; + if (!gitoxideHelperCapability) { + throw new Error( + 'managed_workspace_profile_unavailable: packaged Gitoxide helper authority is unavailable', + ); + } + return inspectGitoxideManagedContinuationBoundaryInternal({ + storageRootLease: context.owner.lease, + stores, + invocationOwnerToken: gitoxideInvocationOwnerToken, + helperCapability: gitoxideHelperCapability, + sourceRoot: header.cwd, + sessionId, + }); + }, }), runBackendActivation: (operation) => runtimePolicyActivation.runBackendActivation(operation), messageAuthority: runtimeAuthority, @@ -2042,6 +2080,20 @@ function requireSessionManager(manager: SessionManager | undefined): SessionMana return manager; } +function observeContinuationSafetyErrors( + inspectContinuationSafety: ReturnType, + onError: ((error: unknown) => void) | undefined, +): ReturnType { + return async (sessionId) => { + try { + return await inspectContinuationSafety(sessionId); + } catch (error) { + onError?.(error); + throw error; + } + }; +} + function requireGraphCoordinator( coordinator: AgentGraphCoordinator | undefined, ): AgentGraphCoordinator { diff --git a/packages/runtime-host/src/server/gitoxide-managed-session-owner-internal.ts b/packages/runtime-host/src/server/gitoxide-managed-session-owner-internal.ts index b78d7c51df..1f69ccfe04 100644 --- a/packages/runtime-host/src/server/gitoxide-managed-session-owner-internal.ts +++ b/packages/runtime-host/src/server/gitoxide-managed-session-owner-internal.ts @@ -21,14 +21,18 @@ import { createHash } from 'node:crypto'; import { mkdir, realpath } from 'node:fs/promises'; import { dirname, join } from 'node:path'; import { MANAGED_MUTATION_EXECUTION_PROFILE_V1_DIGEST } from '@maka/core/runtime-event'; +import type { ManagedWorkspaceContinuationBoundaryV1 } from '@maka/core/runtime-boundary'; import { WORKSPACE_MATERIALIZATION_SEMANTICS_V1, + workspaceMutationPolicyHashV1, type WorkspaceBaselineAuthorityInput, } from '@maka/core/workspace-version-authority'; import type { InteractiveExecutionStoresWriter } from '@maka/storage/execution-stores'; import { issueExecutionStoresWorkspaceBaselineAuthorityInternal, + issueExecutionStoresWorkspaceContinuationAuthorityInternal, requireExecutionStoresWorkspaceBaselineAuthorityInternal, + requireExecutionStoresWorkspaceContinuationAuthorityInternal, } from '@maka/storage/execution-stores-workspace-authority-internal'; import { runWithStorageRootLease, type StorageRootLease } from '@maka/storage/root-authority'; import type { GitoxideHelperInvocationCapability } from './gitoxide-helper-artifact-authority-internal.js'; @@ -39,10 +43,12 @@ import { import { admitGitoxideRepositoryInternal, importAdmittedGitoxideRepositoryInternal, + reopenGitoxideAcceptedRepositoryInternal, requireGitoxideRepositoryAdmissionInternal, } from './gitoxide-repository-admission-authority-internal.js'; const MANAGED_REPOSITORY_DIRECTORY = 'gitoxide-managed-repositories'; +const ACCEPTED_REF = 'refs/maka/accepted'; export interface GitoxideManagedSessionOwnerInternal { readonly repositoryPath: string; @@ -54,6 +60,94 @@ export interface GitoxideManagedSessionOwnerInternal { export type GitoxideManagedSessionOwnerFailpoint = 'after_repository_import'; +/** + * Re-observes an existing managed workspace without gaining mutation authority. + * The returned value is derived from one SQLite read transaction, then bound to + * the exact source HEAD and accepted Gitoxide ref before it reaches Runtime. + */ +export async function inspectGitoxideManagedContinuationBoundaryInternal(input: { + readonly storageRootLease: StorageRootLease<'interactive', 'write'>; + readonly stores: InteractiveExecutionStoresWriter; + readonly invocationOwnerToken: object; + readonly helperCapability: GitoxideHelperInvocationCapability; + readonly sourceRoot: string; + readonly sessionId: string; + readonly abortSignal?: AbortSignal; +}): Promise { + input.abortSignal?.throwIfAborted(); + const [storageRoot, sourceRoot] = await Promise.all([ + runWithStorageRootLease(input.storageRootLease, 'interactive', 'write', async (root) => root), + realpath(input.sourceRoot), + ]); + const identity = deriveManagedSessionIdentity(sourceRoot, input.sessionId); + const continuationOwnerToken = {}; + const capability = issueExecutionStoresWorkspaceContinuationAuthorityInternal({ + ownerToken: continuationOwnerToken, + stores: input.stores, + }); + const authority = requireExecutionStoresWorkspaceContinuationAuthorityInternal( + continuationOwnerToken, + capability, + ); + const boundary = await authority.readContinuationBoundary( + identity.workspaceId, + identity.workspaceEpochId, + MANAGED_MUTATION_EXECUTION_PROFILE_V1_DIGEST, + ); + if (!boundary) return undefined; + if ( + boundary.repositoryId !== identity.repositoryId || + boundary.workspaceId !== identity.workspaceId || + boundary.workspaceEpochId !== identity.workspaceEpochId || + boundary.workspaceInstanceId !== identity.workspaceInstanceId || + boundary.objectFormat !== 'sha1' + ) { + throw new Error('Gitoxide managed continuation boundary conflicts with session identity'); + } + + const sourceOwnerToken = {}; + const sourceAdmission = await admitGitoxideRepositoryInternal({ + invocationOwnerToken: input.invocationOwnerToken, + helperCapability: input.helperCapability, + admissionOwnerToken: sourceOwnerToken, + repositoryPath: sourceRoot, + ...(input.abortSignal ? { abortSignal: input.abortSignal } : {}), + }); + if (sourceAdmission.kind !== 'accepted') { + throw new Error(`Gitoxide managed continuation rejected source: ${sourceAdmission.reason}`); + } + const source = requireGitoxideRepositoryAdmissionInternal( + sourceOwnerToken, + sourceAdmission.capability, + ); + if ( + source.headCommitOid !== boundary.sourceCommitOid || + source.headTreeOid !== boundary.sourceTreeOid + ) { + throw new Error('Gitoxide managed continuation source has drifted'); + } + + const repositoryPath = join( + storageRoot, + MANAGED_REPOSITORY_DIRECTORY, + identity.workspaceEpochId, + 'repository.git', + ); + await reopenGitoxideAcceptedRepositoryInternal({ + invocationOwnerToken: input.invocationOwnerToken, + helperCapability: input.helperCapability, + acceptedRepositoryOwnerToken: {}, + repositoryPath, + acceptedRef: ACCEPTED_REF, + expectedAcceptedCommitOid: boundary.commitOid, + expectedAcceptedTreeOid: boundary.treeOid, + managedTreePolicyVersion: 3, + ...(input.abortSignal ? { abortSignal: input.abortSignal } : {}), + }); + input.abortSignal?.throwIfAborted(); + return boundary; +} + export async function openGitoxideManagedSessionOwnerInternal(input: { readonly storageRootLease: StorageRootLease<'interactive', 'write'>; readonly stores: InteractiveExecutionStoresWriter; @@ -96,8 +190,9 @@ export async function openGitoxideManagedSessionOwnerInternal(input: { const materializationProfileDigest = sha256( `maka-gitoxide-materialization-v3\0${source.helperArtifactSha256}\0`, ); - const policyHash = sha256( - `maka-managed-write-edit-policy-v1\0${materializationProfileDigest}\0${MANAGED_MUTATION_EXECUTION_PROFILE_V1_DIGEST}\0`, + const policyHash = workspaceMutationPolicyHashV1( + materializationProfileDigest, + MANAGED_MUTATION_EXECUTION_PROFILE_V1_DIGEST, ); const baselineOwnerToken = {}; const verifiedBaselines = new WeakMap(); diff --git a/packages/runtime-host/src/server/hosted-execution-projection.ts b/packages/runtime-host/src/server/hosted-execution-projection.ts index b5ab57eb2d..e1f6292d14 100644 --- a/packages/runtime-host/src/server/hosted-execution-projection.ts +++ b/packages/runtime-host/src/server/hosted-execution-projection.ts @@ -73,7 +73,7 @@ export class HostedExecutionProjectionReader { !start || start.claimId !== execution.claimId || start.boundaryDigest !== execution.boundaryDigest || - start.replayManifestDigest !== execution.boundaryDigest || + start.replayManifestDigest !== (execution.replayManifestDigest ?? execution.boundaryDigest) || start.providerReplayDigest !== execution.providerReplayDigest || start.immediateSource.sessionId !== run.sessionId || start.immediateSource.invocationId !== execution.sourceInvocationId || diff --git a/packages/runtime/src/__tests__/runtime-continuation.test.ts b/packages/runtime/src/__tests__/runtime-continuation.test.ts index a9455be4a6..de66d0cc71 100644 --- a/packages/runtime/src/__tests__/runtime-continuation.test.ts +++ b/packages/runtime/src/__tests__/runtime-continuation.test.ts @@ -25,6 +25,7 @@ import { createRuntimeBoundaryCursor, runtimePrefixSegment, type ImmutableRuntimePrefixV1, + type ManagedWorkspaceContinuationBoundaryV1, } from '@maka/core/runtime-boundary'; import type { RuntimeEvent } from '@maka/core/runtime-event'; import type { AgentRunHeader } from '@maka/core/agent-run'; @@ -882,6 +883,28 @@ function runHeader(runId: string, overrides: Partial = {}): Agen }; } +function managedWorkspaceBoundary(): ManagedWorkspaceContinuationBoundaryV1 { + return { + protocol: 'managed_workspace_continuation_boundary_v1', + storageRootId: 'a'.repeat(64), + 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)}`, + acceptedEventId: 'accepted-event-1', + revision: 2, + objectFormat: 'sha1', + sourceCommitOid: '9'.repeat(40), + sourceTreeOid: 'a'.repeat(40), + commitOid: 'b'.repeat(40), + treeOid: 'c'.repeat(40), + materializationProfileDigest: `sha256:${'d'.repeat(64)}`, + policyHash: `sha256:${'e'.repeat(64)}`, + executionProfileDigest: `sha256:${'f'.repeat(64)}`, + }; +} + function event(overrides: Partial): RuntimeEvent { return { id: 'event', diff --git a/packages/runtime/src/__tests__/session-manager.test.ts b/packages/runtime/src/__tests__/session-manager.test.ts index 880a44922e..96c892edd9 100644 --- a/packages/runtime/src/__tests__/session-manager.test.ts +++ b/packages/runtime/src/__tests__/session-manager.test.ts @@ -4654,6 +4654,64 @@ describe('SessionManager permission mode updates', () => { expect(plan.rejectionReasons).toEqual(['safety_observation_unavailable']); }); + test('never downgrades a managed session when its workspace boundary is unavailable', async () => { + const store = new MemorySessionStore(); + const runStore = new MemoryAgentRunStore(); + const manager = new SessionManager({ + store, + runStore, + runtimeEventStore: runStore, + backends: new BackendRegistry(), + safeBoundaryResumeEnabled: true, + inspectContinuationSafety: async () => ({ + workspaceIdentity: 'workspace-managed', + backgroundOperationsSettled: true, + availableToolNames: ['Write', 'Edit'], + }), + newId: nextId(), + now: nextNow(6_533), + }); + const session = await manager.createSession(makeInput({ toolProfile: 'managed-coding-v1' })); + const header = await store.readHeader(session.id); + expect(header.toolProfile).toBe('managed-coding-v1'); + const sourceRunId = 'source-run-managed-boundary-missing'; + const sourceTurnId = 'source-turn-managed-boundary-missing'; + await seedRuntimeRun( + runStore, + makeRunHeader({ + runId: sourceRunId, + sessionId: session.id, + turnId: sourceTurnId, + status: 'failed', + cwd: header.cwd, + workspaceIdentity: 'workspace-managed', + createdAt: 1, + updatedAt: 2, + completedAt: 2, + failureClass: 'app_restarted', + }), + [ + runtimeEvent({ + id: 'source-terminal-managed-boundary-missing', + invocationId: 'source-invocation-managed-boundary-missing', + runId: sourceRunId, + sessionId: session.id, + turnId: sourceTurnId, + ts: 2, + status: 'failed', + actions: { endInvocation: true, stateDelta: { failureClass: 'app_restarted' } }, + }), + ], + ); + + const plan = await manager.planAuthoritativeSafeBoundaryContinuation(session.id, { + sourceRunId, + }); + + expect(plan.disposition).toBe('park'); + expect(plan.rejectionReasons).toEqual(['workspace_boundary_unavailable']); + }); + test('keeps the authoritative continuation entry disabled unless the host enables it', async () => { const store = new MemorySessionStore(); const backends = new BackendRegistry(); @@ -13890,6 +13948,7 @@ class MemorySessionStore implements SessionStore { permissionMode: input.permissionMode, collaborationMode: input.collaborationMode ?? 'agent', orchestrationMode: input.orchestrationMode ?? 'default', + ...(input.toolProfile !== undefined ? { toolProfile: input.toolProfile } : {}), schemaVersion: 1, }; this.headers.set(header.id, header); @@ -14426,6 +14485,9 @@ class MemoryAgentRunStore async claimContinuation(input: { claim: ContinuationClaimV1 }) { const claim = decodeContinuationClaim(input.claim); + if (claim.protocol !== 'continuation_claim_v1') { + throw new Error('Memory continuation authority only supports v1 claims'); + } const existing = this.continuationClaims.get(claim.boundaryDigest); if (existing) return { kind: 'existing' as const, claim: existing }; const conflict = [...this.continuationClaims.values()].find( diff --git a/packages/runtime/src/agent-run.ts b/packages/runtime/src/agent-run.ts index ba68ad62c5..0f6fbe5a2a 100644 --- a/packages/runtime/src/agent-run.ts +++ b/packages/runtime/src/agent-run.ts @@ -25,6 +25,7 @@ import type { } from '@maka/core/agent-run'; import type { RuntimeEvent, ToolBoundaryProtocol } from '@maka/core/runtime-event'; import type { RuntimeEventStore } from '@maka/core/runtime-event-store'; +import { digestWorkspaceBoundContinuationBoundary } from '@maka/core/runtime-boundary'; import type { RunCompositionSnapshot } from '@maka/core/run-composition'; import { decodeRunCompositionSnapshot } from '@maka/core/run-composition'; import { DurableStoreWriteError, RunSealedError } from '@maka/core/runtime-event-store'; @@ -1117,9 +1118,16 @@ export class AgentRun { continuationSource: continuation.claimId && continuation.boundary ? { - protocol: 'continuation_source_v2' as const, + protocol: continuation.workspaceBoundary + ? ('continuation_source_v3' as const) + : ('continuation_source_v2' as const), claimId: continuation.claimId, - boundaryDigest: continuation.boundary.manifestDigest, + boundaryDigest: continuation.workspaceBoundary + ? digestWorkspaceBoundContinuationBoundary( + continuation.boundary, + continuation.workspaceBoundary, + ) + : continuation.boundary.manifestDigest, sourceInvocationId: continuation.sourceInvocationId, sourceRunId: continuation.sourceRunId, sourceTurnId: continuation.sourceTurnId, diff --git a/packages/runtime/src/continuation-safety.ts b/packages/runtime/src/continuation-safety.ts index 1e5b06934c..26f773df81 100644 --- a/packages/runtime/src/continuation-safety.ts +++ b/packages/runtime/src/continuation-safety.ts @@ -32,6 +32,9 @@ export interface LocalContinuationSafetyInspectorDeps { readWorkspaceCheckpoint?: ( sessionId: string, ) => Promise; + readManagedWorkspaceBoundary?: ( + sessionId: string, + ) => Promise; } export function createLocalContinuationSafetyInspector( @@ -39,19 +42,26 @@ export function createLocalContinuationSafetyInspector( ): (sessionId: string) => Promise { return async (sessionId) => { const cwd = await deps.readSessionCwd(sessionId); - const [workspace, availableToolNames, hasPendingBackgroundOperations, workspaceCheckpoint] = - await Promise.all([ - deps.resolveWorkspaceIdentity(cwd), - deps.listAvailableToolNames(sessionId), - deps.hasPendingBackgroundOperations(sessionId), - deps.readWorkspaceCheckpoint?.(sessionId), - ]); + const [ + workspace, + availableToolNames, + hasPendingBackgroundOperations, + workspaceCheckpoint, + workspaceBoundary, + ] = await Promise.all([ + deps.resolveWorkspaceIdentity(cwd), + deps.listAvailableToolNames(sessionId), + deps.hasPendingBackgroundOperations(sessionId), + deps.readWorkspaceCheckpoint?.(sessionId), + deps.readManagedWorkspaceBoundary?.(sessionId), + ]); return { workspaceIdentity: workspace.workspaceIdentity, workspacePath: workspace.canonicalPath, backgroundOperationsSettled: !hasPendingBackgroundOperations, availableToolNames: [...new Set(availableToolNames)].sort(), ...(workspaceCheckpoint ? { workspaceCheckpoint } : {}), + ...(workspaceBoundary ? { workspaceBoundary } : {}), }; }; } diff --git a/packages/runtime/src/runtime-kernel.ts b/packages/runtime/src/runtime-kernel.ts index 1463202268..08a35e0169 100644 --- a/packages/runtime/src/runtime-kernel.ts +++ b/packages/runtime/src/runtime-kernel.ts @@ -20,7 +20,10 @@ import type { AgentRunHeader, AgentRunStore } from '@maka/core/agent-run'; import { decodeRuntimeBoundaryCursor, + digestWorkspaceBoundContinuationBoundary, + type ContinuationClaim, type ContinuationClaimV1, + type ContinuationClaimV2, type ImmutableRuntimePrefixV1, } from '@maka/core/runtime-boundary'; import { @@ -31,6 +34,7 @@ import { import type { RuntimeContinuationAuthorityStore, RuntimeEventStore, + RuntimeWorkspaceBoundContinuationAuthorityStore, } from '@maka/core/runtime-event-store'; import { type ActiveInteractionRequestEvent, @@ -758,7 +762,14 @@ export class RuntimeKernel implements RuntimeKernelLike { claimedAt, }); const claim = continuationClaimForExecution(continuation, claimedAt, targetRunHeader); - const claimResult = await continuationAuthority.claimContinuation({ claim }); + const workspaceContinuationAuthority = + claim.protocol === 'continuation_claim_v2' + ? requireRuntimeWorkspaceBoundContinuationAuthority(this.deps.runtimeEventStore) + : undefined; + const claimResult = + claim.protocol === 'continuation_claim_v2' + ? await workspaceContinuationAuthority!.claimWorkspaceBoundContinuation({ claim }) + : await continuationAuthority.claimContinuation({ claim }); if (claimResult.kind !== 'acquired') { throw new RuntimeContinuationRevalidationError( 'continuation_claim_conflict', @@ -813,7 +824,9 @@ export class RuntimeKernel implements RuntimeKernelLike { commitContinuationStart: async (startedAt) => { const source = claim.boundary.segments.at(-1)!; const eventId = this.deps.newId(); - const result = await continuationAuthority.commitContinuationStart({ + const result = await commitContinuationStartForClaim({ + continuationAuthority, + workspaceContinuationAuthority, claim, event: { id: eventId, @@ -2690,6 +2703,45 @@ function requireRuntimeContinuationAuthority( return candidate as RuntimeContinuationAuthorityStore; } +function requireRuntimeWorkspaceBoundContinuationAuthority( + store: RuntimeEventStore, +): RuntimeWorkspaceBoundContinuationAuthorityStore { + const candidate = store as Partial; + if ( + candidate.workspaceBoundContinuationAuthorityCapability !== + 'runtime_workspace_bound_continuation_authority_v1' || + typeof candidate.claimWorkspaceBoundContinuation !== 'function' || + typeof candidate.readWorkspaceBoundContinuationClaimStateByBoundary !== 'function' || + typeof candidate.listWorkspaceBoundContinuationClaimsForRecovery !== 'function' || + typeof candidate.commitWorkspaceBoundContinuationStart !== 'function' || + typeof candidate.commitWorkspaceBoundContinuationRepairStart !== 'function' + ) { + throw new Error('Managed Runtime continuation requires workspace-bound SQLite authority'); + } + return candidate as RuntimeWorkspaceBoundContinuationAuthorityStore; +} + +async function commitContinuationStartForClaim(input: { + continuationAuthority: RuntimeContinuationAuthorityStore; + workspaceContinuationAuthority?: RuntimeWorkspaceBoundContinuationAuthorityStore; + claim: ContinuationClaim; + event: RuntimeEvent; +}): Promise<{ created: boolean; runtimeEventSeq: number }> { + if (input.claim.protocol === 'continuation_claim_v2') { + if (!input.workspaceContinuationAuthority) { + throw new Error('Workspace-bound continuation start authority is unavailable'); + } + return input.workspaceContinuationAuthority.commitWorkspaceBoundContinuationStart({ + claim: input.claim, + event: input.event, + }); + } + return input.continuationAuthority.commitContinuationStart({ + claim: input.claim, + event: input.event, + }); +} + async function revalidateContinuationBoundary( store: RuntimeContinuationAuthorityStore, continuation: RuntimeContinuation, @@ -2749,7 +2801,7 @@ function continuationClaimForExecution( continuation: RuntimeContinuation, claimedAt: number, targetRunHeader: AgentRunHeader, -): ContinuationClaimV1 { +): ContinuationClaim { if ( !continuation.claimId || !continuation.boundary || @@ -2761,10 +2813,14 @@ function continuationClaimForExecution( 'Runtime continuation is missing its durable claim identity', ); } - return { - protocol: 'continuation_claim_v1', + const common = { claimId: continuation.claimId, - boundaryDigest: continuation.boundary.manifestDigest, + boundaryDigest: continuation.workspaceBoundary + ? digestWorkspaceBoundContinuationBoundary( + continuation.boundary, + continuation.workspaceBoundary, + ) + : continuation.boundary.manifestDigest, boundary: continuation.boundary, providerProjectionVersion: continuation.providerProjectionVersion, providerReplayDigest: continuation.providerReplayDigest, @@ -2777,6 +2833,17 @@ function continuationClaimForExecution( targetRunHeader, claimedAt, }; + if (continuation.workspaceBoundary) { + return { + protocol: 'continuation_claim_v2', + ...common, + workspaceBoundary: continuation.workspaceBoundary, + } satisfies ContinuationClaimV2; + } + return { + protocol: 'continuation_claim_v1', + ...common, + } satisfies ContinuationClaimV1; } function continuationTargetRunHeaderForExecution(input: { @@ -2803,6 +2870,12 @@ function continuationTargetRunHeaderForExecution(input: { ); } const source = continuation.boundary.segments.at(-1)!; + const boundaryDigest = continuation.workspaceBoundary + ? digestWorkspaceBoundContinuationBoundary( + continuation.boundary, + continuation.workspaceBoundary, + ) + : continuation.boundary.manifestDigest; return { runId: continuation.runId, invocationId: continuation.invocationId, @@ -2836,9 +2909,11 @@ function continuationTargetRunHeaderForExecution(input: { ...(userInput.agentId ? { agentId: userInput.agentId } : {}), ...(userInput.agentName ? { agentName: userInput.agentName } : {}), continuationSource: { - protocol: 'continuation_source_v2', + protocol: continuation.workspaceBoundary + ? 'continuation_source_v3' + : 'continuation_source_v2', claimId: continuation.claimId, - boundaryDigest: continuation.boundary.manifestDigest, + boundaryDigest, sourceInvocationId: source.identity.invocationId, sourceRunId: source.identity.runId, sourceTurnId: source.identity.turnId, @@ -3020,6 +3095,14 @@ function assertContinuationSafetyUnchanged( 'Runtime continuation workspace identity changed after planning', ); } + if ( + !isDeepStrictEqual(observation.workspaceBoundary ?? null, snapshot.workspaceBoundary ?? null) + ) { + throw new RuntimeContinuationRevalidationError( + 'workspace_version_changed', + 'Runtime continuation accepted workspace boundary changed after planning', + ); + } if (!observation.backgroundOperationsSettled) { throw new RuntimeContinuationRevalidationError( 'background_operation_started', diff --git a/packages/runtime/src/runtime-resume.ts b/packages/runtime/src/runtime-resume.ts index 6e861da307..80f4106c25 100644 --- a/packages/runtime/src/runtime-resume.ts +++ b/packages/runtime/src/runtime-resume.ts @@ -27,13 +27,18 @@ import { type RuntimeEventFunctionResponseContent, } from '@maka/core/runtime-event'; import type { - ContinuationClaimV1, + ContinuationClaim, ImmutableRuntimePrefixV1, + ManagedWorkspaceContinuationBoundaryV1, RuntimeBoundaryCursorV1, RuntimeBoundaryDigest, } from '@maka/core/runtime-boundary'; +import { digestWorkspaceBoundContinuationBoundary } from '@maka/core/runtime-boundary'; import type { AgentRunHeader } from '@maka/core/agent-run'; -import type { ContinuationClaimStateV1 } from '@maka/core/runtime-event-store'; +import type { + ContinuationClaimStateV1, + ContinuationClaimStateV2, +} from '@maka/core/runtime-event-store'; import { isDeepStrictEqual } from 'node:util'; import { buildContinuationReplayPlan, @@ -77,6 +82,7 @@ export type RuntimeContinuationRevalidationCode = | 'source_ledger_identity_changed' | 'source_replay_changed' | 'workspace_identity_changed' + | 'workspace_version_changed' | 'background_operation_started' | 'tool_catalog_changed' | 'workspace_checkpoint_changed'; @@ -123,6 +129,7 @@ export type ResumePlanDiagnosticCode = | 'source_run_unreadable' | 'continuation_already_exists' | 'continuation_authority_unavailable' + | 'workspace_boundary_unavailable' | 'continuation_claim_repair_required' | 'continuation_started_indeterminate' | 'workspace_identity_missing' @@ -166,6 +173,7 @@ export type ResumeRejectionReason = | 'source_run_unreadable' | 'continuation_already_exists' | 'continuation_authority_unavailable' + | 'workspace_boundary_unavailable' | 'continuation_claim_repair_required' | 'continuation_started_indeterminate' | 'workspace_identity_missing' @@ -300,6 +308,7 @@ export interface SafeBoundaryContinuationFacts { restored: boolean; runtimeEventHighWater: number; }; + workspaceBoundary?: ManagedWorkspaceContinuationBoundaryV1; } export interface RuntimeContinuation { @@ -322,6 +331,7 @@ export interface RuntimeContinuation { /** Identity of the exact provider-facing replay projection. */ providerReplayDigest?: RuntimeBoundaryDigest; providerProjectionVersion?: typeof PROVIDER_REPLAY_PROJECTION_VERSION; + workspaceBoundary?: ManagedWorkspaceContinuationBoundaryV1; safetySnapshot: RuntimeContinuationSafetySnapshot; } @@ -333,6 +343,7 @@ export interface RuntimeContinuationSafetySnapshot { ref: string; runtimeEventHighWater: number; }; + workspaceBoundary?: ManagedWorkspaceContinuationBoundaryV1; } export interface RuntimeContinuationSafetyObservation { @@ -346,6 +357,7 @@ export interface RuntimeContinuationSafetyObservation { restored: boolean; runtimeEventHighWater: number; }; + workspaceBoundary?: ManagedWorkspaceContinuationBoundaryV1; } export interface SafeBoundaryContinuationPlan { @@ -365,6 +377,9 @@ export interface RuntimeContinuationPlannerInput { availableToolNames: readonly string[]; expectedRuntimeEventHighWater?: number; workspaceCheckpoint?: SafeBoundaryContinuationFacts['workspaceCheckpoint']; + workspaceBoundary?: ManagedWorkspaceContinuationBoundaryV1; + /** Managed coding requires an authenticated workspace boundary and may never downgrade to v1. */ + workspaceBoundaryRequirement?: 'optional' | 'required'; } export interface RuntimeContinuationPlannerDeps { @@ -377,6 +392,9 @@ export interface RuntimeContinuationPlannerDeps { readContinuationClaimStateByBoundary?( boundaryDigest: RuntimeBoundaryDigest, ): Promise; + readWorkspaceBoundContinuationClaimStateByBoundary?( + boundaryDigest: RuntimeBoundaryDigest, + ): Promise; findExistingContinuation?( sessionId: string, sourceRunId: string, @@ -389,6 +407,12 @@ export class RuntimeContinuationPlanner { constructor(private readonly deps: RuntimeContinuationPlannerDeps) {} async plan(input: RuntimeContinuationPlannerInput): Promise { + if (input.workspaceBoundaryRequirement === 'required' && !input.workspaceBoundary) { + return parkedPlan( + 'workspace_boundary_unavailable', + 'managed continuation requires an authoritative workspace boundary', + ); + } let sourceRun: Awaited>; try { sourceRun = await this.deps.readSourceRun(input.sessionId, input.sourceRunId); @@ -435,11 +459,22 @@ export class RuntimeContinuationPlanner { `continuation replay segment ${replay.segmentIndex} is not replayable: ${replay.reason}`, ); } - let durableClaimState: ContinuationClaimStateV1 | undefined; - try { - durableClaimState = await this.deps.readContinuationClaimStateByBoundary?.( - replay.plan.boundary.manifestDigest, + const plannedBoundaryDigest = input.workspaceBoundary + ? digestWorkspaceBoundContinuationBoundary(replay.plan.boundary, input.workspaceBoundary) + : replay.plan.boundary.manifestDigest; + if (input.workspaceBoundary && !this.deps.readWorkspaceBoundContinuationClaimStateByBoundary) { + return parkedPlan( + 'continuation_authority_unavailable', + 'workspace-bound continuation authority is unavailable', ); + } + let durableClaimState: ContinuationClaimStateV1 | ContinuationClaimStateV2 | undefined; + try { + durableClaimState = input.workspaceBoundary + ? await this.deps.readWorkspaceBoundContinuationClaimStateByBoundary?.( + plannedBoundaryDigest, + ) + : await this.deps.readContinuationClaimStateByBoundary?.(plannedBoundaryDigest); } catch { return parkedPlan( 'continuation_authority_unavailable', @@ -449,8 +484,11 @@ export class RuntimeContinuationPlanner { if (durableClaimState) { const claim = durableClaimState.claim; if ( - claim.boundaryDigest !== replay.plan.boundary.manifestDigest || + claim.boundaryDigest !== plannedBoundaryDigest || !isDeepStrictEqual(claim.boundary, replay.plan.boundary) || + (claim.protocol === 'continuation_claim_v2' && + !isDeepStrictEqual(claim.workspaceBoundary, input.workspaceBoundary)) || + (claim.protocol === 'continuation_claim_v1' && input.workspaceBoundary !== undefined) || claim.providerProjectionVersion !== replay.plan.providerProjectionVersion || claim.providerReplayDigest !== replay.plan.providerReplayDigest ) { @@ -500,12 +538,15 @@ export class RuntimeContinuationPlanner { ...(input.workspaceCheckpoint !== undefined ? { workspaceCheckpoint: input.workspaceCheckpoint } : {}), + ...(input.workspaceBoundary !== undefined + ? { workspaceBoundary: input.workspaceBoundary } + : {}), }); } private async classifyExistingClaim( sessionId: string, - state: ContinuationClaimStateV1, + state: ContinuationClaimStateV1 | ContinuationClaimStateV2, ): Promise { const { claim } = state; const detail = { @@ -583,6 +624,17 @@ export class RuntimeContinuationPlanner { isTerminalRunStatus(targetRun.status) && terminalRunHeaderMatchesFact(targetRun, terminalClassification.fact) ) { + if ( + state.startKind === 'runtime_admission' && + terminalClassification.fact.runStatus === 'failed' && + terminalClassification.fact.failureClass === 'app_restarted' + ) { + return parkedPlan( + 'continuation_started_indeterminate', + 'continuation-start is durable and Host restart closure does not prove provider absence', + detail, + ); + } return parkedPlan( 'continuation_already_exists', 'source boundary already has a terminal continuation', @@ -621,13 +673,15 @@ export class RuntimeContinuationPlanner { }); const segments: ImmutableRuntimePrefixV1[] = [immediate]; const seen = new Set([sourceRunId]); - const v2Edges: Array<{ + const canonicalEdges: Array<{ childRunId: string; childRunHeader: AgentRunHeader; startEvent: RuntimeEvent; startKind: 'runtime_admission' | 'claim_repair'; + protocol: 'continuation_source_v2' | 'continuation_source_v3'; claimId: string; boundaryDigest: RuntimeBoundaryDigest; + replayManifestDigest: RuntimeBoundaryDigest; providerProjectionVersion: typeof PROVIDER_REPLAY_PROJECTION_VERSION; providerReplayDigest: RuntimeBoundaryDigest; }> = []; @@ -638,41 +692,46 @@ export class RuntimeContinuationPlanner { while (true) { const current = childRun.continuationSource; const start = childPrefix.events[0]?.actions?.continuationStart; - const currentV2 = - current && 'protocol' in current && current.protocol === 'continuation_source_v2' + const canonicalSource = + current && + 'protocol' in current && + (current.protocol === 'continuation_source_v2' || + current.protocol === 'continuation_source_v3') ? current : undefined; - if (start && !currentV2) { + if (start && !canonicalSource) { throw new RuntimeLineageError( 'runtime_lineage_start_mismatch', `canonical continuation-start cannot be downgraded to legacy lineage for ${childRunId}`, ); } - if (currentV2) { + if (canonicalSource) { if ( !start || - start.claimId !== currentV2.claimId || - start.boundaryDigest !== currentV2.boundaryDigest || - start.replayManifestDigest !== currentV2.replayManifestDigest || + start.claimId !== canonicalSource.claimId || + start.boundaryDigest !== canonicalSource.boundaryDigest || + start.replayManifestDigest !== canonicalSource.replayManifestDigest || start.immediateSource.sessionId !== sessionId || - start.immediateSource.invocationId !== currentV2.sourceInvocationId || - start.immediateSource.runId !== currentV2.sourceRunId || - start.immediateSource.turnId !== currentV2.sourceTurnId || - start.immediateSource.highWater !== currentV2.sourceRuntimeEventHighWater || - start.immediateSource.prefixDigest !== currentV2.sourcePrefixDigest + start.immediateSource.invocationId !== canonicalSource.sourceInvocationId || + start.immediateSource.runId !== canonicalSource.sourceRunId || + start.immediateSource.turnId !== canonicalSource.sourceTurnId || + start.immediateSource.highWater !== canonicalSource.sourceRuntimeEventHighWater || + start.immediateSource.prefixDigest !== canonicalSource.sourcePrefixDigest ) { throw new RuntimeLineageError( 'runtime_lineage_start_mismatch', `continuation-start does not authenticate lineage edge for ${childRunId}`, ); } - v2Edges.push({ + canonicalEdges.push({ childRunId, childRunHeader: childRun, startEvent: childPrefix.events[0]!, startKind: start.provenance, + protocol: canonicalSource.protocol, claimId: start.claimId, - boundaryDigest: currentV2.boundaryDigest, + boundaryDigest: canonicalSource.boundaryDigest, + replayManifestDigest: canonicalSource.replayManifestDigest, providerProjectionVersion: start.providerProjectionVersion, providerReplayDigest: start.providerReplayDigest, }); @@ -722,7 +781,8 @@ export class RuntimeContinuationPlanner { } if ( 'protocol' in current && - current.protocol === 'continuation_source_v2' && + (current.protocol === 'continuation_source_v2' || + current.protocol === 'continuation_source_v3') && current.sourcePrefixDigest !== prefix.prefixDigest ) { throw new RuntimeLineageError( @@ -736,7 +796,7 @@ export class RuntimeContinuationPlanner { childRunId = current.sourceRunId; depth += 1; } - for (const edge of v2Edges) { + for (const edge of canonicalEdges) { const childIndex = segments.findIndex((prefix) => prefix.identity.runId === edge.childRunId); if (childIndex <= 0) { throw new RuntimeLineageError( @@ -753,7 +813,7 @@ export class RuntimeContinuationPlanner { }); if ( edgeReplay.kind !== 'replayable' || - edgeReplay.plan.boundary.manifestDigest !== edge.boundaryDigest || + edgeReplay.plan.boundary.manifestDigest !== edge.replayManifestDigest || edgeReplay.plan.providerReplayDigest !== edge.providerReplayDigest ) { throw new RuntimeLineageError( @@ -761,15 +821,19 @@ export class RuntimeContinuationPlanner { `continuation provider replay changed before ${edge.childRunId}`, ); } - if (!this.deps.readContinuationClaimStateByBoundary) { + const readClaimState = + edge.protocol === 'continuation_source_v3' + ? this.deps.readWorkspaceBoundContinuationClaimStateByBoundary + : this.deps.readContinuationClaimStateByBoundary; + if (!readClaimState) { throw new RuntimeLineageError( 'continuation_authority_unavailable', `durable continuation authority is unavailable for ${edge.childRunId}`, ); } - let state: ContinuationClaimStateV1 | undefined; + let state: ContinuationClaimStateV1 | ContinuationClaimStateV2 | undefined; try { - state = await this.deps.readContinuationClaimStateByBoundary(edge.boundaryDigest); + state = await readClaimState(edge.boundaryDigest); } catch { throw new RuntimeLineageError( 'continuation_authority_unavailable', @@ -1169,6 +1233,7 @@ export function buildSafeBoundaryContinuationPlan( : {}), runtimeContext: [...modelRuntimeContext], ...(facts.continuationClaimId ? { claimId: facts.continuationClaimId } : {}), + ...(facts.workspaceBoundary ? { workspaceBoundary: facts.workspaceBoundary } : {}), ...(compositeReplay ? { boundary: compositeReplay.boundary, @@ -1188,6 +1253,7 @@ export function buildSafeBoundaryContinuationPlan( }, } : {}), + ...(facts.workspaceBoundary ? { workspaceBoundary: facts.workspaceBoundary } : {}), }, }, }; @@ -1466,7 +1532,7 @@ function hasMatchingCall( return call !== undefined && call.name === response.name; } -function claimTargetRunHeaderMatches(actual: AgentRunHeader, claim: ContinuationClaimV1): boolean { +function claimTargetRunHeaderMatches(actual: AgentRunHeader, claim: ContinuationClaim): boolean { const candidate = actual as unknown as Record; const expected = claim.targetRunHeader as unknown as Record; const immutable = (header: Record) => { @@ -1487,7 +1553,7 @@ function claimTargetRunHeaderMatches(actual: AgentRunHeader, claim: Continuation function continuationStartMatchesClaim( event: RuntimeEvent | undefined, - claim: ContinuationClaimV1, + claim: ContinuationClaim, startKind: ContinuationClaimStateV1['startKind'], ): boolean { const start = event?.actions?.continuationStart; diff --git a/packages/runtime/src/session-manager.ts b/packages/runtime/src/session-manager.ts index b284c62328..d3a4f8144e 100644 --- a/packages/runtime/src/session-manager.ts +++ b/packages/runtime/src/session-manager.ts @@ -122,10 +122,13 @@ import type { RootExecutionDescriptor, } from '@maka/core/agent-run'; import type { ArtifactRecord } from '@maka/core/artifacts'; -import type { ContinuationClaimV1 } from '@maka/core/runtime-boundary'; +import type { ContinuationClaim } from '@maka/core/runtime-boundary'; import type { + ContinuationClaimStateV1, + ContinuationClaimStateV2, RuntimeEventStore, RuntimeContinuationAuthorityStore, + RuntimeWorkspaceBoundContinuationAuthorityStore, } from '@maka/core/runtime-event-store'; import type { RuntimeEvent, ToolBoundaryProtocol } from '@maka/core/runtime-event'; import type { RunCompositionSnapshot } from '@maka/core/run-composition'; @@ -249,6 +252,21 @@ function runtimeContinuationAuthority( : undefined; } +function runtimeWorkspaceBoundContinuationAuthority( + store: RuntimeEventStore | undefined, +): RuntimeWorkspaceBoundContinuationAuthorityStore | undefined { + const candidate = store as Partial | undefined; + return candidate?.workspaceBoundContinuationAuthorityCapability === + 'runtime_workspace_bound_continuation_authority_v1' && + typeof candidate.claimWorkspaceBoundContinuation === 'function' && + typeof candidate.readWorkspaceBoundContinuationClaimStateByBoundary === 'function' && + typeof candidate.listWorkspaceBoundContinuationClaimsForRecovery === 'function' && + typeof candidate.commitWorkspaceBoundContinuationStart === 'function' && + typeof candidate.commitWorkspaceBoundContinuationRepairStart === 'function' + ? (candidate as RuntimeWorkspaceBoundContinuationAuthorityStore) + : undefined; +} + function runtimeCommitSinkFromEventStore( store: RuntimeEventStore | undefined, ): RuntimeCommitSink | undefined { @@ -1431,6 +1449,7 @@ export class SessionManager { continuationClaimRecovered = await this.recoverContinuationClaimsBeforeProvider( session.id, continuationAuthority, + runtimeWorkspaceBoundContinuationAuthority(this.deps.runtimeEventStore), policy, ); } catch (error) { @@ -2013,6 +2032,13 @@ export class SessionManager { if (!authority) throw new Error('Continuation authority is not configured'); return authority.readContinuationClaimStateByBoundary(boundaryDigest); }, + readWorkspaceBoundContinuationClaimStateByBoundary: async (boundaryDigest) => { + const authority = runtimeWorkspaceBoundContinuationAuthority(this.deps.runtimeEventStore); + if (!authority) { + throw new Error('Workspace-bound continuation authority is not configured'); + } + return authority.readWorkspaceBoundContinuationClaimStateByBoundary(boundaryDigest); + }, findExistingContinuation: async ( targetSessionId, sourceRunId, @@ -2108,12 +2134,17 @@ export class SessionManager { currentWorkspaceIdentity: observation.workspaceIdentity, backgroundOperationsSettled: observation.backgroundOperationsSettled, availableToolNames: observation.availableToolNames, + workspaceBoundaryRequirement: + header.toolProfile === 'managed-coding-v1' ? 'required' : 'optional', ...(input.expectedRuntimeEventHighWater !== undefined ? { expectedRuntimeEventHighWater: input.expectedRuntimeEventHighWater } : {}), ...(observation.workspaceCheckpoint ? { workspaceCheckpoint: observation.workspaceCheckpoint } : {}), + ...(observation.workspaceBoundary + ? { workspaceBoundary: observation.workspaceBoundary } + : {}), }); } @@ -3641,9 +3672,17 @@ export class SessionManager { const continuationAuthority = runtimeContinuationAuthority(this.deps.runtimeEventStore); if (continuationAuthority) { - const claimState = ( - await continuationAuthority.listContinuationClaimsForRecovery(input.sessionId) - ).find( + const workspaceAuthority = runtimeWorkspaceBoundContinuationAuthority( + this.deps.runtimeEventStore, + ); + const claimState = [ + ...(await continuationAuthority.listContinuationClaimsForRecovery(input.sessionId)), + ...(workspaceAuthority + ? await workspaceAuthority.listWorkspaceBoundContinuationClaimsForRecovery( + input.sessionId, + ) + : []), + ].find( (candidate) => candidate.claim.target.runId === input.runId || candidate.claim.target.turnId === input.turnId, @@ -4535,13 +4574,30 @@ export class SessionManager { private async recoverContinuationClaimsBeforeProvider( sessionId: string, authority: RuntimeContinuationAuthorityStore, + workspaceAuthority: RuntimeWorkspaceBoundContinuationAuthorityStore | undefined, _policy: RecoveryPolicy, ): Promise { if (!this.deps.runStore) return false; - const states = await authority.listContinuationClaimsForRecovery(sessionId); + const states: Array = [ + ...(await authority.listContinuationClaimsForRecovery(sessionId)), + ...(workspaceAuthority + ? await workspaceAuthority.listWorkspaceBoundContinuationClaimsForRecovery(sessionId) + : []), + ]; let recovered = false; for (const initialState of states) { const { claim } = initialState; + if (claim.protocol === 'continuation_claim_v2') { + if (!workspaceAuthority || !this.deps.inspectContinuationSafety) { + throw new Error('Workspace-bound continuation recovery authority is unavailable'); + } + const observation = await this.deps.inspectContinuationSafety(sessionId); + if (!isDeepStrictEqual(observation.workspaceBoundary, claim.workspaceBoundary)) { + throw new Error( + `Workspace-bound continuation claim ${claim.claimId} no longer matches accepted head`, + ); + } + } let run: AgentRunHeader; try { run = await this.deps.runStore.readRun(sessionId, claim.target.runId); @@ -4559,9 +4615,13 @@ export class SessionManager { } recovered = true; } - let state = - (await authority.readContinuationClaimStateByBoundary(claim.boundaryDigest)) ?? - initialState; + const readState = () => + claim.protocol === 'continuation_claim_v2' + ? workspaceAuthority!.readWorkspaceBoundContinuationClaimStateByBoundary( + claim.boundaryDigest, + ) + : authority.readContinuationClaimStateByBoundary(claim.boundaryDigest); + let state = (await readState()) ?? initialState; if ( state.startEventId ? !claimTargetRunHeaderIsCompatible(run, claim.targetRunHeader) @@ -4584,9 +4644,16 @@ export class SessionManager { ); } const repairStart = buildContinuationRepairStartEvent(claim); - await authority.commitContinuationRepairStart({ claim, event: repairStart }); + if (claim.protocol === 'continuation_claim_v2') { + await workspaceAuthority!.commitWorkspaceBoundContinuationRepairStart({ + claim, + event: repairStart, + }); + } else { + await authority.commitContinuationRepairStart({ claim, event: repairStart }); + } state = - (await authority.readContinuationClaimStateByBoundary(claim.boundaryDigest)) ?? + (await readState()) ?? (() => { throw new Error(`Continuation claim ${claim.claimId} disappeared during repair`); })(); @@ -4971,7 +5038,7 @@ function continuationRepairEventId( .slice(0, 32)}`; } -function buildContinuationRepairStartEvent(claim: ContinuationClaimV1): RuntimeEvent { +function buildContinuationRepairStartEvent(claim: ContinuationClaim): RuntimeEvent { const source = claim.boundary.segments.at(-1)!; return { id: continuationRepairEventId('start', claim.claimId), @@ -5009,7 +5076,7 @@ function assertClaimOwnsHostedLinkedChildAdmission( runId: string; execution: Exclude; }, - claim: ContinuationClaimV1, + claim: ContinuationClaim, ): void { if ( input.execution.kind !== 'linked_child_resume' && @@ -5027,10 +5094,11 @@ function assertClaimOwnsHostedLinkedChildAdmission( const header = claim.targetRunHeader; const source = claim.boundary.segments.at(-1)!; const continuationSource = header.continuationSource; - const continuationSourceV2 = + const canonicalContinuationSource = continuationSource !== undefined && 'protocol' in continuationSource && - continuationSource.protocol === 'continuation_source_v2' + (continuationSource.protocol === 'continuation_source_v2' || + continuationSource.protocol === 'continuation_source_v3') ? continuationSource : undefined; if ( @@ -5043,10 +5111,10 @@ function assertClaimOwnsHostedLinkedChildAdmission( header.agentName !== input.execution.agentName || source.identity.sessionId !== input.sessionId || source.identity.runId !== input.execution.sourceRunId || - !continuationSourceV2 || - continuationSourceV2.claimId !== claim.claimId || - continuationSourceV2.boundaryDigest !== claim.boundaryDigest || - continuationSourceV2.sourceRunId !== input.execution.sourceRunId + !canonicalContinuationSource || + canonicalContinuationSource.claimId !== claim.claimId || + canonicalContinuationSource.boundaryDigest !== claim.boundaryDigest || + canonicalContinuationSource.sourceRunId !== input.execution.sourceRunId ) { throw new Error('Linked child admission continuation claim identity is inconsistent'); } diff --git a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts index be1d955a10..578865897b 100644 --- a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts @@ -29,10 +29,18 @@ import { canonicalToolArgsHash } from '@maka/core/tool-args-identity'; import { buildImmutableRuntimePrefix, createRuntimeBoundaryCursor, + digestWorkspaceBoundContinuationBoundary, runtimePrefixSegment, + type ContinuationClaim, + type ContinuationClaimV2, type ContinuationClaimV1, type ImmutableRuntimePrefixV1, + type ManagedWorkspaceContinuationBoundaryV1, } from '@maka/core/runtime-boundary'; +import { + workspaceMutationPolicyHashV1, + type WorkspaceBaselineAuthorityInput, +} from '@maka/core/workspace-version-authority'; import { ToolLedgerCorruptionError, ToolLedgerRejectionError, @@ -42,6 +50,12 @@ import { createSqliteRuntimeStore, type SqliteRuntimeStoreFailpoint, } from '../sqlite-runtime-store.js'; +import { + bindWorkspaceBaselineAuthorityStoreRootInternal, + commitWorkspaceBaselineInternal, +} from '../workspace-version-authority-internal.js'; + +const TEST_STORAGE_ROOT_ID = 'a'.repeat(64); describe('SqliteRuntimeStore', () => { it('applies versioned migrations and reopens the same database without rewriting schema', async () => { @@ -61,6 +75,64 @@ describe('SqliteRuntimeStore', () => { }); }); + it('migrates a populated schema 14 continuation authority without changing its v1 claim', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-continuation-schema-14-')); + const dbPath = join(root, 'runtime.sqlite'); + const claim = continuationClaim(); + try { + const current = createSqliteRuntimeStore(dbPath); + await persistImmutablePrefix(current, continuationSourcePrefix()); + assert.equal((await current.claimContinuation({ claim })).kind, 'acquired'); + current.close(); + + const rewind = new DatabaseSync(dbPath); + try { + rewindContinuationClaimsToSchema14(rewind); + } finally { + rewind.close(); + } + + const upgraded = createSqliteRuntimeStore(dbPath); + try { + assert.equal(upgraded.schemaVersion(), SQLITE_RUNTIME_SCHEMA_VERSION); + assert.deepEqual( + await upgraded.readContinuationClaimByBoundary(claim.boundaryDigest), + claim, + ); + } finally { + upgraded.close(); + } + } finally { + await rm(root, { recursive: true, force: true }); + } + }); + + it('fails closed when the workspace-bound continuation capability is missing', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-workspace-continuation-capability-')); + const dbPath = join(root, 'runtime.sqlite'); + const store = createSqliteRuntimeStore(dbPath); + store.close(); + try { + const raw = new DatabaseSync(dbPath); + try { + raw + .prepare( + `DELETE FROM runtime_capabilities + WHERE capability = 'runtime_workspace_bound_continuation_authority'`, + ) + .run(); + } finally { + raw.close(); + } + assert.throws( + () => createSqliteRuntimeStore(dbPath), + /runtime_workspace_bound_continuation_authority@1 is unavailable/, + ); + } finally { + await rm(root, { recursive: true, force: true }); + } + }); + it('refuses every post-terminal append as the typed sealed-run boundary', async () => { await withStore(async (store) => { const opening = functionCallEvent({ @@ -816,6 +888,148 @@ describe('SqliteRuntimeStore', () => { }); }); + it('atomically binds a continuation claim to the exact accepted workspace head', async () => { + await withStore(async (store) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const claim = workspaceBoundContinuationClaim(workspaceBoundary); + await persistImmutablePrefix(store, continuationSourcePrefix()); + + const acquired = await store.claimWorkspaceBoundContinuation({ claim }); + const existing = await store.claimWorkspaceBoundContinuation({ + claim: structuredClone(claim), + }); + + assert.equal(acquired.kind, 'acquired'); + assert.equal(existing.kind, 'existing'); + assert.deepEqual(existing.claim, claim); + assert.deepEqual( + await store.readWorkspaceBoundContinuationClaimByBoundary(claim.boundaryDigest), + claim, + ); + }); + }); + + it('commits a workspace-bound continuation start only through its dedicated authority', async () => { + await withStore(async (store) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const claim = workspaceBoundContinuationClaim(workspaceBoundary); + await persistImmutablePrefix(store, continuationSourcePrefix()); + assert.equal((await store.claimWorkspaceBoundContinuation({ claim })).kind, 'acquired'); + const event = continuationStartEvent(claim); + + const committed = await store.commitWorkspaceBoundContinuationStart({ claim, event }); + const retry = await store.commitWorkspaceBoundContinuationStart({ claim, event }); + const state = await store.readWorkspaceBoundContinuationClaimStateByBoundary( + claim.boundaryDigest, + ); + + assert.equal(committed.created, true); + assert.equal(retry.created, false); + assert.deepEqual(state, { + claim, + startEventId: event.id, + startKind: 'runtime_admission', + }); + assert.equal(await store.readContinuationClaimByBoundary(claim.boundaryDigest), undefined); + }); + }); + + it('revalidates the accepted workspace boundary when committing continuation start', async () => { + await withStore(async (store, dbPath) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const claim = workspaceBoundContinuationClaim(workspaceBoundary); + await persistImmutablePrefix(store, continuationSourcePrefix()); + assert.equal((await store.claimWorkspaceBoundContinuation({ claim })).kind, 'acquired'); + + const tamper = new DatabaseSync(dbPath); + try { + tamper + .prepare( + `UPDATE runtime_workspace_heads + SET revision = revision + 1 + WHERE workspace_id = ? AND workspace_epoch_id = ?`, + ) + .run(workspaceBoundary.workspaceId, workspaceBoundary.workspaceEpochId); + } finally { + tamper.close(); + } + + await assert.rejects( + store.commitWorkspaceBoundContinuationStart({ + claim, + event: continuationStartEvent(claim), + }), + /workspace version projection is incomplete|workspace boundary no longer matches accepted authority/i, + ); + assert.deepEqual( + await store.readRuntimeEvents(claim.target.sessionId, claim.target.runId), + [], + ); + }); + }); + + it('rejects a workspace-bound claim when caller evidence differs from accepted authority', async () => { + await withStore(async (store) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const staleBoundary = { ...workspaceBoundary, revision: workspaceBoundary.revision + 1 }; + const claim = workspaceBoundContinuationClaim(staleBoundary); + await persistImmutablePrefix(store, continuationSourcePrefix()); + + await assert.rejects( + store.claimWorkspaceBoundContinuation({ claim }), + /workspace boundary no longer matches accepted authority/i, + ); + assert.equal( + await store.readWorkspaceBoundContinuationClaimByBoundary(claim.boundaryDigest), + undefined, + ); + }); + }); + + it('rejects a caller-asserted execution profile that is not bound by epoch policy', async () => { + await withStore(async (store) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const claim = workspaceBoundContinuationClaim({ + ...workspaceBoundary, + executionProfileDigest: `sha256:${'c'.repeat(64)}`, + }); + await persistImmutablePrefix(store, continuationSourcePrefix()); + + await assert.rejects( + store.claimWorkspaceBoundContinuation({ claim }), + /execution profile conflict/i, + ); + }); + }); + + it('fails closed when stored workspace-bound claim JSON is changed independently', async () => { + await withStore(async (store, dbPath) => { + const workspaceBoundary = await openContinuationWorkspaceBoundary(store); + const claim = workspaceBoundContinuationClaim(workspaceBoundary); + await persistImmutablePrefix(store, continuationSourcePrefix()); + assert.equal((await store.claimWorkspaceBoundContinuation({ claim })).kind, 'acquired'); + + const tamper = new DatabaseSync(dbPath); + try { + tamper + .prepare( + 'UPDATE runtime_continuation_claims SET workspace_boundary_json = ? WHERE claim_id = ?', + ) + .run( + JSON.stringify({ ...workspaceBoundary, revision: workspaceBoundary.revision + 1 }), + claim.claimId, + ); + } finally { + tamper.close(); + } + + await assert.rejects( + store.readWorkspaceBoundContinuationClaimByBoundary(claim.boundaryDigest), + /workspace boundary digest mismatch|workspace row\/payload mismatch/i, + ); + }); + }); + it('rejects a continuation claim whose immediate source boundary is not durable', async () => { await withStore(async (store) => { const claim = continuationClaim(); @@ -1899,6 +2113,192 @@ function continuationClaimForBoundary( }; } +function workspaceBoundContinuationClaim( + workspaceBoundary: ManagedWorkspaceContinuationBoundaryV1, +): ContinuationClaimV2 { + const legacy = continuationClaim(); + const boundaryDigest = digestWorkspaceBoundContinuationBoundary( + legacy.boundary, + workspaceBoundary, + ); + const source = legacy.boundary.segments.at(-1)!; + return { + ...legacy, + protocol: 'continuation_claim_v2', + boundaryDigest, + workspaceBoundary, + targetRunHeader: { + ...legacy.targetRunHeader, + continuationSource: { + protocol: 'continuation_source_v3', + claimId: legacy.claimId, + boundaryDigest, + sourceInvocationId: source.identity.invocationId, + sourceRunId: source.identity.runId, + sourceTurnId: source.identity.turnId, + sourceRuntimeEventHighWater: source.position.lastEventSeq, + sourcePrefixDigest: source.prefixDigest, + replayManifestDigest: legacy.boundary.manifestDigest, + }, + }, + }; +} + +async function openContinuationWorkspaceBoundary( + store: Store, +): Promise { + const input = continuationWorkspaceBaselineInput(); + const executionProfileDigest = `sha256:${'b'.repeat(64)}` as const; + bindWorkspaceBaselineAuthorityStoreRootInternal(store, TEST_STORAGE_ROOT_ID); + const opened = await commitWorkspaceBaselineInternal(store, input); + return { + protocol: 'managed_workspace_continuation_boundary_v1', + storageRootId: TEST_STORAGE_ROOT_ID, + repositoryId: input.epoch.repositoryId, + workspaceId: input.epoch.workspaceId, + workspaceEpochId: input.epoch.workspaceEpochId, + workspaceInstanceId: input.epoch.workspaceInstanceId, + workspaceVersionId: opened.head.workspaceVersionId, + acceptedEventId: opened.head.acceptedEventId, + revision: opened.head.revision, + objectFormat: input.epoch.objectFormat, + sourceCommitOid: input.epoch.sourceCommitOid, + sourceTreeOid: input.epoch.sourceTreeOid, + commitOid: opened.head.commitOid, + treeOid: opened.head.treeOid, + materializationProfileDigest: input.epoch.materializationProfileDigest, + policyHash: input.epoch.policyHash, + executionProfileDigest, + }; +} + +function continuationWorkspaceBaselineInput(): WorkspaceBaselineAuthorityInput { + const materializationProfileDigest = `sha256:${'3'.repeat(64)}` as const; + const executionProfileDigest = `sha256:${'b'.repeat(64)}` as const; + return { + epochOpenedEventId: 'continuation-workspace-epoch-event', + baselineAcceptedEventId: 'continuation-workspace-version-event', + committedAt: 5, + epoch: { + repositoryId: 'repository_11111111111111111111111111111111', + workspaceId: 'workspace_22222222222222222222222222222222', + workspaceEpochId: 'epoch_33333333333333333333333333333333', + workspaceInstanceId: 'instance_44444444444444444444444444444444', + mode: 'managed_worktree', + objectFormat: 'sha1', + sourceCommitOid: '1'.repeat(40), + sourceTreeOid: '2'.repeat(40), + materializationProfileDigest, + materializationSemantics: 'git_tree_materialized_with_fixed_config_v1', + policyHash: workspaceMutationPolicyHashV1( + materializationProfileDigest, + executionProfileDigest, + ), + }, + baseline: { + workspaceVersionId: 'version_55555555555555555555555555555555', + commitOid: '5'.repeat(40), + treeOid: '2'.repeat(40), + treeDeltaDigest: `sha256:${'6'.repeat(64)}`, + changedFileCount: 1, + deletedFileCount: 0, + }, + }; +} + +function rewindContinuationClaimsToSchema14(db: DatabaseSync): void { + db.exec(` + PRAGMA foreign_keys = OFF; + BEGIN IMMEDIATE; + ALTER TABLE runtime_continuation_claims + RENAME TO runtime_continuation_claims_schema_15; + + CREATE TABLE runtime_continuation_claims ( + claim_id TEXT PRIMARY KEY, + source_session_id TEXT NOT NULL, + source_invocation_id TEXT NOT NULL, + source_run_id TEXT NOT NULL, + source_turn_id TEXT NOT NULL, + source_event_high_water INTEGER NOT NULL CHECK (source_event_high_water > 0), + source_prefix_digest TEXT NOT NULL, + boundary_digest TEXT NOT NULL UNIQUE, + boundary_json TEXT NOT NULL, + provider_projection_version INTEGER NOT NULL CHECK (provider_projection_version = 1), + provider_replay_digest TEXT NOT NULL, + target_session_id TEXT NOT NULL, + target_invocation_id TEXT NOT NULL UNIQUE, + target_run_id TEXT NOT NULL UNIQUE, + target_turn_id TEXT NOT NULL, + target_run_header_json TEXT NOT NULL, + claimed_at INTEGER NOT NULL, + start_event_id TEXT UNIQUE REFERENCES runtime_events(event_id), + start_kind TEXT CHECK ( + start_kind IS NULL OR start_kind IN ('runtime_admission', 'claim_repair') + ), + protocol_version INTEGER NOT NULL CHECK (protocol_version = 1), + UNIQUE ( + source_session_id, + source_run_id, + source_event_high_water, + source_prefix_digest + ), + UNIQUE (target_session_id, target_turn_id) + ); + + INSERT INTO runtime_continuation_claims ( + claim_id, + source_session_id, + source_invocation_id, + source_run_id, + source_turn_id, + source_event_high_water, + source_prefix_digest, + boundary_digest, + boundary_json, + provider_projection_version, + provider_replay_digest, + target_session_id, + target_invocation_id, + target_run_id, + target_turn_id, + target_run_header_json, + claimed_at, + start_event_id, + start_kind, + protocol_version + ) + SELECT + claim_id, + source_session_id, + source_invocation_id, + source_run_id, + source_turn_id, + source_event_high_water, + source_prefix_digest, + boundary_digest, + boundary_json, + provider_projection_version, + provider_replay_digest, + target_session_id, + target_invocation_id, + target_run_id, + target_turn_id, + target_run_header_json, + claimed_at, + start_event_id, + start_kind, + protocol_version + FROM runtime_continuation_claims_schema_15; + + DROP TABLE runtime_continuation_claims_schema_15; + DELETE FROM runtime_capabilities + WHERE capability = 'runtime_workspace_bound_continuation_authority'; + PRAGMA user_version = 14; + COMMIT; + PRAGMA foreign_keys = ON; + `); +} + function continuationSourcePrefix(): ImmutableRuntimePrefixV1 { return buildImmutableRuntimePrefix( { @@ -1960,7 +2360,7 @@ async function persistImmutablePrefix( } function continuationStartEvent( - claim: ContinuationClaimV1, + claim: ContinuationClaim, overrides: { id?: string; provenance?: 'runtime_admission' | 'claim_repair'; 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 31611f44d0..8f1c9de05b 100644 --- a/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts +++ b/packages/storage/src/__tests__/workspace-version-authority-persistence.test.ts @@ -665,7 +665,7 @@ describe('workspace version persistence authority', () => { bindWorkspaceBaselineAuthorityStoreRootInternal(upgraded, TEST_STORAGE_ROOT_ID); registerWorkspaceSuccessorCandidateVerifierInternal(upgraded, verifyTestCandidate); try { - assert.equal(upgraded.schemaVersion(), 14); + assert.equal(upgraded.schemaVersion(), 15); assert.equal( ( await upgraded.readWorkspaceHead( @@ -1350,6 +1350,8 @@ function recreateWorkspaceTablesAsSchema12(database: DatabaseSync): void { DROP TABLE runtime_workspace_heads_schema_13; DROP TABLE runtime_workspace_versions_schema_13; DROP TABLE runtime_managed_mutation_reservations; + DELETE FROM runtime_capabilities + WHERE capability = 'runtime_workspace_bound_continuation_authority'; PRAGMA user_version = 12; COMMIT; PRAGMA foreign_keys = ON; diff --git a/packages/storage/src/agent-run-store.ts b/packages/storage/src/agent-run-store.ts index 27481f3fec..6ff0a18e83 100644 --- a/packages/storage/src/agent-run-store.ts +++ b/packages/storage/src/agent-run-store.ts @@ -2025,7 +2025,7 @@ function normalizeRootExecutionDescriptor(value: unknown): RootExecutionDescript }); } if (value.kind === 'safe_boundary_continuation') { - const keys = [ + const legacyKeys = [ 'kind', 'sourceInvocationId', 'sourceRunId', @@ -2037,6 +2037,8 @@ function normalizeRootExecutionDescriptor(value: unknown): RootExecutionDescript 'safetyDigest', 'targetInvocationId', ]; + const hasReplayManifestDigest = Object.hasOwn(value, 'replayManifestDigest'); + const keys = hasReplayManifestDigest ? [...legacyKeys, 'replayManifestDigest'] : legacyKeys; if ( !hasExactKeys(value, keys) || typeof value.sourceInvocationId !== 'string' || @@ -2050,6 +2052,7 @@ function normalizeRootExecutionDescriptor(value: unknown): RootExecutionDescript typeof value.claimId !== 'string' || !isSafeId(value.claimId) || !isSha256Digest(value.boundaryDigest) || + (hasReplayManifestDigest && !isSha256Digest(value.replayManifestDigest)) || !isSha256Digest(value.providerReplayDigest) || !isSha256Digest(value.safetyDigest) || typeof value.targetInvocationId !== 'string' || @@ -2065,6 +2068,9 @@ function normalizeRootExecutionDescriptor(value: unknown): RootExecutionDescript sourceRuntimeEventHighWater: value.sourceRuntimeEventHighWater as number, claimId: value.claimId, boundaryDigest: value.boundaryDigest, + ...(hasReplayManifestDigest + ? { replayManifestDigest: value.replayManifestDigest as `sha256:${string}` } + : {}), providerReplayDigest: value.providerReplayDigest, safetyDigest: value.safetyDigest, targetInvocationId: value.targetInvocationId, diff --git a/packages/storage/src/execution-stores-workspace-authority-internal.ts b/packages/storage/src/execution-stores-workspace-authority-internal.ts index ba4f55b7aa..999552150c 100644 --- a/packages/storage/src/execution-stores-workspace-authority-internal.ts +++ b/packages/storage/src/execution-stores-workspace-authority-internal.ts @@ -17,6 +17,10 @@ * under the License. */ +import type { + ManagedWorkspaceContinuationBoundaryV1, + RuntimeBoundaryDigest, +} from '@maka/core/runtime-boundary'; import type { WorkspaceBaselineAuthorityInput, WorkspaceBaselineCommitResult, @@ -32,6 +36,7 @@ import { commitVerifiedWorkspaceSuccessorInternal, readActiveManagedMutationInternal, readManagedMutationEvidenceInternal, + readWorkspaceContinuationBoundaryInternal, readWorkspaceEpochInternal, readWorkspaceHeadInternal, readWorkspaceVersionInternal, @@ -53,6 +58,18 @@ export interface ExecutionStoresWorkspaceBaselineAuthorityCapabilityInternal { readonly kind: 'execution_stores_workspace_baseline_authority_v1'; } +export interface ExecutionStoresWorkspaceContinuationAuthorityCapabilityInternal { + readonly kind: 'execution_stores_workspace_continuation_authority_v1'; +} + +export interface ExecutionStoresWorkspaceContinuationAuthorityInternal { + readContinuationBoundary( + workspaceId: string, + workspaceEpochId: string, + executionProfileDigest: RuntimeBoundaryDigest, + ): Promise; +} + export interface ExecutionStoresWorkspaceBaselineAuthorityInternal { commitBaseline(importedRepositoryProof: object): Promise; readEpoch( @@ -122,6 +139,10 @@ const sources = new WeakMap(); const capabilities = new WeakMap(); const noEffectCapabilities = new WeakMap(); const baselineCapabilities = new WeakMap(); +const continuationCapabilities = new WeakMap< + object, + AuthoritySource & { readonly ownerToken: object } +>(); const projectionCapabilities = new WeakMap< object, { @@ -206,6 +227,46 @@ export function issueExecutionStoresWorkspaceBaselineAuthorityInternal(input: { return capability; } +export function issueExecutionStoresWorkspaceContinuationAuthorityInternal(input: { + readonly ownerToken: object; + readonly stores: object; +}): ExecutionStoresWorkspaceContinuationAuthorityCapabilityInternal { + const source = sources.get(input.stores); + if (!source) throw new Error('Execution stores workspace continuation source is unavailable'); + adoptWorkspaceBaselineAuthorityStoreRootInternal(source.store, source.rootId); + const capability = Object.freeze({ + kind: 'execution_stores_workspace_continuation_authority_v1' as const, + }); + continuationCapabilities.set( + capability, + Object.freeze({ ...source, ownerToken: input.ownerToken }), + ); + return capability; +} + +export function requireExecutionStoresWorkspaceContinuationAuthorityInternal( + ownerToken: object, + capability: ExecutionStoresWorkspaceContinuationAuthorityCapabilityInternal, +): ExecutionStoresWorkspaceContinuationAuthorityInternal { + const record = continuationCapabilities.get(capability); + if (!record || record.ownerToken !== ownerToken) { + throw new Error('Execution stores workspace continuation authority capability is invalid'); + } + return Object.freeze({ + readContinuationBoundary: ( + workspaceId: string, + workspaceEpochId: string, + executionProfileDigest: RuntimeBoundaryDigest, + ) => + readWorkspaceContinuationBoundaryInternal( + record.store, + workspaceId, + workspaceEpochId, + executionProfileDigest, + ), + }); +} + export function requireExecutionStoresWorkspaceBaselineAuthorityInternal( ownerToken: object, capability: ExecutionStoresWorkspaceBaselineAuthorityCapabilityInternal, diff --git a/packages/storage/src/execution-stores.ts b/packages/storage/src/execution-stores.ts index 2cbc0463d0..2654562203 100644 --- a/packages/storage/src/execution-stores.ts +++ b/packages/storage/src/execution-stores.ts @@ -24,7 +24,10 @@ import type { AgentRunProjectionKey, } from '@maka/core/agent-run'; import type { RuntimeEvent, ToolBoundaryProtocol } from '@maka/core/runtime-event'; -import type { RuntimeContinuationAuthorityStore } from '@maka/core/runtime-event-store'; +import type { + RuntimeContinuationAuthorityStore, + RuntimeWorkspaceBoundContinuationAuthorityStore, +} from '@maka/core/runtime-event-store'; import type { SessionHeader, SessionSummary, StoredMessage, TurnRecord } from '@maka/core/session'; import type { SessionListFilter } from '@maka/core/runtime-inputs'; import { @@ -139,7 +142,8 @@ export type { export type ExecutionSessionWriter = SessionAuthorityStore; export type ExecutionAgentRunWriter = DurableAgentRunStore; export type ExecutionRuntimeEventWriter = DurableRuntimeEventStore & - RuntimeContinuationAuthorityStore & { + RuntimeContinuationAuthorityStore & + RuntimeWorkspaceBoundContinuationAuthorityStore & { readonly toolBoundaryProtocol: ToolBoundaryProtocol; commitToolPrepared(input: CommitToolPreparedInput): Promise; commitToolOutcome(input: CommitToolOutcomeInput): Promise; @@ -517,6 +521,8 @@ async function createExecutionStoresForWrite run(() => runtimeEventStore.appendRuntimeEvent(sessionId, runId, event, options)), @@ -551,6 +557,20 @@ async function createExecutionStoresForWrite runtimeEventStore.commitContinuationStart(input)), commitContinuationRepairStart: (input) => run(() => runtimeEventStore.commitContinuationRepairStart(input)), + claimWorkspaceBoundContinuation: (input) => + run(() => runtimeEventStore.claimWorkspaceBoundContinuation(input)), + readWorkspaceBoundContinuationClaimByBoundary: (boundaryDigest) => + run(() => runtimeEventStore.readWorkspaceBoundContinuationClaimByBoundary(boundaryDigest)), + readWorkspaceBoundContinuationClaimStateByBoundary: (boundaryDigest) => + run(() => + runtimeEventStore.readWorkspaceBoundContinuationClaimStateByBoundary(boundaryDigest), + ), + listWorkspaceBoundContinuationClaimsForRecovery: (sessionId) => + run(() => runtimeEventStore.listWorkspaceBoundContinuationClaimsForRecovery(sessionId)), + commitWorkspaceBoundContinuationStart: (input) => + run(() => runtimeEventStore.commitWorkspaceBoundContinuationStart(input)), + commitWorkspaceBoundContinuationRepairStart: (input) => + run(() => runtimeEventStore.commitWorkspaceBoundContinuationRepairStart(input)), readImmutableSteeringMessageProof: (sessionId, messageId) => run(() => runtimeEventStore.readImmutableSteeringMessageProof(sessionId, messageId)), repairImmutableSteeringMessageProofsForRecovery: (sessionId) => diff --git a/packages/storage/src/sqlite-runtime-schema.ts b/packages/storage/src/sqlite-runtime-schema.ts index 47add708ee..f2dd8d9223 100644 --- a/packages/storage/src/sqlite-runtime-schema.ts +++ b/packages/storage/src/sqlite-runtime-schema.ts @@ -19,11 +19,14 @@ import type { DatabaseSync } from 'node:sqlite'; -export const SQLITE_RUNTIME_SCHEMA_VERSION = 14; +export const SQLITE_RUNTIME_SCHEMA_VERSION = 15; 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'; export const RUNTIME_CONTINUATION_AUTHORITY_CAPABILITY_VERSION = 1; +export const RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY = + 'runtime_workspace_bound_continuation_authority'; +export const RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY_VERSION = 1; export const RUNTIME_WORKSPACE_VERSION_AUTHORITY_CAPABILITY = 'runtime_workspace_version_authority'; export const RUNTIME_WORKSPACE_VERSION_AUTHORITY_CAPABILITY_VERSION = 1; const SQLITE_INITIALIZATION_BUSY_TIMEOUT_MS = 5_000; @@ -443,6 +446,102 @@ const MIGRATIONS: ReadonlyMap = new Map([ ); `, ], + [ + 15, + ` + ALTER TABLE runtime_continuation_claims + RENAME TO runtime_continuation_claims_v14; + + CREATE TABLE runtime_continuation_claims ( + claim_id TEXT PRIMARY KEY, + source_session_id TEXT NOT NULL, + source_invocation_id TEXT NOT NULL, + source_run_id TEXT NOT NULL, + source_turn_id TEXT NOT NULL, + source_event_high_water INTEGER NOT NULL CHECK (source_event_high_water > 0), + source_prefix_digest TEXT NOT NULL, + boundary_digest TEXT NOT NULL UNIQUE, + boundary_json TEXT NOT NULL, + workspace_boundary_json TEXT, + provider_projection_version INTEGER NOT NULL CHECK (provider_projection_version = 1), + provider_replay_digest TEXT NOT NULL, + target_session_id TEXT NOT NULL, + target_invocation_id TEXT NOT NULL UNIQUE, + target_run_id TEXT NOT NULL UNIQUE, + target_turn_id TEXT NOT NULL, + target_run_header_json TEXT NOT NULL, + claimed_at INTEGER NOT NULL, + start_event_id TEXT UNIQUE REFERENCES runtime_events(event_id), + start_kind TEXT CHECK ( + start_kind IS NULL OR start_kind IN ('runtime_admission', 'claim_repair') + ), + protocol_version INTEGER NOT NULL CHECK (protocol_version IN (1, 2)), + UNIQUE ( + source_session_id, + source_run_id, + source_event_high_water, + source_prefix_digest + ), + UNIQUE (target_session_id, target_turn_id), + CHECK ( + (protocol_version = 1 AND workspace_boundary_json IS NULL) + OR (protocol_version = 2 AND workspace_boundary_json IS NOT NULL) + ) + ); + + INSERT INTO runtime_continuation_claims ( + claim_id, + source_session_id, + source_invocation_id, + source_run_id, + source_turn_id, + source_event_high_water, + source_prefix_digest, + boundary_digest, + boundary_json, + workspace_boundary_json, + provider_projection_version, + provider_replay_digest, + target_session_id, + target_invocation_id, + target_run_id, + target_turn_id, + target_run_header_json, + claimed_at, + start_event_id, + start_kind, + protocol_version + ) + SELECT + claim_id, + source_session_id, + source_invocation_id, + source_run_id, + source_turn_id, + source_event_high_water, + source_prefix_digest, + boundary_digest, + boundary_json, + NULL, + provider_projection_version, + provider_replay_digest, + target_session_id, + target_invocation_id, + target_run_id, + target_turn_id, + target_run_header_json, + claimed_at, + start_event_id, + start_kind, + protocol_version + FROM runtime_continuation_claims_v14; + + DROP TABLE runtime_continuation_claims_v14; + + INSERT INTO runtime_capabilities(capability, version) + VALUES ('runtime_workspace_bound_continuation_authority', 1); + `, + ], ]); 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 ff2e8f9c36..9c4915f765 100644 --- a/packages/storage/src/sqlite-runtime-store.ts +++ b/packages/storage/src/sqlite-runtime-store.ts @@ -27,6 +27,7 @@ import { buildWorkspaceBaselineAuthorityEvents, buildWorkspaceSuccessorAuthorityEvent, scanWorkspaceBaselineAuthority, + workspaceMutationPolicyHashV1, WORKSPACE_AUTHORITY_SESSION_ID, WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1, type ScannedWorkspaceBaselineAuthority, @@ -53,13 +54,17 @@ import { import { RunSealedError, RUNTIME_CONTINUATION_AUTHORITY_V1, + RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_V1, TOOL_RECOVERY_BUNDLE_CAPABILITY_V1, type ContinuationClaimResult, type ContinuationClaimStateV1, + type ContinuationClaimStateV2, type RuntimeContinuationAuthorityStore, type RuntimeRecoveryBundleCommit, type RuntimeRecoveryBundleStore, type RuntimeWorkspaceVersionAuthorityStore, + type RuntimeWorkspaceBoundContinuationAuthorityStore, + type WorkspaceBoundContinuationClaimResult, } from '@maka/core/runtime-event-store'; import { type ToolRecoveryDecisionFact } from '@maka/core/tool-recovery-fact'; import { canonicalToolArgsHash, stableJsonStringify } from '@maka/core/tool-args-identity'; @@ -77,7 +82,9 @@ import { import { buildImmutableRuntimePrefix, decodeContinuationClaim, + type ContinuationClaim, type ContinuationClaimV1, + type ManagedWorkspaceContinuationBoundaryV1, type ImmutableRuntimePrefixV1, type RuntimeBoundaryDigest, } from '@maka/core/runtime-boundary'; @@ -93,6 +100,8 @@ import { RUNTIME_RECOVERY_AUTHORITY_CAPABILITY_VERSION, RUNTIME_CONTINUATION_AUTHORITY_CAPABILITY, RUNTIME_CONTINUATION_AUTHORITY_CAPABILITY_VERSION, + RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY, + RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY_VERSION, RUNTIME_WORKSPACE_VERSION_AUTHORITY_CAPABILITY, RUNTIME_WORKSPACE_VERSION_AUTHORITY_CAPABILITY_VERSION, SQLITE_RUNTIME_SCHEMA_VERSION, @@ -268,12 +277,15 @@ export class SqliteRuntimeStore implements RuntimeRecoveryBundleStore, RuntimeContinuationAuthorityStore, + RuntimeWorkspaceBoundContinuationAuthorityStore, RuntimeWorkspaceVersionAuthorityStore { readonly durability = 'canonical' as const; readonly toolBoundaryProtocol = 't1_after_preflight_v1' as const; readonly recoveryBundleCapability = TOOL_RECOVERY_BUNDLE_CAPABILITY_V1; readonly continuationAuthorityCapability = RUNTIME_CONTINUATION_AUTHORITY_V1; + readonly workspaceBoundContinuationAuthorityCapability = + RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_V1; readonly workspaceVersionAuthorityCapability = WORKSPACE_VERSION_AUTHORITY_CAPABILITY_V1; private readonly db: DatabaseSync; private readonly databaseLease?: OperationalStateDatabaseLease; @@ -293,6 +305,7 @@ export class SqliteRuntimeStore this.db = options.databaseLease.database; assertRecoveryAuthorityCapability(this.db); assertContinuationAuthorityCapability(this.db); + assertWorkspaceBoundContinuationAuthorityCapability(this.db); assertWorkspaceVersionAuthorityCapability(this.db); if (!options.readOnly) { this.registerWorkspaceBaselineAuthorityWriter(); @@ -319,6 +332,7 @@ export class SqliteRuntimeStore } assertRecoveryAuthorityCapability(this.db); assertContinuationAuthorityCapability(this.db); + assertWorkspaceBoundContinuationAuthorityCapability(this.db); assertWorkspaceVersionAuthorityCapability(this.db); if (!options.readOnly) { this.registerWorkspaceBaselineAuthorityWriter(); @@ -861,8 +875,38 @@ export class SqliteRuntimeStore ); } - async claimContinuation(input: { claim: ContinuationClaimV1 }): Promise { - const claim = decodeContinuationClaim(input.claim); + async claimContinuation(input: { + claim: Extract; + }): Promise { + const result = await this.claimContinuationAuthority(input.claim); + if (result.claim.protocol !== 'continuation_claim_v1') { + throw new Error('Legacy continuation authority conflicts with a workspace-bound claim'); + } + if (result.kind === 'acquired') return { kind: 'acquired', claim: result.claim }; + if (result.kind === 'existing') return { kind: 'existing', claim: result.claim }; + return { kind: 'conflict', claim: result.claim }; + } + + async claimWorkspaceBoundContinuation(input: { + claim: Extract; + }): Promise { + const result = await this.claimContinuationAuthority(input.claim); + if (result.claim.protocol !== 'continuation_claim_v2') { + throw new Error('Workspace-bound continuation authority conflicts with a legacy claim'); + } + if (result.kind === 'acquired') return { kind: 'acquired', claim: result.claim }; + if (result.kind === 'existing') return { kind: 'existing', claim: result.claim }; + return { kind: 'conflict', claim: result.claim }; + } + + private async claimContinuationAuthority( + inputClaim: ContinuationClaim, + ): Promise< + | { kind: 'acquired'; claim: ContinuationClaim } + | { kind: 'existing'; claim: ContinuationClaim } + | { kind: 'conflict'; claim: ContinuationClaim } + > { + const claim = decodeContinuationClaim(inputClaim); if ( claim.target.sessionId === WORKSPACE_AUTHORITY_SESSION_ID || claim.boundary.segments.some( @@ -872,13 +916,20 @@ export class SqliteRuntimeStore throw new Error('Continuation cannot target the reserved workspace authority stream'); } const boundaryJson = stableJsonStringify(claim.boundary); + const workspaceBoundaryJson = + claim.protocol === 'continuation_claim_v2' + ? stableJsonStringify(claim.workspaceBoundary) + : null; return this.transaction(() => { this.assertContinuationAuthorityIntegrity(); this.assertContinuationBoundaryMatchesLedger(claim); const byBoundary = this.readContinuationClaimRow('boundary_digest = ?', claim.boundaryDigest); if (byBoundary) { const existing = decodeContinuationClaimRow(byBoundary); - if (byBoundary.boundary_json !== boundaryJson) { + if ( + byBoundary.boundary_json !== boundaryJson || + byBoundary.workspace_boundary_json !== workspaceBoundaryJson + ) { throw new Error('Continuation claim boundary digest has conflicting canonical JSON'); } return { kind: 'existing', claim: existing }; @@ -910,6 +961,9 @@ export class SqliteRuntimeStore if (this.continuationTargetHasRuntimeState(claim)) { throw new Error('Continuation claim target RuntimeEvent ledger is not empty'); } + if (claim.protocol === 'continuation_claim_v2') { + this.assertContinuationWorkspaceBoundaryMatchesAuthority(claim.workspaceBoundary); + } try { this.db @@ -924,6 +978,7 @@ export class SqliteRuntimeStore source_prefix_digest, boundary_digest, boundary_json, + workspace_boundary_json, provider_projection_version, provider_replay_digest, target_session_id, @@ -933,7 +988,7 @@ export class SqliteRuntimeStore target_run_header_json, claimed_at, protocol_version - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1) + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `) .run( claim.claimId, @@ -945,6 +1000,7 @@ export class SqliteRuntimeStore source.prefixDigest, claim.boundaryDigest, boundaryJson, + workspaceBoundaryJson, claim.providerProjectionVersion, claim.providerReplayDigest, claim.target.sessionId, @@ -953,6 +1009,7 @@ export class SqliteRuntimeStore claim.target.turnId, stableJsonStringify(claim.targetRunHeader), claim.claimedAt, + claim.protocol === 'continuation_claim_v2' ? 2 : 1, ); } catch (error) { const raced = @@ -989,17 +1046,40 @@ export class SqliteRuntimeStore async readContinuationClaimByBoundary( boundaryDigest: RuntimeBoundaryDigest, - ): Promise { + ): Promise | undefined> { return (await this.readContinuationClaimStateByBoundary(boundaryDigest))?.claim; } + async readWorkspaceBoundContinuationClaimByBoundary( + boundaryDigest: RuntimeBoundaryDigest, + ): Promise | undefined> { + return (await this.readWorkspaceBoundContinuationClaimStateByBoundary(boundaryDigest))?.claim; + } + + async readWorkspaceBoundContinuationClaimStateByBoundary( + boundaryDigest: RuntimeBoundaryDigest, + ): Promise { + if (!/^sha256:[0-9a-f]{64}$/.test(boundaryDigest)) { + throw new Error('Invalid continuation boundary digest'); + } + const row = this.readContinuationClaimRow( + 'boundary_digest = ? AND protocol_version = 2', + boundaryDigest, + ); + if (!row) return undefined; + return this.decodeWorkspaceBoundContinuationClaimStateRow(row); + } + async readContinuationClaimStateByBoundary( boundaryDigest: RuntimeBoundaryDigest, ): Promise { if (!/^sha256:[0-9a-f]{64}$/.test(boundaryDigest)) { throw new Error('Invalid continuation boundary digest'); } - const row = this.readContinuationClaimRow('boundary_digest = ?', boundaryDigest); + const row = this.readContinuationClaimRow( + 'boundary_digest = ? AND protocol_version = 1', + boundaryDigest, + ); return row ? this.decodeContinuationClaimStateRow(row) : undefined; } @@ -1016,6 +1096,7 @@ export class SqliteRuntimeStore source_prefix_digest, boundary_digest, boundary_json, + workspace_boundary_json, provider_projection_version, provider_replay_digest, target_session_id, @@ -1028,13 +1109,20 @@ export class SqliteRuntimeStore start_kind, protocol_version FROM runtime_continuation_claims - WHERE target_session_id = ? + WHERE target_session_id = ? AND protocol_version = 1 ORDER BY claimed_at ASC, claim_id ASC `) .all(sessionId) as unknown as ContinuationClaimStorageRow[]; return rows.map((row) => this.decodeContinuationClaimStateRow(row)); } + async listWorkspaceBoundContinuationClaimsForRecovery( + sessionId: string, + ): Promise { + const rows = this.readContinuationClaimRowsForSession(sessionId, 2); + return rows.map((row) => this.decodeWorkspaceBoundContinuationClaimStateRow(row)); + } + async commitContinuationStart(input: { claim: ContinuationClaimV1; event: RuntimeEvent; @@ -1049,14 +1137,31 @@ export class SqliteRuntimeStore return this.commitContinuationStartOfKind(input, 'claim_repair'); } + async commitWorkspaceBoundContinuationStart(input: { + claim: Extract; + event: RuntimeEvent; + }): Promise { + return this.commitContinuationStartOfKind(input, 'runtime_admission'); + } + + async commitWorkspaceBoundContinuationRepairStart(input: { + claim: Extract; + event: RuntimeEvent; + }): Promise { + return this.commitContinuationStartOfKind(input, 'claim_repair'); + } + private commitContinuationStartOfKind( input: { - claim: ContinuationClaimV1; + claim: ContinuationClaim; event: RuntimeEvent; }, startKind: 'runtime_admission' | 'claim_repair', ): ToolCommitResult { const claim = decodeContinuationClaim(input.claim); + if (claim.protocol !== input.claim.protocol) { + throw new Error('Continuation start claim protocol changed during decoding'); + } const event = canonicalizeRuntimeEventForStorage(input.event); assertNoReservedWorkspaceAuthorityAppend(event); assertContinuationStartEvent(claim, event, startKind); @@ -1069,6 +1174,12 @@ export class SqliteRuntimeStore if (!isDeepStrictEqual(storedClaim, claim)) { throw new Error('Continuation start claim identity conflict'); } + if (claim.protocol === 'continuation_claim_v2') { + // Claim acquisition and continuation start are separate durable boundaries. + // Re-observe the accepted head in this transaction so a stale managed claim + // can never start merely because an in-process execution lease survived. + this.assertContinuationWorkspaceBoundaryMatchesAuthority(claim.workspaceBoundary); + } if (row.start_event_id) { if (row.start_event_id !== event.id || row.start_kind !== startKind) { throw new Error('Continuation claim already has a different start event'); @@ -1561,9 +1672,33 @@ export class SqliteRuntimeStore readWorkspaceVersion, (workspaceInstanceId) => this.#readActiveManagedMutation(workspaceInstanceId), (operationId) => this.#readManagedMutationEvidence(operationId), + (workspaceId, workspaceEpochId, rootId, executionProfileDigest) => + this.#readWorkspaceContinuationBoundary( + workspaceId, + workspaceEpochId, + rootId, + executionProfileDigest, + ), ); } + async #readWorkspaceContinuationBoundary( + workspaceId: string, + workspaceEpochId: string, + rootId: string, + executionProfileDigest: `sha256:${string}`, + ): Promise { + return this.readTransaction(() => { + this.#assertWorkspaceStorageRootBinding(rootId); + return this.currentWorkspaceContinuationBoundarySync( + workspaceId, + workspaceEpochId, + rootId, + executionProfileDigest, + ); + }); + } + async #readManagedMutationEvidence( operationId: string, ): Promise< @@ -3051,6 +3186,7 @@ export class SqliteRuntimeStore source_prefix_digest, boundary_digest, boundary_json, + workspace_boundary_json, provider_projection_version, provider_replay_digest, target_session_id, @@ -3082,6 +3218,7 @@ export class SqliteRuntimeStore source_prefix_digest, boundary_digest, boundary_json, + workspace_boundary_json, provider_projection_version, provider_replay_digest, target_session_id, @@ -3099,13 +3236,48 @@ export class SqliteRuntimeStore .all() as unknown as ContinuationClaimStorageRow[]; } + private readContinuationClaimRowsForSession( + sessionId: string, + protocolVersion: 1 | 2, + ): ContinuationClaimStorageRow[] { + return this.db + .prepare(` + SELECT + claim_id, + source_session_id, + source_invocation_id, + source_run_id, + source_turn_id, + source_event_high_water, + source_prefix_digest, + boundary_digest, + boundary_json, + workspace_boundary_json, + provider_projection_version, + provider_replay_digest, + target_session_id, + target_invocation_id, + target_run_id, + target_turn_id, + target_run_header_json, + claimed_at, + start_event_id, + start_kind, + protocol_version + FROM runtime_continuation_claims + WHERE target_session_id = ? AND protocol_version = ? + ORDER BY claimed_at ASC, claim_id ASC + `) + .all(sessionId, protocolVersion) as unknown as ContinuationClaimStorageRow[]; + } + private assertContinuationAuthorityIntegrity(): void { for (const row of this.readContinuationClaimRows()) { - this.decodeContinuationClaimStateRow(row); + this.decodeAnyContinuationClaimStateRow(row); } } - private continuationTargetHasRuntimeState(claim: ContinuationClaimV1): boolean { + private continuationTargetHasRuntimeState(claim: ContinuationClaim): boolean { const { target } = claim; const values = [ target.invocationId, @@ -3142,6 +3314,28 @@ export class SqliteRuntimeStore private decodeContinuationClaimStateRow( row: ContinuationClaimStorageRow, ): ContinuationClaimStateV1 { + const state = this.decodeAnyContinuationClaimStateRow(row); + if (state.claim.protocol !== 'continuation_claim_v1') { + throw new Error(`Legacy continuation reader encountered ${state.claim.protocol}`); + } + return state as ContinuationClaimStateV1; + } + + private decodeWorkspaceBoundContinuationClaimStateRow( + row: ContinuationClaimStorageRow, + ): ContinuationClaimStateV2 { + const state = this.decodeAnyContinuationClaimStateRow(row); + if (state.claim.protocol !== 'continuation_claim_v2') { + throw new Error(`Workspace-bound continuation reader encountered ${state.claim.protocol}`); + } + return state as ContinuationClaimStateV2; + } + + private decodeAnyContinuationClaimStateRow(row: ContinuationClaimStorageRow): { + claim: ContinuationClaim; + startEventId?: string; + startKind?: 'runtime_admission' | 'claim_repair'; + } { const claim = decodeContinuationClaimRow(row); if (!row.start_event_id) { if (row.start_kind !== null) { @@ -3160,7 +3354,7 @@ export class SqliteRuntimeStore return { claim, startEventId: row.start_event_id, startKind: row.start_kind }; } - private assertContinuationBoundaryMatchesLedger(claim: ContinuationClaimV1): void { + private assertContinuationBoundaryMatchesLedger(claim: ContinuationClaim): void { const lastIndex = claim.boundary.segments.length - 1; for (const [index, segment] of claim.boundary.segments.entries()) { let prefix: ImmutableRuntimePrefixV1; @@ -3208,6 +3402,94 @@ export class SqliteRuntimeStore } } + private assertContinuationWorkspaceBoundaryMatchesAuthority( + boundary: ManagedWorkspaceContinuationBoundaryV1, + ): void { + const storageRoot = this.#readWorkspaceStorageRootBinding(); + if ( + !storageRoot || + storageRoot.protocol_version !== 1 || + storageRoot.root_id !== boundary.storageRootId + ) { + throw new Error('Continuation workspace boundary storage-root identity conflict'); + } + + const expected = this.currentWorkspaceContinuationBoundarySync( + boundary.workspaceId, + boundary.workspaceEpochId, + storageRoot.root_id, + boundary.executionProfileDigest, + ); + if (!expected || !isDeepStrictEqual(boundary, expected)) { + throw new Error('Continuation workspace boundary no longer matches accepted authority'); + } + } + + private currentWorkspaceContinuationBoundarySync( + workspaceId: string, + workspaceEpochId: string, + storageRootId: string, + executionProfileDigest: `sha256:${string}`, + ): ManagedWorkspaceContinuationBoundaryV1 | undefined { + const authority = this.readCanonicalWorkspaceAuthoritySync(); + this.assertWorkspaceProjectionsMatchSync(authority); + const epoch = authority.baselines.find( + (candidate) => + candidate.epoch.workspaceId === workspaceId && + candidate.epoch.workspaceEpochId === workspaceEpochId, + ); + const head = authority.heads.find( + (candidate) => + candidate.workspaceId === workspaceId && candidate.workspaceEpochId === workspaceEpochId, + ); + if (!epoch || !head) return undefined; + + const baseline = authority.baselines.find( + (candidate) => + candidate.baseline.workspaceVersionId === head.workspaceVersionId && + candidate.baselineAcceptedEventId === head.acceptedEventId, + ); + const successor = authority.successors.find( + (candidate) => + candidate.successor.workspaceVersionId === head.workspaceVersionId && + candidate.acceptedEventId === head.acceptedEventId, + ); + if ((baseline === undefined) === (successor === undefined)) { + throw new Error('Continuation workspace head has ambiguous accepted evidence'); + } + const version = baseline?.baseline ?? successor!.successor; + if ( + epoch.epoch.policyHash !== version.policyHash || + workspaceMutationPolicyHashV1( + epoch.epoch.materializationProfileDigest, + executionProfileDigest, + ) !== epoch.epoch.policyHash || + (successor !== undefined && + successor.successor.executionProfileDigest !== executionProfileDigest) + ) { + throw new Error('Continuation workspace boundary execution profile conflict'); + } + return { + protocol: 'managed_workspace_continuation_boundary_v1', + storageRootId, + repositoryId: epoch.epoch.repositoryId, + workspaceId: epoch.epoch.workspaceId, + workspaceEpochId: epoch.epoch.workspaceEpochId, + workspaceInstanceId: epoch.epoch.workspaceInstanceId, + workspaceVersionId: head.workspaceVersionId, + acceptedEventId: head.acceptedEventId, + revision: head.revision, + objectFormat: epoch.epoch.objectFormat, + sourceCommitOid: epoch.epoch.sourceCommitOid, + sourceTreeOid: epoch.epoch.sourceTreeOid, + commitOid: head.commitOid, + treeOid: head.treeOid, + materializationProfileDigest: epoch.epoch.materializationProfileDigest, + policyHash: version.policyHash, + executionProfileDigest, + }; + } + private assertToolLedgerTransition( candidateEvents: readonly RuntimeEvent[], expectedTransition: Parameters[0]['expectedTransition'], @@ -4037,7 +4319,7 @@ function assertRecoveryAuthorityCapability(db: DatabaseSync): void { } function assertContinuationStartEvent( - claim: ContinuationClaimV1, + claim: ContinuationClaim, event: RuntimeEvent, startKind: 'runtime_admission' | 'claim_repair', ): void { @@ -4099,6 +4381,19 @@ function assertContinuationAuthorityCapability(db: DatabaseSync): void { } } +function assertWorkspaceBoundContinuationAuthorityCapability(db: DatabaseSync): void { + const row = db + .prepare('SELECT version FROM runtime_capabilities WHERE capability = ?') + .get(RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY) as + | { version?: unknown } + | undefined; + if (row?.version !== RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY_VERSION) { + throw new Error( + `SQLite runtime workspace-bound continuation capability ${RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY}@${RUNTIME_WORKSPACE_BOUND_CONTINUATION_AUTHORITY_CAPABILITY_VERSION} is unavailable`, + ); + } +} + function assertRuntimeStorageSafeId(value: string, message: string): void { if (!isRuntimeStorageSafeId(value)) throw new Error(message); } @@ -4433,6 +4728,7 @@ interface ContinuationClaimStorageRow { source_prefix_digest: string; boundary_digest: string; boundary_json: string; + workspace_boundary_json: string | null; provider_projection_version: number; provider_replay_digest: string; target_session_id: string; @@ -4621,8 +4917,8 @@ function hasOnlyKeys(value: object, allowed: readonly string[]): boolean { return Object.keys(value).every((key) => allowedSet.has(key)); } -function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): ContinuationClaimV1 { - if (row.protocol_version !== 1) { +function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): ContinuationClaim { + if (row.protocol_version !== 1 && row.protocol_version !== 2) { throw new Error(`Unsupported continuation claim protocol ${row.protocol_version}`); } const boundary = JSON.parse(row.boundary_json) as unknown; @@ -4630,10 +4926,17 @@ function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): Continuat markPersisted(JSON.parse(row.target_run_header_json)), ); const claim = decodeContinuationClaim({ - protocol: 'continuation_claim_v1', + protocol: row.protocol_version === 1 ? 'continuation_claim_v1' : 'continuation_claim_v2', claimId: row.claim_id, boundaryDigest: row.boundary_digest, boundary, + ...(row.protocol_version === 2 + ? { + workspaceBoundary: row.workspace_boundary_json + ? (JSON.parse(row.workspace_boundary_json) as unknown) + : undefined, + } + : {}), providerProjectionVersion: row.provider_projection_version, providerReplayDigest: row.provider_replay_digest, target: { @@ -4645,6 +4948,14 @@ function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): Continuat targetRunHeader, claimedAt: row.claimed_at, }); + if ( + (row.protocol_version === 1 && row.workspace_boundary_json !== null) || + (row.protocol_version === 2 && + (claim.protocol !== 'continuation_claim_v2' || + row.workspace_boundary_json !== stableJsonStringify(claim.workspaceBoundary))) + ) { + throw new Error(`Continuation claim workspace row/payload mismatch for ${row.claim_id}`); + } const source = claim.boundary.segments.at(-1)!; if ( row.source_session_id !== source.identity.sessionId || diff --git a/packages/storage/src/workspace-version-authority-internal.ts b/packages/storage/src/workspace-version-authority-internal.ts index d779f9308e..db0db3b6dc 100644 --- a/packages/storage/src/workspace-version-authority-internal.ts +++ b/packages/storage/src/workspace-version-authority-internal.ts @@ -26,6 +26,10 @@ import type { WorkspaceVersionRecordV1, } from '@maka/core/workspace-version-authority'; import type { RuntimeEvent } from '@maka/core/runtime-event'; +import type { + ManagedWorkspaceContinuationBoundaryV1, + RuntimeBoundaryDigest, +} from '@maka/core/runtime-boundary'; type WorkspaceBaselineAuthorityWriter = ( input: WorkspaceBaselineAuthorityInput, @@ -123,6 +127,12 @@ type ManagedMutationReservationReader = ( type ManagedMutationEvidenceReader = ( operationId: string, ) => Promise; +type WorkspaceContinuationBoundaryReader = ( + workspaceId: string, + workspaceEpochId: string, + rootId: string, + executionProfileDigest: RuntimeBoundaryDigest, +) => Promise; interface WorkspaceBaselineAuthorityRegistration { readonly writer: WorkspaceBaselineAuthorityWriter; @@ -135,6 +145,7 @@ interface WorkspaceBaselineAuthorityRegistration { readonly readVersion: WorkspaceVersionReader; readonly readActiveManagedMutation: ManagedMutationReservationReader; readonly readManagedMutationEvidence: ManagedMutationEvidenceReader; + readonly readContinuationBoundary: WorkspaceContinuationBoundaryReader; readonly bindStorageRoot: WorkspaceStorageRootBinder; readonly adoptStorageRoot: WorkspaceStorageRootAdopter; boundRootId?: string; @@ -157,6 +168,7 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( readVersion: WorkspaceVersionReader, readActiveManagedMutation: ManagedMutationReservationReader, readManagedMutationEvidence: ManagedMutationEvidenceReader, + readContinuationBoundary: WorkspaceContinuationBoundaryReader, ): void { if (workspaceBaselineAuthorityWriters.has(store)) { throw new Error('Workspace baseline authority writer is already registered'); @@ -170,11 +182,31 @@ export function registerWorkspaceBaselineAuthorityWriterInternal( readVersion, readActiveManagedMutation, readManagedMutationEvidence, + readContinuationBoundary, bindStorageRoot, adoptStorageRoot, }); } +export function readWorkspaceContinuationBoundaryInternal( + store: object, + workspaceId: string, + workspaceEpochId: string, + executionProfileDigest: RuntimeBoundaryDigest, +): Promise { + const registration = workspaceBaselineAuthorityWriters.get(store); + if (!registration) throw new Error('Workspace continuation boundary reader is unavailable'); + if (!registration.boundRootId) { + throw new Error('Workspace continuation boundary store has no durable storage-root binding'); + } + return registration.readContinuationBoundary( + workspaceId, + workspaceEpochId, + registration.boundRootId, + executionProfileDigest, + ); +} + export function readActiveManagedMutationInternal( store: object, workspaceInstanceId: string,