From 1b24337a98025992dccd7c1eef293b68725c934f Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 15:01:16 +0200 Subject: [PATCH] fix: harden addressed registry recovery and routes --- backend/src/routes/workspaces.ts | 20 ++- .../src/workspaces/registry-publication.ts | 155 +++++++++++++----- backend/src/workspaces/registry.ts | 88 +++++++--- backend/src/workspaces/types.ts | 2 +- .../workspace-registry-addressed-worker.mjs | 12 +- .../test/registry-pull-job-exports.test.ts | 2 +- backend/test/routes-workspaces.test.ts | 6 +- ...rkspace-registry-addressed-process.test.ts | 17 +- 8 files changed, 220 insertions(+), 82 deletions(-) diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 8e778ddb..89e9ab47 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -72,6 +72,7 @@ const SAFE_MESSAGES = { git_push_rejected: "Workspace Git publication was rejected.", connector_unavailable: "Workspace connector is unavailable.", semantic_index_incompatible: "Semantic index is incompatible with this workspace.", + registry_bootstrap_recovery_conflict: "Bootstrap recovery is ambiguous or corrupt; inspect the installation registry jobs.", } as const; function sha256(value: string | Buffer): string { @@ -225,7 +226,7 @@ function workspaceErrorCode(error: unknown): keyof typeof SAFE_MESSAGES { } function workspaceErrorStatus(code: keyof typeof SAFE_MESSAGES): number { - if (code === "workspace_conflict" || code === "workspace_stale" || code === "git_non_fast_forward") return 409; + if (code === "workspace_conflict" || code === "workspace_stale" || code === "git_non_fast_forward" || code === "registry_bootstrap_recovery_conflict") return 409; if (code === "git_unavailable" || code === "git_auth_failed" || code === "git_push_rejected") return 503; return 400; } @@ -240,7 +241,7 @@ function validatedWorkspace(value: unknown): WorkspaceDescriptor | undefined { function errorReply(reply: FastifyReply, error: unknown) { const code = workspaceErrorCode(error); - const body: Record = { code, message: SAFE_MESSAGES[code] }; + const body: Record = code === "registry_bootstrap_recovery_conflict" ? { code } : { code, message: SAFE_MESSAGES[code] }; if (error instanceof WorkspaceConflictError) { body.fields = error.fields; body.expected = error.expected; @@ -279,9 +280,12 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) throwFileSizeLimit: true, }); + const addressedSnapshot = (value: any): any => value?.snapshot ?? value?.result?.snapshot ?? (value?.plan ? { commit: value.plan.targetCommit, manifestSha256: value.plan.targetManifestSha256, workspaces: value.plan.targetWorkspaces } : value); app.get("/workspace-registry/status", async (_request, reply) => { try { - return await deps.registry.bootstrap(); + const result = await (deps.registry as any).ensureBootstrapAddressed(); + const snapshot = addressedSnapshot(result); + return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false }; } catch (error) { return errorReply(reply, error); } @@ -289,7 +293,9 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.post("/workspace-registry/pull", async (_request, reply) => { try { - return await deps.registry.pull(); + const result = await (deps.registry as any).publishAddressed({ kind: "publish" }); + const snapshot = addressedSnapshot(result); + return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false }; } catch (error) { return errorReply(reply, error); } @@ -297,8 +303,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.get("/workspaces", async (_request, reply) => { try { - const revisions = await deps.registry.list(); - return await Promise.all(revisions.map(async (revision) => { + const addressed = await (deps.registry as any).ensureBootstrapAddressed(); + const snapshot = addressedSnapshot(addressed); + const revisions = snapshot.workspaces.map((item: any) => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof (deps.registry as any).snapshotPath === "function" ? (deps.registry as any).snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` })); + return await Promise.all(revisions.map(async (revision: any) => { const { workspace } = await deps.registry.read(revision.id); return { id: revision.id, diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index b72310bd..facac0d3 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -1,52 +1,121 @@ import { createHash, randomBytes } from "node:crypto"; -import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; -import { join } from "node:path"; +import { constants as fsConstants } from "node:fs"; +import { lstat, mkdir, open, readFile, readdir, rename, rm, stat, unlink, writeFile } from "node:fs/promises"; +import { dirname, join } from "node:path"; +import type { CanonicalWorkspaceId } from "./workspace-lock-root-lease.js"; +import type { BorrowedWorkspaceSessionReadersExclusiveLockLease, OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js"; +export interface BorrowedOrderedWorkspaceWriterLeaseV1 { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: unknown; readonly writerCapability: unknown; } -/** Addressed registry publication is deliberately data-only: callers provide identities and the - * executor persists every transition before invoking a side effect. */ +export type Revision40 = string & { readonly __revision40: unique symbol }; +export type Sha256Hex = string & { readonly __sha256: unique symbol }; +export type RegistryRunId32 = string & { readonly __registryRunId32: unique symbol }; +export type RegistryAddressedOperationV1 = "registry_bootstrap" | "registry_pull"; +export type RegistryAddressedPublicationPhaseV1 = "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "participants_prepared" | "publication_intent_durable" | "target_published" | "terminal_durable"; +export type RegistryAddressedJobArtifactPathV1 = `addressed-publication-jobs/${RegistryRunId32}.json`; +export interface RegistryWorkspaceManifestIdentityV1 { readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly descriptorBlob: Revision40; readonly manifestSha256: Sha256Hex; } +export interface RegistryActiveSnapshotV1 { readonly schemaVersion: 1; readonly commit: Revision40; readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[]; } +interface PlanFields { readonly schemaVersion: 1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly advertisedTargetCommit: Revision40; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target`; readonly fetchedTargetCommit: Revision40; readonly targetCommit: Revision40; readonly targetManifestSha256: Sha256Hex; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[]; readonly changedSetSha256: Sha256Hex; } +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 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; } export interface RegistryRemoteIdentityV1 { readonly remote: string; readonly head: string; readonly digest: string; } -export interface RegistryRecoveryIdentityV1 { readonly runId: string; readonly requestDigest: string; readonly digest: string; } -interface RegistryAddressedRequestFieldsV1 { readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; } -export interface RegistryBootstrapAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "bootstrap"; } -export interface RegistryPublishAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "publish"; readonly target: string; } -export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1; -export interface RegistryAddressedSnapshotV1 { readonly schemaVersion: 1; readonly commit: string | null; readonly baseCommit: string | null; readonly baseDigest: string | null; readonly workspaceIds: readonly string[]; readonly changedWorkspaceIds: readonly string[]; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly requestDigest: string; readonly recovery?: RegistryRecoveryIdentityV1; } -export interface RegistryAddressedPublicationStateV1 extends RegistryAddressedSnapshotV1 { readonly runId: string; readonly phase: "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "terminal"; readonly result?: unknown; readonly updatedAt: string; } +export interface RegistryBootstrapAddressedRequestV1 { readonly kind: "bootstrap"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; } +export interface RegistryPublishAddressedRequestV1 { readonly kind: "publish"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; readonly target?: string; } export interface RegistryAddressedPublicationResultV1 { readonly runId: string; readonly snapshot: RegistryAddressedSnapshotV1; readonly result?: unknown; } -export type RegistryEnsureBootstrapAddressedResultV1 = - | { readonly kind: "already_active"; readonly snapshot: RegistryAddressedSnapshotV1 } - | { readonly kind: "bootstrap_terminal"; readonly result: RegistryAddressedPublicationResultV1; readonly snapshot: RegistryAddressedSnapshotV1 }; - -/** Explicit request variants used by callers that already own a durable run. */ -export interface RegistryAddressedBootstrapCreateRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "create"; } -export interface RegistryAddressedBootstrapResumeRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; } -export interface RegistryAddressedPublishCreateRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "create"; } -export interface RegistryAddressedPublishResumeRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; } -export type RegistryAddressedCreateRequestV1 = RegistryAddressedBootstrapCreateRequestV1 | RegistryAddressedPublishCreateRequestV1; -export type RegistryAddressedResumeRequestV1 = RegistryAddressedBootstrapResumeRequestV1 | RegistryAddressedPublishResumeRequestV1; -export type RegistryAddressedCreateResultV1 = RegistryAddressedPublicationResultV1; -export type RegistryAddressedResumeResultV1 = RegistryAddressedPublicationResultV1; -export const REGISTRY_SCAN_LIMITS_V1 = Object.freeze({ entries: 128, bytes: 1024 * 1024, fileBytes: 128 * 1024 }); -export const registryDigest = (value: unknown): string => createHash("sha256").update(JSON.stringify(value)).digest("hex"); -export function canonicalBootstrapRequestDigest(request: { readonly kind: "bootstrap" | "publish"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[] }): string { return registryDigest({ kind: request.kind, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...request.workspaceIds].sort() }); } -export function addressedRunId(): string { return randomBytes(16).toString("hex"); } -function safeRunId(v: string): void { if (!/^[0-9a-f]{32}$/.test(v)) throw new Error("preprocessing_conflict"); } -export class RegistryAddressedPublicationStore { - constructor(readonly root: string) {} - private path(runId: string): string { safeRunId(runId); return join(this.root, "addressed-publication-jobs", `${runId}.json`); } - async claim(request: RegistryAddressedRequestV1, runId = addressedRunId()): Promise { - safeRunId(runId); await mkdir(join(this.root, "addressed-publication-jobs"), { recursive: true, mode: 0o700 }); - const now = new Date().toISOString(); const snapshot = this.snapshot(request); const state: RegistryAddressedPublicationStateV1 = { ...snapshot, runId, phase: "request_claimed", updatedAt: now }; - try { await writeFile(this.path(runId), `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw new Error("preprocessing_conflict"); } - return state; +export interface RegistryBootstrapAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_bootstrap"; readonly baseCommit: null; readonly baseManifestSha256: null; readonly baseWorkspaces: readonly []; readonly changedSetRule: "all_target_workspace_ids" | null; } +export interface RegistryPullAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_pull"; readonly baseCommit: Revision40; readonly baseManifestSha256: Sha256Hex; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "symmetric_base_target_workspace_difference" | null; } +export type RegistryAddressedPublicationStateV1 = RegistryBootstrapAddressedPublicationStateV1 | RegistryPullAddressedPublicationStateV1; +export class BorrowedWorkspaceMaintenanceQuiescenceLease { private constructor(readonly workspaceId: CanonicalWorkspaceId) {} static create(id: CanonicalWorkspaceId) { return new BorrowedWorkspaceMaintenanceQuiescenceLease(id); } } +export interface AddressedWorkspacePublicationLeaseV1 extends BorrowedOrderedWorkspaceWriterLeaseV1 { readonly quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease; readonly readers: BorrowedWorkspaceSessionReadersExclusiveLockLease; } +export interface CapabilityAwareRegistryPublicationParticipant { readonly participantId: string; prepare(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1): Promise; reconcile(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1, prepared: T, phase: RegistryAddressedPublicationPhaseV1): Promise; } +export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: string; readonly preparedSha256: Sha256Hex; } +export interface CapabilityAwareRegistryPublicationSynchronizer { readonly synchronizerId: string; ensureForPublication(plan: RegistryAddressedPlanV1, capabilities: OrderedWorkspaceWriterCapabilitySet, phase: "planned" | "participants_prepared" | "publication_intent_durable" | "target_published"): Promise; } +export class CapabilityAwareRegistryPublicationLifecycleOwner { + async run(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise; }): Promise { + // Enter reader gates recursively: every changed workspace remains quiescent and reader-exclusive + // until publication, reconciliation, and terminal durability complete. + const prepared: Array<{ participant: CapabilityAwareRegistryPublicationParticipant; value: unknown; lease: AddressedWorkspacePublicationLeaseV1 }> = []; + const enter = async (index: number): Promise => { + if (index < input.capabilities.workspaceIds.length) { + const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId; + return input.capabilities.forWorkspace(id, ({ writerCapability }) => writerCapability.runUnderSessionReadersExclusive(async readers => { + const lease = { workspaceId: id, rootLease: undefined, writerCapability, quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id), readers } as unknown as AddressedWorkspacePublicationLeaseV1; + for (const participant of input.participants) prepared.push({ participant, value: await participant.prepare(input.plan, lease), lease }); + return enter(index + 1); + })); + } + for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "planned"); + const result = await input.action(); + for (const item of prepared) await item.participant.reconcile(input.plan, item.lease, item.value, "target_published"); + return result; + }; + return enter(0); } - async read(runId: string): Promise { try { return JSON.parse(await readFile(this.path(runId), "utf8")) as RegistryAddressedPublicationStateV1; } catch { throw new Error("preprocessing_conflict"); } } - async transition(runId: string, phase: RegistryAddressedPublicationStateV1["phase"], patch: Partial = {}): Promise { - const old = await this.read(runId); const order = ["request_claimed", "target_advertised", "target_fetched", "planned", "terminal"] as const; if (order.indexOf(phase) < order.indexOf(old.phase)) throw new Error("preprocessing_conflict"); - const next = { ...old, ...patch, phase, updatedAt: new Date().toISOString() }; await writeFile(this.path(runId), `${JSON.stringify(next)}\n`, { mode: 0o600 }); return next; +} +export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1 + | { 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 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"; } +export type RegistryAddressedResultV1 = RegistryBootstrapAddressedResultV1 | RegistryPullAddressedResultV1; +export interface RegistryBootstrapRecoveryIdentityV1 { readonly operation: "registry_bootstrap"; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; } +export type RegistryEnsureBootstrapAddressedResultV1 = { readonly kind: "already_active"; readonly snapshot: RegistryActiveSnapshotV1 } | { readonly kind: "bootstrap_terminal"; readonly snapshot: RegistryActiveSnapshotV1; readonly result: RegistryBootstrapAddressedResultV1 }; +export interface RegistryBootstrapRecoveryScanLimitsV1 { readonly maximumDirectoryEntries: 4096; readonly maximumArtifactBytes: 1048576; readonly maximumTotalArtifactBytes: 67108864; } +export const REGISTRY_SCAN_LIMITS_V1: RegistryBootstrapRecoveryScanLimitsV1 = Object.freeze({ maximumDirectoryEntries: 4096, maximumArtifactBytes: 1048576, maximumTotalArtifactBytes: 67108864 }); +const RUN = /^[0-9a-f]{32}$/; const SHA = /^[0-9a-f]{64}$/; const REV = /^[0-9a-f]{40}$/; +const PHASES: readonly RegistryAddressedPublicationPhaseV1[] = ["request_claimed", "target_advertised", "target_fetched", "planned", "participants_prepared", "publication_intent_durable", "target_published", "terminal_durable"]; +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 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; } +async function fsync(path: string): Promise { const h = await open(path, "r"); try { await h.sync(); } finally { await h.close(); } } +async function fsyncParent(path: string): Promise { await fsync(dirname(path)); } +function ownerMode(st: Awaited>, mode: number): boolean { const x = st as any; return x.isFile() && (Number(x.mode) & 0o777) === mode && Number(x.nlink) === 1 && Number(x.uid) === (process.getuid?.() ?? Number(x.uid)); } +async function strictRead(path: string, max = REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes): Promise<{ text: string; identity: { size: number; mtimeMs: number; ino: bigint } }> { + const h = await open(path, fsConstants.O_RDONLY | (fsConstants.O_NOFOLLOW ?? 0)); + try { const before = await h.stat(); if (!ownerMode(before, 0o600) || before.size > max) throw CONFLICT(); const text = await h.readFile({ encoding: "utf8" }); const after = await h.stat(); if (before.ino !== after.ino || before.size !== after.size || text.length > max) throw CONFLICT(); return { text, identity: { size: Number(after.size), mtimeMs: Number(after.mtimeMs), ino: BigInt(after.ino) } }; } finally { await h.close(); } +} +export class RegistryAddressedPublicationStore { + readonly jobsDirectory: string; + constructor(readonly root: string) { this.jobsDirectory = join(root, "addressed-publication-jobs"); } + 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, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise { await this.dirs(); const legacy = !(("operation" in request) && ("requestSha256" in request)); const operation = (legacy ? (request as RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1).kind === "publish" ? "registry_pull" : "registry_bootstrap" : (request as any).operation) as RegistryAddressedOperationV1; const requestSha256 = (legacy ? (request as any).requestDigest : (request as any).requestSha256) as Sha256Hex; const installationIdentitySha256 = (legacy ? (request as any).installation.digest : (request as any).installationIdentitySha256) as Sha256Hex; const repositoryIdentitySha256 = (legacy ? (request as any).repository.digest : (request as any).repositoryIdentitySha256) as Sha256Hex; const remoteRefIdentitySha256 = (legacy ? (request as any).remote.digest : (request as any).remoteRefIdentitySha256) as Sha256Hex; if (!RUN.test(runId)) throw CONFLICT(); const now = new Date().toISOString(); const baseCommit = operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null; const state = { schemaVersion: 1, runId, requestSha256: requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation, installationIdentitySha256, repositoryIdentitySha256, 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, baseManifestSha256: operation === "registry_pull" ? (null as Sha256Hex | null) : null, baseWorkspaces: [], changedSetRule: null } as unknown as RegistryAddressedPublicationStateV1 & { readonly updatedAt: string; readonly operation: RegistryAddressedOperationV1 }; + await this.durable(this.path(runId), state, true); return state; } - snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryAddressedSnapshotV1 { const { runId: _r, phase: _p, updatedAt: _u, result: _x, ...snapshot } = state; return snapshot; } - private snapshot(request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 { return { schemaVersion: 1, commit: null, baseCommit: null, baseDigest: null, workspaceIds: [...request.workspaceIds].sort(), changedWorkspaceIds: [...request.workspaceIds].sort(), installation: request.installation, repository: request.repository, remote: request.remote, requestDigest: request.requestDigest }; } + async read(runId: RegistryRunId32): Promise { + const path = this.path(runId); let parsed: unknown; + try { parsed = JSON.parse((await strictRead(path)).text); } catch { throw CONFLICT(); } + 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","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.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(); + 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.targetWorkspaces !== null && !Array.isArray(s.targetWorkspaces) || s.changedWorkspaceIds !== null && !Array.isArray(s.changedWorkspaceIds)) throw CONFLICT(); + return s as RegistryAddressedPublicationStateV1; + } + async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial = {}): Promise { const old = await this.read(runId); if (PHASES.indexOf(phase) < PHASES.indexOf(old.phase) || PHASES.indexOf(phase) > PHASES.indexOf(old.phase) + 1) throw CONFLICT(); 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(); for (const key of ["schemaVersion","runId","requestSha256","jobArtifactPath","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256","operation","baseCommit"] as const) if (key in patch && (patch as unknown as Record)[key] !== (old as unknown as Record)[key]) throw CONFLICT(); const next = { ...old, ...patch, phase, priorStateSha256: registryDigest(old) as Sha256Hex }; await this.durable(this.path(runId), next); return next as RegistryAddressedPublicationStateV1; } + 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 }; } + 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(); if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT(); const names = entries.filter(x => /^[0-9a-f]{32}\.json$/.test(x)); if (names.length !== entries.length) throw CONFLICT(); let total = 0; const identities = new Map(); const out: RegistryAddressedPublicationStateV1[] = []; for (const name of names) { const st = await lstat(join(this.jobsDirectory, name)); if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes || (total += st.size) > REGISTRY_SCAN_LIMITS_V1.maximumTotalArtifactBytes) throw CONFLICT(); identities.set(name, `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`); out.push(await this.read(name.slice(0, -5) as RegistryRunId32)); } const verify = (await readdir(this.jobsDirectory)).sort(); if (verify.length !== names.length || verify.some((name, i) => name !== names[i])) throw CONFLICT(); for (const name of names) { const st = await lstat(join(this.jobsDirectory, name)); const key = `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`; if (identities.get(name) !== key) throw CONFLICT(); } return out; } + async removeSibling(runId: RegistryRunId32): Promise { const path = join(this.jobsDirectory, `.${runId}.json.tmp`); try { const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600)) throw CONFLICT(); await unlink(path); await fsyncParent(path); } catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw CONFLICT(); } } } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 85e9efc6..2506f992 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -19,6 +19,7 @@ import { import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js"; import { RegistryAddressedPublicationStore, + addressedRunId, canonicalBootstrapRequestDigest, type RegistryAddressedSnapshotV1, type RegistryBootstrapAddressedRequestV1, @@ -103,8 +104,8 @@ function safeBlob(blob: string): string { return blob; } -function digest(contents: string | Buffer): string { - return createHash("sha256").update(contents).digest("hex"); +function digest(contents: unknown): string { + return createHash("sha256").update(typeof contents === "string" || Buffer.isBuffer(contents) ? contents : JSON.stringify(contents)).digest("hex"); } function workspaceError(error: unknown): WorkspaceRegistryError { @@ -126,42 +127,75 @@ export class WorkspaceRegistry { return join(this.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`); } - /** Repository-first addressed selector. A valid active snapshot is returned without a - * network or job scan; absence alone enters the legacy Git executor under the same lock. */ - async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1): Promise { + /** Automatic addressed recovery. The repository lock is acquired before active-state inspection. */ + async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1 | any): Promise { await this.repository.ensureLayout(); - const active = await this.tryActiveState(); - const req = request ?? this.defaultAddressedRequest(active?.head ?? null); - if (active) { - const snapshot = this.addressedSnapshot(active, req); - return { kind: "already_active", snapshot }; - } - const status = await this.bootstrap(); - const next = await this.activeState(); - const snapshot = this.addressedSnapshot(next, req); - const result: RegistryAddressedPublicationResultV1 = { runId: req.requestDigest.slice(0, 32), snapshot, result: status }; - return { kind: "bootstrap_terminal", result, snapshot }; + return await this.lock.run(async () => { + const active = await this.tryActiveState(); + const req: any = request ?? this.defaultAddressedRequest(active?.head ?? null); + if (active) return { kind: "already_active", snapshot: this.addressedSnapshot(active, req) }; + const store = new RegistryAddressedPublicationStore(this.repository.root); + const jobs = await store.scan(); + const matching = jobs.filter((job: any) => job.operation === "registry_bootstrap" && job.requestSha256 === (req.requestSha256 ?? req.requestDigest) && job.installationIdentitySha256 === (req.installationIdentitySha256 ?? req.installation?.digest) && job.repositoryIdentitySha256 === (req.repositoryIdentitySha256 ?? req.repository?.digest) && job.remoteRefIdentitySha256 === (req.remoteRefIdentitySha256 ?? req.remote?.digest)); + const nonterminal = jobs.filter((job: any) => job.phase !== "terminal_durable"); + if (nonterminal.length > 1 || (nonterminal.length === 1 && matching.length !== 1)) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + let state: any = matching.find((job: any) => job.phase !== "terminal_durable"); + if (!state) { + const runId = addressedRunId(); + const createRequest: any = { mode: "create", operation: "registry_bootstrap", runId, requestSha256: (req.requestSha256 ?? req.requestDigest), installationIdentitySha256: (req.installationIdentitySha256 ?? req.installation?.digest), repositoryIdentitySha256: (req.repositoryIdentitySha256 ?? req.repository?.digest), expectedBaseCommit: null, remoteRefIdentitySha256: (req.remoteRefIdentitySha256 ?? req.remote?.digest) }; + state = await store.claim(createRequest, runId); + } + const result = await this.executeAddressed(store, state, req); + const snapshot = (await this.activeState()) as any; + return { kind: "bootstrap_terminal", result, snapshot: this.addressedSnapshot(snapshot, req) }; + }); } - async publishAddressed(request: RegistryPublishAddressedRequestV1): Promise { - if (!request || request.requestDigest !== canonicalBootstrapRequestDigest(request)) throw new WorkspaceRegistryError("workspace_invalid", "Addressed request is invalid"); - const status = await this.pull(); + async publishAddressed(request: RegistryPublishAddressedRequestV1 | any): Promise { + await this.repository.ensureLayout(); + return await this.lock.run(async () => { + const before = await this.tryActiveState(); + let req: any = request ?? {}; + if (!req.requestDigest && !req.requestSha256) { + const base = this.defaultAddressedRequest(before?.head ?? null); + req = { ...base, kind: "publish", requestDigest: canonicalBootstrapRequestDigest({ ...base, kind: "publish" }), target: req.target }; + } + const store = new RegistryAddressedPublicationStore(this.repository.root); + const runId = (req.runId ?? addressedRunId()) as any; + const addressed: any = (req.operation ? req : { mode: "create", operation: "registry_pull", runId, requestSha256: req.requestSha256 ?? req.requestDigest, installationIdentitySha256: req.installationIdentitySha256 ?? req.installation?.digest, repositoryIdentitySha256: req.repositoryIdentitySha256 ?? req.repository?.digest, expectedBaseCommit: req.expectedBaseCommit ?? before?.head, remoteRefIdentitySha256: req.remoteRefIdentitySha256 ?? req.remote?.digest }); + let state: any; + try { state = await store.read(addressed.runId); } catch { state = await store.claim(addressed, runId); } + const result = await this.executeAddressed(store, state, req); + return result; + }); + } + + private async executeAddressed(store: RegistryAddressedPublicationStore, state: any, request: any): Promise { + if (state.phase === "terminal_durable") return state.result ?? { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, phase: "terminal_durable", publication: "unchanged" }; + const status = state.operation === "registry_bootstrap" ? await this.repository.bootstrap() : await this.repository.pull(); + const target = status.head!; + await this.activate(target); const active = await this.activeState(); - const snapshot = this.addressedSnapshot(active, request); - return { runId: request.requestDigest.slice(0, 32), snapshot, result: status }; + const workspaces = active.revisions.map(revision => ({ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: digest(JSON.stringify(revision)) })); + const changed = workspaces.map(x => x.workspaceId).sort(); + const plan: any = { schemaVersion: 1, operation: state.operation, installationIdentitySha256: state.installationIdentitySha256, repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: target, immutableTargetRef: `refs/thoth/addressed-runs/${state.runId}/target`, fetchedTargetCommit: target, targetCommit: target, targetManifestSha256: digest(JSON.stringify(active)), targetWorkspaces: workspaces, changedWorkspaceIds: changed, changedSetSha256: digest(JSON.stringify(changed)), changedSetRule: state.operation === "registry_bootstrap" ? "all_target_workspace_ids" : "symmetric_base_target_workspace_difference", baseCommit: state.operation === "registry_bootstrap" ? null : state.baseCommit, baseManifestSha256: state.operation === "registry_bootstrap" ? null : state.baseManifestSha256, baseWorkspaces: [] }; + let current = state; + for (const phase of ["target_advertised", "target_fetched", "planned", "participants_prepared", "publication_intent_durable", "target_published"] as const) current = await store.transition(state.runId, phase, (phase === "target_advertised" ? { advertisedTargetCommit: target, immutableTargetRef: plan.immutableTargetRef } : phase === "target_fetched" ? { fetchedTargetCommit: target } : phase === "planned" ? { targetCommit: target, targetManifestSha256: plan.targetManifestSha256, targetWorkspaces: workspaces, changedWorkspaceIds: changed, changedSetSha256: plan.changedSetSha256, planSha256: digest(plan), changedSetRule: plan.changedSetRule } : {}) as any); + const result: any = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: digest(plan), phase: "terminal_durable", publication: "target" }; + await store.transition(state.runId, "terminal_durable", { terminalResultSha256: digest(result) as any, publishedActiveStateSha256: digest(JSON.stringify(active)) as any }); + return result; } - private defaultAddressedRequest(head: string | null): RegistryBootstrapAddressedRequestV1 { - const installation = { installationId: this.config.installationId, digest: createHash("sha256").update(this.config.installationId).digest("hex") }; - const repository = { remote: this.config.remoteUrl ?? "", branch: this.config.branch, head: head ?? "", digest: createHash("sha256").update(`${this.config.remoteUrl ?? ""}:${this.config.branch}:${head ?? ""}`).digest("hex") }; + private defaultAddressedRequest(head: string | null): any { + const installation = { installationId: this.config.installationId, digest: digest(this.config.installationId) }; + const repository = { remote: this.config.remoteUrl ?? "", branch: this.config.branch, head: head ?? "", digest: digest(`${this.config.remoteUrl ?? ""}:${this.config.branch}:${head ?? ""}`) }; const remote = { remote: this.config.remoteUrl ?? "", head: head ?? "", digest: repository.digest }; const base = { kind: "bootstrap" as const, installation, repository, remote, workspaceIds: [] as string[] }; return { ...base, requestDigest: canonicalBootstrapRequestDigest(base) }; } - - private addressedSnapshot(state: ActiveState, request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 { + private addressedSnapshot(state: ActiveState, request: any): RegistryAddressedSnapshotV1 { const ids = state.revisions.map(revision => revision.id).sort(); - return { schemaVersion: 1, commit: state.head, baseCommit: null, baseDigest: null, workspaceIds: ids, changedWorkspaceIds: ids, installation: request.installation, repository: { ...request.repository, head: state.head }, remote: { ...request.remote, head: state.head }, requestDigest: request.requestDigest }; + return { schemaVersion: 1, commit: state.head as any, manifestSha256: digest(JSON.stringify(state)) as any, workspaces: ids.map(id => { const revision = state.revisions.find(x => x.id === id)!; return { workspaceId: id as any, revision: revision.commit as any, descriptorBlob: revision.blob as any, manifestSha256: digest(JSON.stringify(revision)) as any }; }), requestDigest: request.requestDigest ?? request.requestSha256 }; } async bootstrap(): Promise { diff --git a/backend/src/workspaces/types.ts b/backend/src/workspaces/types.ts index 3017d6c7..7df43b39 100644 --- a/backend/src/workspaces/types.ts +++ b/backend/src/workspaces/types.ts @@ -14,6 +14,6 @@ export type WorkspaceErrorCode = | "workspace_invalid" | "binding_missing" | "workspace_not_activatable" | "workspace_stale" | "workspace_conflict" | "git_unavailable" | "git_auth_failed" | "git_non_fast_forward" | "git_push_rejected" - | "connector_unavailable" | "semantic_index_incompatible"; + | "connector_unavailable" | "semantic_index_incompatible" | "registry_bootstrap_recovery_conflict"; export type { WorkspaceV3 } from "./schema.js"; diff --git a/backend/test/fixtures/workspace-registry-addressed-worker.mjs b/backend/test/fixtures/workspace-registry-addressed-worker.mjs index cc7d872c..86d5c2a2 100644 --- a/backend/test/fixtures/workspace-registry-addressed-worker.mjs +++ b/backend/test/fixtures/workspace-registry-addressed-worker.mjs @@ -1 +1,11 @@ -import { mkdir, writeFile } from "node:fs/promises"; const root=process.env.JOB_ROOT; if (!root) throw new Error("JOB_ROOT required"); await mkdir(root,{recursive:true,mode:0o700}); const id=process.env.RUN_ID??"a".repeat(32); await writeFile(`${root}/${id}.json`,JSON.stringify({runId:id,phase:"request_claimed"})+"\n",{flag:"wx",mode:0o600}); process.stdout.write(JSON.stringify({ok:true,runId:id})); +import { mkdir, open, writeFile, readFile, rm } from "node:fs/promises"; +import { constants } from "node:fs"; +const jobs = process.env.JOB_ROOT; +if (!jobs) throw new Error("JOB_ROOT required"); +const id = process.env.RUN_ID ?? "a".repeat(32); +const barrier = process.env.BARRIER; +await mkdir(jobs, { recursive: true, mode: 0o700 }); +if (barrier) { await writeFile(`${barrier}/${process.pid}.ready`, "ready", { flag: "wx", mode: 0o600 }); while (true) { try { await readFile(`${barrier}/release`); break; } catch { await new Promise(r => setTimeout(r, 5)); } } } +let winner = false; +try { const h = await open(`${jobs}/${id}.json`, constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL | constants.O_NOFOLLOW, 0o600); await h.writeFile(JSON.stringify({ schemaVersion: 1, runId: id, phase: "request_claimed" }) + "\n"); await h.sync(); await h.close(); winner = true; } catch (e) { if (e?.code !== "EEXIST") throw e; } +process.stdout.write(JSON.stringify({ ok: true, runId: id, winner }) + "\n"); diff --git a/backend/test/registry-pull-job-exports.test.ts b/backend/test/registry-pull-job-exports.test.ts index 6eccd852..e94c8ad2 100644 --- a/backend/test/registry-pull-job-exports.test.ts +++ b/backend/test/registry-pull-job-exports.test.ts @@ -1 +1 @@ -import {describe,it,expect} from "vitest"; import * as publication from "../src/workspaces/registry-publication.js"; describe("pull job exports",()=>it("exports addressed state",()=>expect(publication.REGISTRY_SCAN_LIMITS_V1.entries).toBeGreaterThan(0))); +import {describe,it,expect} from "vitest"; import * as publication from "../src/workspaces/registry-publication.js"; describe("pull job exports",()=>it("exports addressed state and exact scan bounds",()=>expect(publication.REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries).toBe(4096))); diff --git a/backend/test/routes-workspaces.test.ts b/backend/test/routes-workspaces.test.ts index 645a2d18..86abf229 100644 --- a/backend/test/routes-workspaces.test.ts +++ b/backend/test/routes-workspaces.test.ts @@ -116,7 +116,7 @@ const revision: WorkspaceRevision = { snapshotPath: "/registry/snapshots/psd-clinical.yaml", }; -type RegistryFake = Pick; +type RegistryFake = Pick & { ensureBootstrapAddressed: ReturnType; publishAddressed: ReturnType }; function registryFake(overrides: Partial = {}): RegistryFake { return { @@ -129,6 +129,8 @@ function registryFake(overrides: Partial = {}): RegistryFake { list: vi.fn(async () => [revision]), read: vi.fn(async () => ({ workspace, revision })), publish: vi.fn(async () => revision), + ensureBootstrapAddressed: vi.fn(async () => ({ kind: "already_active", snapshot: { schemaVersion: 1, commit: revision.commit, manifestSha256: "a".repeat(64), workspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })), + publishAddressed: vi.fn(async () => ({ plan: { targetCommit: revision.commit, targetManifestSha256: "a".repeat(64), targetWorkspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })), ...overrides, }; } @@ -218,7 +220,7 @@ test("returns a redacted registry status and pulls without Git credential detail expect(status.statusCode).toBe(200); expect(status.json()).toEqual({ - branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: true, lastError: "git_auth_failed", + branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: false, }); expect(pull.statusCode).toBe(200); expect(JSON.stringify([status.json(), pull.json()])).not.toMatch(/token|password|ssh:\/\//i); diff --git a/backend/test/workspace-registry-addressed-process.test.ts b/backend/test/workspace-registry-addressed-process.test.ts index c73251d4..605fea41 100644 --- a/backend/test/workspace-registry-addressed-process.test.ts +++ b/backend/test/workspace-registry-addressed-process.test.ts @@ -1 +1,16 @@ -import {describe,it,expect} from "vitest"; describe("addressed publication process contract",()=>{it("has a bounded run id",()=>expect("a".repeat(32)).toHaveLength(32));}); +import { describe, expect, it } from "vitest"; +import { mkdtemp, mkdir, readdir, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { spawn } from "node:child_process"; +const worker = join(process.cwd(), "test/fixtures/workspace-registry-addressed-worker.mjs"); +async function run(env: Record) { return await new Promise((resolve, reject) => { const p = spawn(process.execPath, [worker], { env: { ...process.env, ...env }, stdio: ["ignore", "pipe", "pipe"] }); let out = ""; p.stdout.on("data", b => out += b); p.on("error", reject); p.on("exit", c => c === 0 ? resolve(out) : reject(new Error(`worker ${c}`))); }); } +describe("addressed publication process ownership", () => { + it("serializes two claimers and leaves one durable final artifact", async () => { + const root = await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-addressed-process-")); const barrier = await mkdir(join(root, "barrier"), { recursive: true }).then(() => join(root, "barrier")); const jobs = join(root, "jobs"); + const env = { JOB_ROOT: jobs, BARRIER: barrier, RUN_ID: "a".repeat(32) }; + const a = run(env), b = run(env); + for (let i = 0; i < 100; i++) { if ((await readdir(barrier)).filter(x => x.endsWith(".ready")).length === 2) break; await new Promise(r => setTimeout(r, 5)); } + await writeFile(join(barrier, "release"), "go", { flag: "wx", mode: 0o600 }); const [one, two] = await Promise.all([a,b]); const results = [JSON.parse(one), JSON.parse(two)]; + expect(results.filter(x => x.winner)).toHaveLength(1); expect(await readdir(jobs)).toEqual([`${"a".repeat(32)}.json`]); + }, 5000); +});