From 01c67d4c705d7cc9e1486e3b48d2931bf1bebc21 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 20:02:26 +0200 Subject: [PATCH] fix: enforce addressed registry lifecycle invariants --- backend/src/app.ts | 8 +- backend/src/routes/workspaces.ts | 14 +- .../src/workspaces/registry-publication.ts | 98 ++++++++++---- backend/src/workspaces/registry.ts | 127 ++++++------------ 4 files changed, 124 insertions(+), 123 deletions(-) diff --git a/backend/src/app.ts b/backend/src/app.ts index 6fd937c9..f948453b 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -47,6 +47,7 @@ export interface BuildAppDeps { function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry { const registryConfig = config.workspaceRegistry; + mkdirSync(registryConfig.root, { recursive: true, mode: 0o700 }); const repository = new GitWorkspaceRepository(registryConfig); const sessionsRoot = join(registryConfig.root, "sessions"); mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 }); @@ -56,9 +57,10 @@ function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry { provisionedWorkspaceMode: 0o700, }); const hash = (value: string) => createHash("sha256").update(value).digest("hex"); - const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(`${registryConfig.remoteUrl ?? ""}:${registryConfig.branch}`) }; - const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: repositoryIdentity.digest }; - const request = { kind: "bootstrap" as const, installation: { installationId: registryConfig.installationId, digest: hash(registryConfig.installationId) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] }; + const volumeIdentity = realpathSync(registryConfig.root); + const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(JSON.stringify({ schemaVersion: 1, volume: volumeIdentity, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) }; + const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: hash(JSON.stringify({ schemaVersion: 1, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) }; + const request = { kind: "bootstrap" as const, operation: "registry_bootstrap" as const, installation: { installationId: registryConfig.installationId, volume: volumeIdentity, digest: hash(JSON.stringify({ schemaVersion: 1, installationId: registryConfig.installationId, volume: volumeIdentity })) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] }; return new WorkspaceRegistry({ rootLeaseFactory, lifecycleOwner: new CapabilityAwareRegistryPublicationLifecycleOwner(), participants: [], synchronizers: [], repository, installationIdentity: { operation: "registry_bootstrap", requestSha256: canonicalBootstrapRequestDigest(request) as never, installationIdentitySha256: request.installation.digest as never, repositoryIdentitySha256: repositoryIdentity.digest as never, remoteRefIdentitySha256: remoteIdentity.digest as never }, repositoryIdentity, remoteIdentity }); } diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 93b54ffe..104591f8 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -13,7 +13,7 @@ import { type PublishWorkspaceRequest, type WorkspaceRegistry, } from "../workspaces/registry.js"; -import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1, RegistryWorkspaceMutationV1 } from "../workspaces/registry-publication.js"; +import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js"; import type { Revision40 } from "../workspaces/workspace-lock-root-lease.js"; import { resolveRuntimeBindings } from "../workspaces/bindings.js"; import { buildInstallationContract, renderWorkspaceDocs } from "../workspaces/contracts.js"; @@ -320,8 +320,12 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); const revisions = snapshot.workspaces.map(item => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof deps.registry.snapshotPath === "function" ? deps.registry.snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` })); return await Promise.all(revisions.map(async revision => { - const { workspace } = await deps.registry.read(revision.id); - return { id: revision.id, name: revision.id, file: `${revision.id}.yaml`, displayName: workspace.workspace.name, description: workspace.workspace.description, language: workspace.workspace.language, workspace, revision }; + const pinned = typeof deps.registry.readPinned === "function" + ? await deps.registry.readPinned(revision.id, revision.commit) + : await deps.registry.read(revision.id); + const workspace = pinned.workspace; + const exactRevision = { ...revision, snapshotPath: "workspaceConfigPath" in pinned ? pinned.workspaceConfigPath : revision.snapshotPath }; + return { id: revision.id, name: revision.id, file: `${revision.id}.yaml`, displayName: workspace.workspace.name, description: workspace.workspace.description, language: workspace.workspace.language, workspace, revision: exactRevision }; })); } catch (error) { return errorReply(reply, error); } }); @@ -369,7 +373,6 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) const requestValue = publishRequest(request.body); const identity = recoveryIdentity(); const id = requestValue.action === "delete" ? requestValue.id : requestValue.workspace.workspace.id; - const mutation = requestValue as unknown as RegistryWorkspaceMutationV1; const addressed = { mode: "create" as const, operation: "registry_pull" as const, runId: addressedRunId(), requestSha256: sha256(JSON.stringify(requestValue)) as never, @@ -377,9 +380,8 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, expectedBaseCommit: requestValue.baseCommit as Revision40, - mutation, }; - const result = await deps.registry.publishAddressed(addressed); + const result = await deps.registry.publishAddressed(addressed, { authoring: requestValue }); return { revision: revisionFromPublication(result, id) }; } catch (error) { return errorReply(reply, error); diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index 530c5237..72be6bc1 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -16,7 +16,7 @@ interface PlanFields { readonly schemaVersion: 1; readonly installationIdentityS export interface RegistryBootstrapAddressedPlanV1 extends PlanFields { readonly operation: "registry_bootstrap"; readonly changedSetRule: "all_target_workspace_ids"; readonly baseCommit: null; readonly baseManifestSha256: null; readonly baseWorkspaces: readonly []; } export interface RegistryPullAddressedPlanV1 extends PlanFields { readonly operation: "registry_pull"; readonly changedSetRule: "symmetric_base_target_workspace_difference"; readonly baseCommit: Revision40; readonly baseManifestSha256: Sha256Hex; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; } export type RegistryAddressedPlanV1 = RegistryBootstrapAddressedPlanV1 | RegistryPullAddressedPlanV1; -interface StateFields { readonly schemaVersion: 1; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly phase: RegistryAddressedPublicationPhaseV1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly advertisedTargetCommit: Revision40 | null; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target` | null; readonly fetchedTargetCommit: Revision40 | null; readonly targetCommit: Revision40 | null; readonly targetManifestSha256: Sha256Hex | null; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[] | null; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[] | null; readonly planSha256: Sha256Hex | null; readonly changedSetSha256: Sha256Hex | null; readonly participantsSha256: Sha256Hex | null; readonly synchronizersSha256: Sha256Hex | null; readonly publicationIntentSha256: Sha256Hex | null; readonly publishedActiveStateSha256: Sha256Hex | null; readonly terminalResultSha256: Sha256Hex | null; readonly terminalResult: RegistryAddressedResultV1 | null; readonly priorStateSha256: Sha256Hex | null; readonly baseCommit: Revision40 | null; readonly baseManifestSha256: Sha256Hex | null; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "all_target_workspace_ids" | "symmetric_base_target_workspace_difference" | null; } +interface StateFields { readonly schemaVersion: 1; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly phase: RegistryAddressedPublicationPhaseV1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly advertisedTargetCommit: Revision40 | null; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target` | null; readonly fetchedTargetCommit: Revision40 | null; readonly targetCommit: Revision40 | null; readonly targetManifestSha256: Sha256Hex | null; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[] | null; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[] | null; readonly planSha256: Sha256Hex | null; readonly changedSetSha256: Sha256Hex | null; readonly participantsSha256: Sha256Hex | null; readonly synchronizersSha256: Sha256Hex | null; readonly publicationIntentSha256: Sha256Hex | null; readonly publishedActiveStateSha256: Sha256Hex | null; readonly terminalResultSha256: Sha256Hex | null; readonly priorStateSha256: Sha256Hex | null; readonly baseCommit: Revision40 | null; readonly baseManifestSha256: Sha256Hex | null; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "all_target_workspace_ids" | "symmetric_base_target_workspace_difference" | null; } export interface RegistryAddressedSnapshotV1 extends RegistryActiveSnapshotV1 { readonly requestDigest: string; } export interface RegistryInstallationIdentityV1 { readonly installationId: string; readonly digest: string; } export interface RegistryRepositoryIdentityV1 { readonly remote: string; readonly branch: string; readonly head: string; readonly digest: string; } @@ -32,7 +32,7 @@ export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: strin export interface CapabilityAwareRegistryPublicationSynchronizer { readonly synchronizerId: string; ensureForPublication(plan: RegistryAddressedPlanV1, capabilities: OrderedWorkspaceWriterCapabilitySet, phase: "planned" | "participants_prepared" | "publication_intent_durable" | "target_published"): Promise; } export class CapabilityAwareRegistryPublicationLifecycleOwner { constructor() {} - async run(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise; }): Promise { + async run(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise; readonly afterPublication?: (result: T) => Promise; }): Promise { // Reader gates are recursively nested so every changed workspace remains quiescent and // reader-exclusive until the publication callback, reconciliation, and terminal durability // have all settled. No participant or synchronizer is allowed to acquire a capability. @@ -45,6 +45,7 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner { const result = await input.action(); for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "target_published"); for (const item of prepared) await item.participant.reconcile(input.plan, item.lease, item.value, "target_published"); + if (input.afterPublication) await input.afterPublication(result); return result; } const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId; @@ -69,15 +70,10 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner { } } -export type RegistryWorkspaceMutationV1 = - | { readonly action: "create"; readonly workspace: CanonicalWorkspace; readonly baseCommit: Revision40 } - | { readonly action: "update"; readonly workspace: CanonicalWorkspace; readonly baseCommit: Revision40; readonly baseBlob: Revision40 } - | { readonly action: "delete"; readonly id: CanonicalWorkspaceId; readonly baseCommit: Revision40; readonly baseBlob: Revision40 }; - export type RegistryAddressedRequestV1 = | { readonly mode: "create"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: null; readonly remoteRefIdentitySha256: Sha256Hex } | { readonly mode: "resume"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex } - | { readonly mode: "create"; readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: Revision40; readonly remoteRefIdentitySha256: Sha256Hex; readonly mutation?: RegistryWorkspaceMutationV1 } + | { readonly mode: "create"; readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: Revision40; readonly remoteRefIdentitySha256: Sha256Hex } | { readonly mode: "resume"; readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex }; export interface RegistryBootstrapAddressedResultV1 { readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly plan: RegistryBootstrapAddressedPlanV1; readonly planSha256: Sha256Hex; readonly phase: "terminal_durable"; readonly publication: "target" | "reconciled_target" | "unchanged"; } export interface RegistryPullAddressedResultV1 { readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly plan: RegistryPullAddressedPlanV1; readonly planSha256: Sha256Hex; readonly phase: "terminal_durable"; readonly publication: "target" | "reconciled_target" | "unchanged"; } @@ -91,7 +87,10 @@ const PHASES: readonly RegistryAddressedPublicationPhaseV1[] = ["request_claimed const CONFLICT = () => Object.assign(new Error("preprocessing_conflict"), { code: "preprocessing_conflict" }); const canonical = (v: unknown): string => JSON.stringify(v, (_k, x) => x && typeof x === "object" && !Array.isArray(x) ? Object.fromEntries(Object.keys(x).sort().map(k => [k, x[k]])) : x); export const registryDigest = (value: unknown): string => createHash("sha256").update(canonical(value)).digest("hex"); -export function canonicalBootstrapRequestDigest(request: { readonly kind?: "bootstrap" | "publish"; readonly operation?: RegistryAddressedOperationV1; readonly installation?: unknown; readonly repository?: unknown; readonly remote?: unknown; readonly workspaceIds?: readonly string[] }): string { return registryDigest({ kind: request.kind ?? (request.operation === "registry_pull" ? "publish" : "bootstrap"), operation: request.operation, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...(request.workspaceIds ?? [])].sort() }); } +export function canonicalBootstrapRequestDigest(request: { readonly kind?: "bootstrap" | "publish"; readonly operation?: RegistryAddressedOperationV1; readonly installation?: unknown; readonly repository?: unknown; readonly remote?: unknown; readonly workspaceIds?: readonly string[] }): string { + const operation = request.operation ?? (request.kind === "publish" ? "registry_pull" : "registry_bootstrap"); + return registryDigest({ schemaVersion: 1, operation, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...(request.workspaceIds ?? [])].sort() }); +} export function addressedRunId(): RegistryRunId32 { return randomBytes(16).toString("hex") as RegistryRunId32; } function failIfBadIdentity(s: StateFields, runId: string): void { if (s.schemaVersion !== 1 || s.runId !== runId || !RUN.test(s.runId) || s.jobArtifactPath !== `addressed-publication-jobs/${s.runId}.json` || !SHA.test(s.requestSha256) || !SHA.test(s.installationIdentitySha256) || !SHA.test(s.repositoryIdentitySha256) || !SHA.test(s.remoteRefIdentitySha256) || !PHASES.includes(s.phase)) throw CONFLICT(); } function immutable(s: StateFields): unknown { const { phase: _p, priorStateSha256: _h, ...rest } = s; return rest; } @@ -108,31 +107,52 @@ export class RegistryAddressedPublicationStore { private path(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.json`); } private async dirs(): Promise { await mkdir(this.jobsDirectory, { recursive: true, mode: 0o700 }); const st = await lstat(this.jobsDirectory); if (st.isSymbolicLink() || !st.isDirectory() || (Number((st as any).mode) & 0o777) !== 0o700 || Number((st as any).uid) !== (process.getuid?.() ?? Number((st as any).uid))) throw CONFLICT(); } private async durable(path: string, value: unknown, exclusive = false): Promise { const name = path.split("/").pop()!; const tmp = join(this.jobsDirectory, `.${name}.tmp`); if (exclusive) { try { await stat(path); throw CONFLICT(); } catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw CONFLICT(); } } const bytes = `${canonical(value)}\n`; try { const h = await open(tmp, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | (fsConstants.O_NOFOLLOW ?? 0), 0o600); try { await h.writeFile(bytes); await h.sync(); } finally { await h.close(); } const st = await stat(tmp); if (!ownerMode(st, 0o600)) throw CONFLICT(); if (!exclusive) { const current = await lstat(path); if (current.isSymbolicLink() || !ownerMode(current, 0o600)) throw CONFLICT(); } await rename(tmp, path); await fsyncParent(path); } catch (e) { await rm(tmp, { force: true }).catch(() => undefined); if (exclusive && (e as NodeJS.ErrnoException)?.code === "EEXIST") throw CONFLICT(); throw CONFLICT(); } } - async claim(request: RegistryAddressedRequestV1 | { readonly kind: "bootstrap" | "publish"; readonly installation: { readonly digest: string }; readonly repository: { readonly digest: string }; readonly remote: { readonly digest: string }; readonly requestDigest: string }, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise { + async claim(request: RegistryAddressedRequestV1 | { readonly kind: "bootstrap" | "publish"; readonly installation: { readonly digest: string }; readonly repository: { readonly digest: string }; readonly remote: { readonly digest: string }; readonly requestDigest: string }, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId()), base?: { readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] }): Promise { await this.dirs(); - if (!("operation" in request)) request = { mode: "create", operation: request.kind === "publish" ? "registry_pull" : "registry_bootstrap", runId, requestSha256: request.requestDigest as Sha256Hex, installationIdentitySha256: request.installation.digest as Sha256Hex, repositoryIdentitySha256: request.repository.digest as Sha256Hex, remoteRefIdentitySha256: request.remote.digest as Sha256Hex, expectedBaseCommit: request.kind === "publish" ? "0".repeat(40) as Revision40 : null } as RegistryAddressedRequestV1; + if (!("operation" in request)) { + const legacyPublish = request.kind === "publish"; + request = { mode: "create", operation: legacyPublish ? "registry_pull" : "registry_bootstrap", runId, requestSha256: request.requestDigest as Sha256Hex, installationIdentitySha256: request.installation.digest as Sha256Hex, repositoryIdentitySha256: request.repository.digest as Sha256Hex, remoteRefIdentitySha256: request.remote.digest as Sha256Hex, expectedBaseCommit: legacyPublish ? "0".repeat(40) as Revision40 : null } as RegistryAddressedRequestV1; + if (legacyPublish && !base) base = { manifestSha256: "0".repeat(64) as Sha256Hex, workspaces: [] }; + } if (!RUN.test(runId)) throw CONFLICT(); + if (request.operation === "registry_pull" && !base) throw CONFLICT(); const now = new Date().toISOString(); - const state = { schemaVersion: 1, runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation: request.operation, installationIdentitySha256: request.installationIdentitySha256, repositoryIdentitySha256: request.repositoryIdentitySha256, remoteRefIdentitySha256: request.remoteRefIdentitySha256, advertisedTargetCommit: null, immutableTargetRef: null, fetchedTargetCommit: null, targetCommit: null, targetManifestSha256: null, targetWorkspaces: null, changedWorkspaceIds: null, planSha256: null, changedSetSha256: null, participantsSha256: null, synchronizersSha256: null, publicationIntentSha256: null, publishedActiveStateSha256: null, terminalResultSha256: null, terminalResult: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null, baseManifestSha256: null, baseWorkspaces: [], changedSetRule: null }; + const state = { schemaVersion: 1, runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation: request.operation, installationIdentitySha256: request.installationIdentitySha256, repositoryIdentitySha256: request.repositoryIdentitySha256, remoteRefIdentitySha256: request.remoteRefIdentitySha256, advertisedTargetCommit: null, immutableTargetRef: null, fetchedTargetCommit: null, targetCommit: null, targetManifestSha256: null, targetWorkspaces: null, changedWorkspaceIds: null, planSha256: null, changedSetSha256: null, participantsSha256: null, synchronizersSha256: null, publicationIntentSha256: null, publishedActiveStateSha256: null, terminalResultSha256: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null, baseManifestSha256: base?.manifestSha256 ?? null, baseWorkspaces: base?.workspaces ?? [], changedSetRule: null }; await this.durable(this.path(runId), state, true); return state as unknown as RegistryAddressedPublicationStateV1; } - async setBase(runId: RegistryRunId32, baseManifestSha256: Sha256Hex, baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]): Promise { - const old = await this.read(runId); - if (old.phase !== "request_claimed" || old.operation !== "registry_pull") throw CONFLICT(); - const next = { ...old, baseManifestSha256, baseWorkspaces, priorStateSha256: registryDigest(old) as Sha256Hex }; - await this.durable(this.path(runId), next); - return next; - } async storeTerminalResult(runId: RegistryRunId32, result: RegistryAddressedResultV1): Promise { const state = await this.read(runId); if (state.phase !== "target_published") throw CONFLICT(); - await this.durable(this.path(runId), { ...state, terminalResult: result, terminalResultSha256: registryDigest(result) as Sha256Hex, priorStateSha256: registryDigest(state) as Sha256Hex }); + // The terminal result is a deterministic projection of the immutable plan. Persist only + // its digest in the fixed state shape; recovery reconstructs and verifies the projection. + const expected = registryDigest(result) as Sha256Hex; + await this.transition(runId, "terminal_durable", { terminalResultSha256: expected }); + } + private planFromState(state: RegistryAddressedPublicationStateV1): RegistryAddressedPlanV1 { + if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces || !state.changedWorkspaceIds || !state.changedSetSha256 || !state.advertisedTargetCommit || !state.immutableTargetRef || !state.fetchedTargetCommit || !state.planSha256) throw CONFLICT(); + const common = { schemaVersion: 1 as const, installationIdentitySha256: state.installationIdentitySha256, + repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256, + jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit, + immutableTargetRef: state.immutableTargetRef, fetchedTargetCommit: state.fetchedTargetCommit, + targetCommit: state.targetCommit, targetManifestSha256: state.targetManifestSha256, + targetWorkspaces: state.targetWorkspaces, changedWorkspaceIds: state.changedWorkspaceIds, + changedSetSha256: state.changedSetSha256 }; + const plan = state.operation === "registry_bootstrap" + ? { ...common, operation: "registry_bootstrap" as const, changedSetRule: "all_target_workspace_ids" as const, baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] as const } + : { ...common, operation: "registry_pull" as const, changedSetRule: "symmetric_base_target_workspace_difference" as const, + baseCommit: state.baseCommit!, baseManifestSha256: state.baseManifestSha256!, baseWorkspaces: state.baseWorkspaces }; + if (registryDigest(plan) !== state.planSha256) throw CONFLICT(); + return plan; } async readTerminalResult(runId: RegistryRunId32, expectedDigest: Sha256Hex): Promise { const state = await this.read(runId); - if (state.phase !== "terminal_durable" || !state.terminalResult || registryDigest(state.terminalResult) !== expectedDigest) throw CONFLICT(); - return state.terminalResult; + if (state.phase !== "terminal_durable" || state.terminalResultSha256 !== expectedDigest) throw CONFLICT(); + const plan = await this.planFromState(state); + const result = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, + plan, planSha256: state.planSha256!, phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1; + if (registryDigest(result) !== expectedDigest) throw CONFLICT(); + return result; } async read(runId: RegistryRunId32): Promise { @@ -141,24 +161,36 @@ export class RegistryAddressedPublicationStore { if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw CONFLICT(); const s = parsed as StateFields & { operation?: unknown }; const keys = Object.keys(parsed).sort(); - const expected = ["schemaVersion","runId","requestSha256","jobArtifactPath","phase","operation","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256","advertisedTargetCommit","immutableTargetRef","fetchedTargetCommit","targetCommit","targetManifestSha256","targetWorkspaces","changedWorkspaceIds","planSha256","changedSetSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","terminalResult","priorStateSha256","baseCommit","baseManifestSha256","baseWorkspaces","changedSetRule"].sort(); + const expected = ["schemaVersion","runId","requestSha256","jobArtifactPath","phase","operation","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256","advertisedTargetCommit","immutableTargetRef","fetchedTargetCommit","targetCommit","targetManifestSha256","targetWorkspaces","changedWorkspaceIds","planSha256","changedSetSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256","baseCommit","baseManifestSha256","baseWorkspaces","changedSetRule"].sort(); if (keys.length !== expected.length || keys.some((v, i) => v !== expected[i])) throw CONFLICT(); failIfBadIdentity(s, runId); + if (s.priorStateSha256 !== null && !SHA.test(s.priorStateSha256)) throw CONFLICT(); if (s.operation !== "registry_bootstrap" && s.operation !== "registry_pull") throw CONFLICT(); for (const key of ["advertisedTargetCommit","fetchedTargetCommit","targetCommit"] as const) if (s[key] !== null && !REV.test(s[key])) throw CONFLICT(); for (const key of ["requestSha256","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256"] as const) if (!SHA.test(s[key])) throw CONFLICT(); - if (s.terminalResult !== null && (typeof s.terminalResult !== "object" || Array.isArray(s.terminalResult))) throw CONFLICT(); for (const key of ["targetManifestSha256","changedSetSha256","planSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256"] as const) if (s[key] !== null && !SHA.test(s[key])) throw CONFLICT(); if (s.immutableTargetRef !== null && s.immutableTargetRef !== `refs/thoth/addressed-runs/${runId}/target`) throw CONFLICT(); if (s.operation === "registry_bootstrap" && (s.baseCommit !== null || s.baseManifestSha256 !== null || s.baseWorkspaces.length !== 0 || (s.changedSetRule !== null && s.changedSetRule !== "all_target_workspace_ids"))) throw CONFLICT(); - if (s.operation === "registry_pull" && (s.baseCommit === null || !REV.test(s.baseCommit) || (s.baseManifestSha256 !== null && !SHA.test(s.baseManifestSha256)) || s.changedSetRule !== null && s.changedSetRule !== "symmetric_base_target_workspace_difference")) throw CONFLICT(); + if (s.operation === "registry_pull" && (s.baseCommit === null || !REV.test(s.baseCommit) || s.baseManifestSha256 === null || !SHA.test(s.baseManifestSha256) || s.changedSetRule !== null && s.changedSetRule !== "symmetric_base_target_workspace_difference")) throw CONFLICT(); if (s.targetWorkspaces !== null && !Array.isArray(s.targetWorkspaces) || s.changedWorkspaceIds !== null && !Array.isArray(s.changedWorkspaceIds)) throw CONFLICT(); const required: Record = { request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], target_fetched: ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit"], planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"], }; for (const phase of PHASES.slice(0, PHASES.indexOf(s.phase) + 1)) for (const key of required[phase]) if (s[key] === null || s[key] === undefined) throw CONFLICT(); - if (s.phase === "terminal_durable" && (!s.terminalResult || registryDigest(s.terminalResult) !== s.terminalResultSha256)) throw CONFLICT(); + const introduced: Record = { + request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], target_fetched: ["fetchedTargetCommit"], + planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], + participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"], + }; + const currentIndex = PHASES.indexOf(s.phase); + let chainPrevious: string | null = null; + for (let index = 0; index <= currentIndex; index++) { + const candidate: Record = { ...s, phase: PHASES[index]!, priorStateSha256: chainPrevious }; + for (const later of PHASES.slice(index + 1)) for (const key of introduced[later]!) candidate[key] = null; + if (index === currentIndex && s.priorStateSha256 !== chainPrevious) throw CONFLICT(); + chainPrevious = registryDigest(candidate); + } return s as RegistryAddressedPublicationStateV1; } async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial = {}): Promise { @@ -171,8 +203,20 @@ export class RegistryAddressedPublicationStore { const from = PHASES.indexOf(old.phase); const to = PHASES.indexOf(phase); if (to !== from + 1) throw CONFLICT(); - const allowed = new Set(["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "terminalResult", "changedSetRule", "baseManifestSha256", "baseWorkspaces"]); + const allowed = new Set(["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "changedSetRule", "baseManifestSha256", "baseWorkspaces"]); for (const key of Object.keys(patch)) if (!allowed.has(key)) throw CONFLICT(); + // All durable facts are append-only. In particular the advertised OID/ref and base + // identity may never be replaced during resume or by a competing caller. + for (const key of ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "baseManifestSha256", "baseWorkspaces", "changedSetRule"] as const) { + if (key in patch && old[key] !== null && old[key] !== undefined && JSON.stringify(patch[key]) !== JSON.stringify(old[key])) throw CONFLICT(); + } + const phaseFields: Record = { + request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], + target_fetched: ["fetchedTargetCommit"], planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], + participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"], + target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"], + }; + if (Object.keys(patch).some(key => !phaseFields[phase].includes(key))) throw CONFLICT(); for (const key of ["schemaVersion", "runId", "requestSha256", "jobArtifactPath", "installationIdentitySha256", "repositoryIdentitySha256", "remoteRefIdentitySha256", "operation", "baseCommit"] as const) { if (key in patch && patch[key] !== old[key]) throw CONFLICT(); } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 6a80de4f..c1218160 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -157,45 +157,6 @@ export class WorkspaceRegistry { /** Automatic addressed recovery. The repository lock is held for selection and execution. */ recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; } - private async bootstrap(): Promise { - await this.repository.ensureLayout(); - return this.lock.run(async () => { try { const status = await this.repository.bootstrap(); await this.materialize(status.head!); await this.publishMaterialized(status.head!); return status; } catch (error) { return this.gitFallback(error); } }); - } - private async pull(): Promise { - await this.repository.ensureLayout(); - return this.lock.run(async () => { try { const status = await this.repository.pull(); await this.materialize(status.head!); await this.publishMaterialized(status.head!); return status; } catch (error) { return this.gitFallback(error); } }); - } - private async list(): Promise { - const active = await this.tryActiveState(); - if (active) return active.revisions; - await this.bootstrap(); - return (await this.activeState()).revisions; - } - private async activate(commit: string): Promise { await this.materialize(commit); await this.publishMaterialized(commit); } - /** Test-only migration seam; production callers use publishAddressed. */ - private async publish(request: PublishWorkspaceRequest): Promise { return this.publishWorkspace(request); } - private async publishWorkspace(request: PublishWorkspaceRequest): Promise { - await this.repository.ensureLayout(); - return this.lock.run(() => this.publishWorkspaceLocked(request)); - } - private async publishWorkspaceLocked(request: PublishWorkspaceRequest): Promise { - const status = await this.repository.pull(); await this.materialize(status.head!); await this.publishMaterialized(status.head!); - const current = await this.activeState(); const id = request.action === "delete" ? request.id : request.workspace.workspace.id; - const existing = current.revisions.find(revision => revision.id === id); const local = request.action === "delete" ? undefined : request.workspace; - if (request.baseCommit !== status.head || (request.action !== "create" && existing?.blob !== request.baseBlob)) { - if (request.action !== "create" && request.baseCommit !== status.head && existing?.blob === request.baseBlob) throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale"); - throw await this.conflictFor(request, status.head!, existing, local); - } - if (request.action === "create" && existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "create" && !existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "delete") await this.assertEvidenceContext(request.workspace, status.head!); - const yamlPath = workspacePath(id); const docs = this.documentationPaths(id); - if (request.action === "delete") { await this.repository.removeRegistryFile(yamlPath); await this.repository.removeRegistryFile(docs.contract); await this.repository.removeRegistryFile(docs.readme); } - else { const source = serializeWorkspaceYaml(request.workspace); const rendered = renderWorkspaceDocs(request.workspace); await this.repository.writeRegistryFile(yamlPath, source); await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown); } - const next = await this.repository.commitAndPush([yamlPath, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`); - await this.materialize(next.head!); await this.publishMaterialized(next.head!); return (await this.activeState()).revisions.find(revision => revision.id === id); - } - async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise { await this.repository.ensureLayout(); return this.lock.run(async () => this.ensureBootstrapAddressedLocked(identity)); @@ -226,23 +187,21 @@ export class WorkspaceRegistry { return { kind: "bootstrap_terminal", result: result as Extract, snapshot: this.addressedSnapshot(published, identity.requestSha256) }; } - async publishAddressed(request: RegistryAddressedRequestV1): Promise { + async publishAddressed(request: RegistryAddressedRequestV1, context?: { readonly authoring?: PublishWorkspaceRequest }): Promise { await this.repository.ensureLayout(); return this.lock.run(async () => { - // Authoring keeps the accepted HTTP CRUD payload, but crosses the same addressed - // boundary as every other publication. The mutation is deliberately handled while - // repository.lock is held; callers never receive the historical publish API. - if ("mutation" in request && request.mutation) return this.publishWorkspaceAddressed(request as Extract & { readonly mutation: PublishWorkspaceRequest }); const store = new RegistryAddressedPublicationStore(this.repository.root); let state: RegistryAddressedPublicationStateV1; try { state = await store.read(request.runId); } catch { - state = await store.claim(request, request.runId); + let base: ActiveState | undefined; if (request.operation === "registry_pull" && request.mode === "create") { - const base = await this.snapshotState(request.expectedBaseCommit); - const baseWorkspaces = base.revisions.map(revision => this.manifestIdentity(revision)); - state = await store.setBase(state.runId, registryDigest(base) as never, baseWorkspaces); + base = await this.snapshotState(request.expectedBaseCommit); } + state = await store.claim(request, request.runId, base ? { + manifestSha256: registryDigest(base) as never, + workspaces: base.revisions.map(revision => this.manifestIdentity(revision)), + } : undefined); } const identity: RegistryBootstrapRecoveryIdentityV1 = { operation: "registry_bootstrap", requestSha256: request.requestSha256, @@ -251,43 +210,30 @@ export class WorkspaceRegistry { remoteRefIdentitySha256: request.remoteRefIdentitySha256, }; if (state.operation !== request.operation || state.requestSha256 !== request.requestSha256 || state.installationIdentitySha256 !== request.installationIdentitySha256 || state.repositoryIdentitySha256 !== request.repositoryIdentitySha256 || state.remoteRefIdentitySha256 !== request.remoteRefIdentitySha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed request identity does not match the durable job"); + if (request.operation === "registry_pull" && request.mode === "create" && state.baseCommit !== request.expectedBaseCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed base does not match the durable job"); + if (context?.authoring && state.phase !== "request_claimed") throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Author mutation cannot be replayed after claim"); + if (context?.authoring) await this.applyAuthorMutation(context.authoring); return this.executeAddressedLocked(store, state, identity); }); } - private async publishWorkspaceAddressed(request: Extract & { readonly mutation: PublishWorkspaceRequest }): Promise { - const mutation = request.mutation; - await this.publishWorkspaceLocked(mutation); - const targetState = await this.activeState(); - const targetWorkspaces = targetState.revisions.map(revision => this.manifestIdentity(revision)); - const base = request.expectedBaseCommit ? await this.snapshotState(request.expectedBaseCommit) : undefined; - const baseWorkspaces = base?.revisions.map(revision => this.manifestIdentity(revision)) ?? []; - const baseIds = new Set(baseWorkspaces.map(item => item.workspaceId)); - const targetIds = new Set(targetWorkspaces.map(item => item.workspaceId)); - const changedWorkspaceIds = [...new Set([...baseIds, ...targetIds])].filter(id => - !baseIds.has(id) || !targetIds.has(id) - || registryDigest(baseWorkspaces.find(item => item.workspaceId === id)) !== registryDigest(targetWorkspaces.find(item => item.workspaceId === id)), - ).sort() as CanonicalWorkspaceId[]; - const plan: RegistryPullAddressedPlanV1 = { - schemaVersion: 1, operation: "registry_pull", - installationIdentitySha256: request.installationIdentitySha256, - repositoryIdentitySha256: request.repositoryIdentitySha256, - remoteRefIdentitySha256: request.remoteRefIdentitySha256, - jobArtifactPath: `addressed-publication-jobs/${request.runId}.json`, - advertisedTargetCommit: targetState.head as Revision40, - immutableTargetRef: `refs/thoth/addressed-runs/${request.runId}/target`, - fetchedTargetCommit: targetState.head as Revision40, - targetCommit: targetState.head as Revision40, - targetManifestSha256: registryDigest(targetState) as never, - targetWorkspaces, - changedWorkspaceIds, - changedSetSha256: registryDigest(changedWorkspaceIds) as never, - changedSetRule: "symmetric_base_target_workspace_difference", - baseCommit: request.expectedBaseCommit ?? targetState.head as Revision40, - baseManifestSha256: registryDigest(base) as never, - baseWorkspaces, - }; - return { operation: "registry_pull", runId: request.runId, jobArtifactPath: plan.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable", publication: "target" }; + private async applyAuthorMutation(request: PublishWorkspaceRequest): Promise { + const status = await this.repository.pull(); + const current = await this.tryActiveState(); + const id = request.action === "delete" ? request.id : request.workspace.workspace.id; + const existing = current?.revisions.find(revision => revision.id === id); + const local = request.action === "delete" ? undefined : request.workspace; + if (request.baseCommit !== status.head || (request.action !== "create" && existing?.blob !== request.baseBlob)) { + if (request.action !== "create" && request.baseCommit !== status.head && existing?.blob === request.baseBlob) throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale"); + throw await this.conflictFor(request, status.head!, existing, local); + } + if (request.action === "create" && existing) throw await this.conflictFor(request, status.head!, existing, local); + if (request.action !== "create" && !existing) throw await this.conflictFor(request, status.head!, existing, local); + if (request.action !== "delete") await this.assertEvidenceContext(request.workspace, status.head!); + const yamlPath = workspacePath(id); const docs = this.documentationPaths(id); + if (request.action === "delete") { await this.repository.removeRegistryFile(yamlPath); await this.repository.removeRegistryFile(docs.contract); await this.repository.removeRegistryFile(docs.readme); } + else { await this.repository.writeRegistryFile(yamlPath, serializeWorkspaceYaml(request.workspace)); const rendered = renderWorkspaceDocs(request.workspace); await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown); } + await this.repository.commitAndPush([yamlPath, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`); } private async executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise { @@ -334,7 +280,8 @@ export class WorkspaceRegistry { const plan = await this.planForState(state); const leases = []; for (const id of plan.changedWorkspaceIds) leases.push(await this.rootLeaseFactory.acquireOrProvision(this.rootLeaseFactory.canonicalInput(id))); - const result = await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => { + let output: RegistryAddressedResultV1 | undefined; + await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => { if (state.phase === "planned") { state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(this.participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(this.synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never }); } @@ -342,14 +289,20 @@ export class WorkspaceRegistry { state = await store.transition(state.runId, "publication_intent_durable", { publicationIntentSha256: registryDigest({ runId: state.runId, planSha256: state.planSha256 }) as never }); } if (state.phase === "publication_intent_durable") { - await this.publishMaterialized(target!); + await this.publishSnapshotPointer(target!); state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.activeState()) as never }); } return undefined; + }, afterPublication: async () => { + // Terminal durability is deliberately the final operation while every writer, root, + // quiescence, and reader-exclusive gate is still owned by this callback. + if (state.phase === "target_published") { + output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1; + await store.storeTerminalResult(state.runId, output); + state = await store.read(state.runId); + } }})); - const output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan), phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1; - await store.storeTerminalResult(state.runId, output); - state = await store.transition(state.runId, "terminal_durable", { terminalResultSha256: registryDigest(output) as never }); + if (!output) output = await store.readTerminalResult(state.runId, state.terminalResultSha256!); return output; } @@ -711,7 +664,7 @@ export class WorkspaceRegistry { } - private async publishMaterialized(commit: string): Promise { + private async publishSnapshotPointer(commit: string): Promise { const state = await this.snapshotState(commit); await this.writeActiveState({ head: state.head, revisions: state.revisions }); }