From 7741e610453ec89b71e359f94ddfaaabc110f79d Mon Sep 17 00:00:00 2001 From: mptyl Date: Wed, 12 Aug 2026 01:21:05 +0200 Subject: [PATCH] fix registry active manifest digest identity --- .../src/workspaces/registry-publication.ts | 17 ++++++++++++- backend/src/workspaces/registry.ts | 25 +++++++++++-------- ...ace-registry-addressed-publication.test.ts | 14 +++++++++++ 3 files changed, 45 insertions(+), 11 deletions(-) diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index 8b023901..5b500969 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -85,6 +85,14 @@ 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"); +/** Digest of the active registry identity, independent of snapshot storage details. */ +export function registryManifestDigest( + commit: Revision40, + workspaces: readonly RegistryWorkspaceManifestIdentityV1[], +): Sha256Hex { + const sorted = [...workspaces].sort((left, right) => left.workspaceId < right.workspaceId ? -1 : left.workspaceId > right.workspaceId ? 1 : 0); + return registryDigest({ commit, workspaces: sorted }) as Sha256Hex; +} 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() }); @@ -127,6 +135,8 @@ export class RegistryAddressedPublicationStore { } 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(); + if (registryManifestDigest(state.targetCommit, state.targetWorkspaces) !== state.targetManifestSha256) throw CONFLICT(); + if (state.operation === "registry_pull" && (!state.baseCommit || registryManifestDigest(state.baseCommit, state.baseWorkspaces) !== state.baseManifestSha256)) throw CONFLICT(); const common = { schemaVersion: 1 as const, installationIdentitySha256: state.installationIdentitySha256, repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit, @@ -234,7 +244,12 @@ export class RegistryAddressedPublicationStore { await this.durable(this.path(runId), next); return next; } - snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryActiveSnapshotV1 { if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces) throw CONFLICT(); return { schemaVersion: 1, commit: state.targetCommit, manifestSha256: state.targetManifestSha256, workspaces: state.targetWorkspaces }; } + snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryActiveSnapshotV1 { + if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces) throw CONFLICT(); + const manifestSha256 = registryManifestDigest(state.targetCommit, state.targetWorkspaces); + if (manifestSha256 !== state.targetManifestSha256) throw CONFLICT(); + return { schemaVersion: 1, commit: state.targetCommit, manifestSha256, workspaces: [...state.targetWorkspaces].sort((left, right) => left.workspaceId < right.workspaceId ? -1 : left.workspaceId > right.workspaceId ? 1 : 0) }; + } async scan(): Promise { await this.dirs(); let entries: string[]; try { entries = (await readdir(this.jobsDirectory)).sort(); } catch { throw CONFLICT(); } if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT(); for (const name of entries.filter(x => /^\.[0-9a-f]{32}\.json\.tmp$/.test(x))) { const path = join(this.jobsDirectory, name); const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes) throw CONFLICT(); await unlink(path); await fsyncParent(path); } entries = (await readdir(this.jobsDirectory)).sort(); diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 0a36617e..b5b0f64d 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -20,7 +20,7 @@ import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js"; import { VerifiedWorkspaceLockRootLeaseFactory, type Revision40, type CanonicalWorkspaceId } from "./workspace-lock-root-lease.js"; import { WorkspaceFsAtV1 } from "./workspace-fs-at.js"; import { runUnderOrderedWorkspaceWriterLocks, type OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js"; -import { RegistryAddressedPublicationStore, addressedRunId, canonicalBootstrapRequestDigest, registryDigest, CapabilityAwareRegistryPublicationLifecycleOwner, reconcilePreparedPublication, type CapabilityAwareRegistryPublicationParticipant, type CapabilityAwareRegistryPublicationSynchronizer, type RegistryAddressedPlanV1, type RegistryPullAddressedPlanV1, type RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.js"; +import { RegistryAddressedPublicationStore, addressedRunId, canonicalBootstrapRequestDigest, registryDigest, registryManifestDigest, CapabilityAwareRegistryPublicationLifecycleOwner, reconcilePreparedPublication, type CapabilityAwareRegistryPublicationParticipant, type CapabilityAwareRegistryPublicationSynchronizer, type RegistryAddressedPlanV1, type RegistryPullAddressedPlanV1, type RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.js"; export type { GitStatus } from "./git-repository.js"; export interface WorkspaceRevision { @@ -326,7 +326,7 @@ export class WorkspaceRegistry { if (existing.length === 1) return existing[0]!; let base: ActiveState | undefined; if (request.operation === "registry_pull") base = await this.#snapshotState(request.expectedBaseCommit); - return store.claim(request, base ? { commit: base.head as Revision40, manifestSha256: registryDigest(base) as never, workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)) } : null); + return store.claim(request, base ? { commit: base.head as Revision40, manifestSha256: this.#manifestDigest(base), workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)) } : null); }); } @@ -344,7 +344,7 @@ export class WorkspaceRegistry { // claim is an exclusive O_CREAT|O_EXCL operation: an existing same-ID run is // always a create conflict, never an implicit replay. state = await store.claim(request, base ? { - commit: base.head as Revision40, manifestSha256: registryDigest(base) as never, + commit: base.head as Revision40, manifestSha256: this.#manifestDigest(base), workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)), } : null); } @@ -391,7 +391,7 @@ export class WorkspaceRegistry { // target-published job whose active pointer has drifted. if (state.phase === "target_published") { const active = await this.#tryActiveState(); - if (!active || active.head !== target || registryDigest(active) !== state.targetManifestSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry target drifted"); + if (!active || active.head !== target || this.#manifestDigest(active) !== state.targetManifestSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry target drifted"); } // Once advertised, the run-specific ref is immutable evidence. Every resume after // the fetch barrier revalidates it before reading or publishing any bytes. @@ -411,11 +411,11 @@ export class WorkspaceRegistry { const plan = operation === "registry_bootstrap" ? { schemaVersion: 1, operation, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit!, immutableTargetRef: state.immutableTargetRef!, fetchedTargetCommit: target as Revision40, targetCommit: target as Revision40, - targetManifestSha256: registryDigest(targetState) as never, targetWorkspaces, changedWorkspaceIds: targetWorkspaces.map(x => x.workspaceId).sort() as CanonicalWorkspaceId[], changedSetSha256: registryDigest(targetWorkspaces.map(x => x.workspaceId).sort()) as never, changedSetRule: "all_target_workspace_ids" as const, baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] as const, + targetManifestSha256: this.#manifestDigest(targetState), targetWorkspaces, changedWorkspaceIds: targetWorkspaces.map(x => x.workspaceId).sort() as CanonicalWorkspaceId[], changedSetSha256: registryDigest(targetWorkspaces.map(x => x.workspaceId).sort()) as never, changedSetRule: "all_target_workspace_ids" as const, baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] as const, } : { schemaVersion: 1, operation, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit!, immutableTargetRef: state.immutableTargetRef!, fetchedTargetCommit: target as Revision40, targetCommit: target as Revision40, - targetManifestSha256: registryDigest(targetState), targetWorkspaces, changedWorkspaceIds, changedSetSha256: registryDigest(changedWorkspaceIds) as never, changedSetRule: "symmetric_base_target_workspace_difference" as const, baseCommit: state.baseCommit!, baseManifestSha256: registryDigest(base) as never, baseWorkspaces, + targetManifestSha256: this.#manifestDigest(targetState), targetWorkspaces, changedWorkspaceIds, changedSetSha256: registryDigest(changedWorkspaceIds) as never, changedSetRule: "symmetric_base_target_workspace_difference" as const, baseCommit: state.baseCommit!, baseManifestSha256: this.#manifestDigest(base!), baseWorkspaces, }; state = await store.transition(state.runId, "planned", { targetCommit: target as Revision40, targetManifestSha256: plan.targetManifestSha256 as never, targetWorkspaces: plan.targetWorkspaces, changedWorkspaceIds: plan.changedWorkspaceIds, changedSetSha256: plan.changedSetSha256 as never, planSha256: registryDigest(plan) as never, changedSetRule: plan.changedSetRule }); } @@ -438,19 +438,19 @@ export class WorkspaceRegistry { } if (state.phase === "publication_intent_durable") { const activeBefore = await this.#tryActiveState(); - const activeDigest = activeBefore ? registryDigest(activeBefore) : undefined; + const activeDigest = activeBefore ? this.#manifestDigest(activeBefore) : undefined; const baseAllowed = activeBefore === undefined || (operation === "registry_pull" && activeBefore.head === state.baseCommit && activeDigest === state.baseManifestSha256); const targetAllowed = activeBefore !== undefined && activeBefore.head === target && activeDigest === state.targetManifestSha256; if (!baseAllowed && !targetAllowed) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry identity conflicts with addressed publication"); if (targetAllowed) publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target"; else await this.#publishSnapshotPointer(target!); - state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.#activeState()) as never }); + state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: this.#manifestDigest(await this.#activeState()) as never }); } if (state.phase === "target_published") { // A restart at this barrier must verify the exact active target before success. const active = await this.#tryActiveState(); - if (!active || active.head !== target || registryDigest(active) !== state.targetManifestSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry target drifted"); + if (!active || active.head !== target || this.#manifestDigest(active) !== state.targetManifestSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry target drifted"); if (resumedAtTargetPublished) publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target"; for (const synchronizer of registryContext(this).synchronizers) await synchronizer.ensureForPublication(plan, capabilities, "target_published"); await reconcilePreparedPublication(registryContext(this).lifecycleOwner, plan); @@ -467,15 +467,20 @@ export class WorkspaceRegistry { async #planForState(state: RegistryAddressedPublicationStateV1): Promise { if (!state.targetCommit || !state.targetWorkspaces || !state.changedWorkspaceIds || !state.planSha256 || !state.changedSetSha256 || !state.targetManifestSha256 || !state.advertisedTargetCommit || !state.immutableTargetRef || !state.fetchedTargetCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + if (registryManifestDigest(state.targetCommit, state.targetWorkspaces) !== state.targetManifestSha256 || (state.operation === "registry_pull" && (!state.baseCommit || registryManifestDigest(state.baseCommit, state.baseWorkspaces) !== state.baseManifestSha256))) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); 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: RegistryAddressedPlanV1 = state.operation === "registry_bootstrap" ? { ...common, operation: "registry_bootstrap", changedSetRule: "all_target_workspace_ids", baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] } : { ...common, operation: "registry_pull", changedSetRule: "symmetric_base_target_workspace_difference", baseCommit: state.baseCommit!, baseManifestSha256: state.baseManifestSha256!, baseWorkspaces: state.baseWorkspaces }; if (registryDigest(plan) !== state.planSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); return plan; } + #manifestDigest(state: ActiveState): RegistryActiveSnapshotV1["manifestSha256"] { + return registryManifestDigest(state.head as Revision40, state.revisions.map(revision => this.#manifestIdentity(revision))); + } + #addressedSnapshot(state: ActiveState, requestDigest: string): RegistryActiveSnapshotV1 { const ids = state.revisions.map(revision => revision.id).sort(); - return { schemaVersion: 1, commit: state.head as RegistryActiveSnapshotV1["commit"], manifestSha256: registryDigest(state) as RegistryActiveSnapshotV1["manifestSha256"], workspaces: ids.map(id => { const revision = state.revisions.find(candidate => candidate.id === id)!; return this.#manifestIdentity(revision); }) }; + return { schemaVersion: 1, commit: state.head as RegistryActiveSnapshotV1["commit"], manifestSha256: this.#manifestDigest(state), workspaces: ids.map(id => { const revision = state.revisions.find(candidate => candidate.id === id)!; return this.#manifestIdentity(revision); }) }; } async #listRetainedSnapshots(): Promise { diff --git a/backend/test/workspace-registry-addressed-publication.test.ts b/backend/test/workspace-registry-addressed-publication.test.ts index 8f9d0e2d..8f0f14cc 100644 --- a/backend/test/workspace-registry-addressed-publication.test.ts +++ b/backend/test/workspace-registry-addressed-publication.test.ts @@ -3,10 +3,24 @@ import { mkdtemp } from "node:fs/promises"; import { join } from "node:path"; import { RegistryAddressedPublicationStore, + registryManifestDigest, type RegistryAddressedRequestV1, } from "../src/workspaces/registry-publication.js"; describe("addressed publication", () => { + it("derives one order-independent digest from the exact commit and workspace identities", () => { + const identities = [ + { workspaceId: "workspace-z", revision: "a".repeat(40), descriptorBlob: "b".repeat(40), manifestSha256: "c".repeat(64) }, + { workspaceId: "workspace-a", revision: "d".repeat(40), descriptorBlob: "e".repeat(40), manifestSha256: "f".repeat(64) }, + ] as const; + expect(registryManifestDigest("1".repeat(40) as any, identities)).toBe( + registryManifestDigest("1".repeat(40) as any, [...identities].reverse()), + ); + expect(registryManifestDigest("2".repeat(40) as any, identities)).not.toBe( + registryManifestDigest("1".repeat(40) as any, identities), + ); + }); + it("claims a complete pull context and advances only through legal phases", async () => { const root = await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-pub-")); const runId = "a".repeat(32);