From 47f3166f89d31e36914aff575797781e97afec35 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 22:25:26 +0200 Subject: [PATCH] fix addressed registry recovery invariants --- backend/src/routes/sessions.ts | 23 +-- backend/src/routes/sql.ts | 7 +- backend/src/routes/workspaces.ts | 42 ++++- backend/src/workspaces/author-git-service.ts | 19 ++- .../src/workspaces/registry-publication.ts | 52 +++--- backend/src/workspaces/registry.ts | 154 ++++++++++++++---- .../registry-addressed-invariants.test.ts | 29 ++++ 7 files changed, 242 insertions(+), 84 deletions(-) create mode 100644 backend/test/registry-addressed-invariants.test.ts diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 3f03958d..0c77844d 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -7,7 +7,7 @@ import { getPrincipal } from "../auth/auth.js"; import type { PrincipalContext } from "../auth/principal.js"; import type { ReadinessManager } from "../runtime/readiness-manager.js"; import type { ListModelsFn } from "./meta.js"; -import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js"; +import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, type WorkspaceRegistry } from "../workspaces/registry.js"; import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js"; import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js"; @@ -43,6 +43,7 @@ export function sessionRoutes( maintenanceBarrier: MaintenanceBarrier; }, ) { + const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(d.workspaceRegistry); } catch { return d.workspaceRegistry as any; } }; const lifecycleTails = new Map>(); const boundRuntimes = new Map< string, @@ -99,9 +100,9 @@ export function sessionRoutes( app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); }); /** Include retained historical descriptors after the same addressed selector used by status/list. */ - const sessionRevisions = async () => { + const sessionRevisions = async (): Promise["listRetainedSnapshots"]>>> => { await d.workspaceRegistry.ensureBootstrapAddressed(d.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(d.workspaceRegistry)); - return d.workspaceRegistry.listRetainedSnapshots(); + return snapshotReader().listRetainedSnapshots(); }; const isNotFound = (error: unknown) => @@ -141,7 +142,7 @@ export function sessionRoutes( throw error; } }; - let revisions: Awaited>; + let revisions: Awaited["listRetainedSnapshots"]>>; try { revisions = await sessionRevisions(); } catch (registryError) { @@ -169,7 +170,7 @@ export function sessionRoutes( const saved = located.manifest as { workspace_id?: string; workspace_revision?: string }; if (!saved.workspace_id || !saved.workspace_revision) return located; try { - const pinned = await d.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision); + const pinned = await snapshotReader().readPinned(saved.workspace_id, saved.workspace_revision); const workspace = validateOperationalWorkspace(pinned.workspace); return { ...located, @@ -320,7 +321,7 @@ export function sessionRoutes( code: "workspace_revision_unavailable", }); } - let revisionLease: Awaited> | undefined; + let revisionLease: Awaited["acquireSessionRevision"]>> | undefined; let manifestPersisted = false; try { let workspaceConfigPath: string | undefined; @@ -331,11 +332,11 @@ export function sessionRoutes( if (requestedWorkspaceId) { try { const registry = d.workspaceRegistry as Partial; - const resolved = typeof registry.acquireSessionRevision === "function" - ? await registry.acquireSessionRevision.call(d.workspaceRegistry, requestedWorkspaceId) - : await d.workspaceRegistry.read(requestedWorkspaceId); + const resolved = typeof (registry as any).acquireSessionRevision === "function" + ? await snapshotReader().acquireSessionRevision(requestedWorkspaceId) + : await snapshotReader().read(requestedWorkspaceId); if ("markPersisted" in resolved && "abort" in resolved) { - revisionLease = resolved as Awaited>; + revisionLease = resolved as Awaited["acquireSessionRevision"]>>; } if (!d.workspaceRuntimeSupport(resolved.workspace)) { return reply.code(409).send({ @@ -473,7 +474,7 @@ export function sessionRoutes( const runner = runnerFor(scopedPrincipal); const revisions = await sessionRevisions(); const lists = await Promise.all(revisions - .map((revision) => runner.sessionList(revision.snapshotPath) as Promise)); + .map((revision: { snapshotPath: string }) => runner.sessionList(revision.snapshotPath) as Promise)); const sessions = new Map(); for (const row of lists.flat()) { if (!sessions.has(row.id)) sessions.set(row.id, row); diff --git a/backend/src/routes/sql.ts b/backend/src/routes/sql.ts index e280fc63..3a1fba10 100644 --- a/backend/src/routes/sql.ts +++ b/backend/src/routes/sql.ts @@ -3,13 +3,14 @@ import type { ThtRunner } from "../tht/tht-runner.js"; import { getPrincipal } from "../auth/auth.js"; import type { PrincipalContext } from "../auth/principal.js"; import type { Settings } from "../settings/settings-store.js"; -import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js"; +import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, type WorkspaceRegistry } from "../workspaces/registry.js"; export function sqlRoutes(app: FastifyInstance, deps: { tht: ThtRunner; getSettings: (principal: PrincipalContext) => Promise; workspaceRegistry: WorkspaceRegistry; workspaceRegistryRecoveryIdentity?: () => ReturnType; }): void { + const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(deps.workspaceRegistry); } catch { return deps.workspaceRegistry as any; } }; const runnerFor = (principal: PrincipalContext): any => { const runner = deps.tht as any; return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner; @@ -21,14 +22,14 @@ export function sqlRoutes(app: FastifyInstance, deps: { const runner = runnerFor(principal); if (typeof runner.sessionShow !== "function") return { manifest: {}, workspace: legacyWorkspace }; await deps.workspaceRegistry.ensureBootstrapAddressed(deps.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(deps.workspaceRegistry)); - const revisions = await deps.workspaceRegistry.listRetainedSnapshots(); + const revisions = await snapshotReader().listRetainedSnapshots(); for (const revision of revisions) { try { const manifest = await runner.sessionShow(id, revision.snapshotPath); if (!manifest) continue; const saved = manifest as { workspace_id?: string; workspace_revision?: string }; if (saved.workspace_id && saved.workspace_revision) { - const pinned = await deps.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision); + const pinned = await snapshotReader().readPinned(saved.workspace_id, saved.workspace_revision); return { manifest, workspace: pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath, diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 1e8d6c02..17d44087 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -6,7 +6,7 @@ import yauzl from "yauzl"; import yazl from "yazl"; import { z } from "zod"; import type { WorkspaceRegistryConfig } from "../workspaces/types.js"; -import type { WorkspaceAuthorGitService } from "../workspaces/author-git-service.js"; +import { publishAddressedByAuthor, type WorkspaceAuthorGitService } from "../workspaces/author-git-service.js"; import { WorkspaceRegistryError } from "../workspaces/git-repository.js"; import { addressedRunId } from "../workspaces/registry-publication.js"; import { @@ -14,6 +14,7 @@ import { type PublishWorkspaceRequest, type WorkspaceRegistry, workspaceRegistryRecoveryIdentity, + workspaceRegistrySnapshotReader, } from "../workspaces/registry.js"; import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js"; import type { Revision40 } from "../workspaces/workspace-lock-root-lease.js"; @@ -283,6 +284,13 @@ function publishRequest(value: unknown): PublishWorkspaceRequest { } export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps): void { + const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(deps.registry); } catch { return deps.registry as any; } }; + const readPinned = async (id: string, commit: string) => { + const reader = snapshotReader(); + if (typeof reader.readPinned === "function") return reader.readPinned(id, commit); + const legacy = await reader.read(id); + return { workspace: legacy.workspace, workspaceConfigPath: legacy.revision.snapshotPath }; + }; app.register(multipart, { limits: { fileSize: deps.config.maxImportBytes, files: 1, fields: 0, parts: 1 }, throwFileSizeLimit: true, @@ -325,7 +333,7 @@ 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: deps.snapshotPath(item.revision, item.workspaceId) })); return await Promise.all(revisions.map(async revision => { - const pinned = await deps.registry.readPinned(revision.id, revision.commit); + const pinned = await readPinned(revision.id, revision.commit); const workspace = pinned.workspace; const exactRevision = { ...revision, snapshotPath: pinned.workspaceConfigPath }; 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 }; @@ -336,7 +344,11 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.get("/workspaces/:id", async (request, reply) => { try { const { id } = z.object({ id: workspaceId }).parse(request.params); - return await deps.registry.read(id); + const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); + const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id); + if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); + const pinned = await readPinned(id, item.revision); + return { workspace: pinned.workspace, revision: { id, commit: item.revision, blob: item.descriptorBlob, snapshotPath: pinned.workspaceConfigPath } }; } catch (error) { return errorReply(reply, error); } @@ -355,7 +367,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.post("/workspaces/:id/test", async (request, reply) => { try { const { id } = z.object({ id: workspaceId }).parse(request.params); - const { workspace } = await deps.registry.read(id); + const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); + const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id); + if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); + const { workspace } = await readPinned(id, item.revision); let operational: CanonicalWorkspace; try { operational = validateOperationalWorkspace(workspace); @@ -387,10 +402,16 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) remoteRefIdentitySha256: identity.remoteRefIdentitySha256, expectedBaseCommit: requestValue.baseCommit as Revision40, }; - // Authoring is intentionally split from activation: this service creates and pushes - // the remote commit only. The addressed registry is the sole pointer owner. - await deps.authorService.publish(requestValue); - const result = await deps.registry.publishAddressed(addressed); + // Claim the addressed identity before any author Git network. A crash after + // push therefore leaves a durable run that can be resumed with this exact ID. + let result: RegistryAddressedResultV1; + if (Object.prototype.hasOwnProperty.call(deps.registry, "snapshotReader")) { + result = await publishAddressedByAuthor(deps.registry, deps.authorService, requestValue, addressed); + } else { + // Test/dry-run registry doubles predate the package-private coordinator. + await deps.authorService.publish(requestValue); + result = await deps.registry.publishAddressed(addressed); + } return { revision: revisionFromPublication(result, id) }; } catch (error) { return errorReply(reply, error); @@ -400,7 +421,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.get("/workspaces/:id/export", async (request, reply) => { try { const { id } = z.object({ id: workspaceId }).parse(request.params); - const { workspace } = await deps.registry.read(id); + const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); + const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id); + if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); + const { workspace } = await readPinned(id, item.revision); const canonical = validateWorkspaceDescriptor(workspace); const bundle = await exportBundle(canonical); return reply diff --git a/backend/src/workspaces/author-git-service.ts b/backend/src/workspaces/author-git-service.ts index 9f97050f..c4ce6a2b 100644 --- a/backend/src/workspaces/author-git-service.ts +++ b/backend/src/workspaces/author-git-service.ts @@ -1,4 +1,5 @@ -import { WorkspaceConflictError, type PublishWorkspaceRequest } from "./registry.js"; +import { WorkspaceConflictError, claimAddressedPublication, type PublishWorkspaceRequest, type WorkspaceRegistry } from "./registry.js"; +import type { RegistryAddressedRequestV1, RegistryAddressedResultV1 } from "./registry-publication.js"; import { GitWorkspaceRepository, WorkspaceRegistryError, WorkspaceRepositoryLock } from "./git-repository.js"; import { parseWorkspaceYaml, renderWorkspaceDocs, serializeWorkspaceYaml, validateOperationalWorkspace, type CanonicalWorkspace, type WorkspaceDescriptor } from "./schema.js"; @@ -104,3 +105,19 @@ export class WorkspaceAuthorGitService { return new WorkspaceConflictError(fields, expected, actual, base, local, remote); } } + +/** + * Repository-author/publication friend. The addressed claim is durable before the + * author Git service performs pull/stage/push network, and activation always resumes + * that exact run identity. + */ +export async function publishAddressedByAuthor( + registry: WorkspaceRegistry, + author: WorkspaceAuthorGitService, + request: PublishWorkspaceRequest, + addressed: Extract, +): Promise { + await claimAddressedPublication(registry, addressed); + await author.publish(request); + return registry.publishAddressed({ ...addressed, mode: "resume" }); +} diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index 72be6bc1..393df6a1 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 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 terminalPublication: "target" | "reconciled_target" | "unchanged" | 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; readonly afterPublication?: (result: T) => Promise; }): Promise { + async run(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => 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,7 +45,6 @@ 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; @@ -107,18 +106,16 @@ 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()), base?: { readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] }): Promise { + async claim(request: RegistryAddressedRequestV1, base: { readonly commit: Revision40; readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] } | null): Promise { await this.dirs(); - 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 (!RUN.test(request.runId) || request.mode !== "create") 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, 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; + if (request.operation === "registry_pull" && base && request.expectedBaseCommit !== base.commit) throw CONFLICT(); + if (request.operation === "registry_bootstrap" && (base !== null || request.expectedBaseCommit !== null)) throw CONFLICT(); + if (request.operation === "registry_pull" && !base) throw CONFLICT(); + const state = { schemaVersion: 1, runId: request.runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${request.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, terminalPublication: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" ? request.expectedBaseCommit : null, baseManifestSha256: base?.manifestSha256 ?? null, baseWorkspaces: base?.workspaces ?? [], changedSetRule: null }; + await this.durable(this.path(request.runId), state, true); + return state as unknown as RegistryAddressedPublicationStateV1; } async storeTerminalResult(runId: RegistryRunId32, result: RegistryAddressedResultV1): Promise { @@ -127,7 +124,7 @@ export class RegistryAddressedPublicationStore { // 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 }); + await this.transition(runId, "terminal_durable", { terminalResultSha256: expected, terminalPublication: result.publication }); } 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(); @@ -150,7 +147,7 @@ export class RegistryAddressedPublicationStore { 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; + plan, planSha256: state.planSha256!, phase: "terminal_durable" as const, publication: state.terminalPublication! } as RegistryAddressedResultV1; if (registryDigest(result) !== expectedDigest) throw CONFLICT(); return result; } @@ -161,12 +158,18 @@ 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","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","terminalPublication","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.terminalPublication !== null && s.terminalPublication !== "target" && s.terminalPublication !== "reconciled_target" && s.terminalPublication !== "unchanged") 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(); + // The advertisement is the sole target identity. Every later durable phase must + // carry exactly that OID; accepting a different fetched or planned OID would make + // resume capable of silently publishing an object other than the claimed target. + if (s.advertisedTargetCommit !== null && s.fetchedTargetCommit !== null && s.fetchedTargetCommit !== s.advertisedTargetCommit) throw CONFLICT(); + if (s.advertisedTargetCommit !== null && s.targetCommit !== null && s.targetCommit !== s.advertisedTargetCommit) 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(); @@ -175,13 +178,13 @@ export class RegistryAddressedPublicationStore { 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"], + planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"], }; 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(); 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"], + participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"], }; const currentIndex = PHASES.indexOf(s.phase); let chainPrevious: string | null = null; @@ -195,26 +198,23 @@ export class RegistryAddressedPublicationStore { } async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial = {}): Promise { const old = await this.read(runId); - // The compatibility-shaped store fixture may advance only the phase; production callers - // always provide the durable phase fields. Keep that fixture deterministic without weakening - // validation of persisted production records. - if (Object.keys(patch).length === 0 && phase === "target_advertised" && old.phase === "request_claimed") patch = { advertisedTargetCommit: "0".repeat(40) as Revision40, immutableTargetRef: `refs/thoth/addressed-runs/${runId}/target` }; - 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", "changedSetRule", "baseManifestSha256", "baseWorkspaces"]); + const allowed = new Set(["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "terminalPublication", "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(); } + if ("fetchedTargetCommit" in patch && old.advertisedTargetCommit !== null && patch.fetchedTargetCommit !== old.advertisedTargetCommit) throw CONFLICT(); + if ("targetCommit" in patch && old.advertisedTargetCommit !== null && patch.targetCommit !== old.advertisedTargetCommit) 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"], + target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"], }; 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) { @@ -229,7 +229,7 @@ export class RegistryAddressedPublicationStore { participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], - terminal_durable: ["terminalResultSha256"], + terminal_durable: ["terminalResultSha256", "terminalPublication"], }; for (const key of required[phase]) if (next[key] === null || next[key] === undefined) throw CONFLICT(); await this.durable(this.path(runId), next); diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 76ecf211..53c2a5e5 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -186,9 +186,48 @@ export function createWorkspaceRegistry(config: WorkspaceRegistryConfig, deps: W lifecycleOwner: input.lifecycleOwner, participants: input.participants, synchronizers: input.synchronizers, installationIdentity, repositoryIdentity, remoteIdentity, }); + // Compatibility bindings are instance-owned and intentionally absent from the + // frozen WorkspaceRegistry prototype. New callers receive the separate reader. + const reader = workspaceRegistrySnapshotReader(registry); + Object.defineProperties(registry, { + listRetainedSnapshots: { value: reader.listRetainedSnapshots.bind(reader) }, + read: { value: reader.read.bind(reader) }, + acquireSessionRevision: { value: reader.acquireSessionRevision.bind(reader) }, + readPinned: { value: reader.readPinned.bind(reader) }, + snapshotReader: { value: reader }, + }); return registry; } +interface SnapshotReaderOps { + listRetainedSnapshots(): Promise; + read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }>; + acquireSessionRevision(id: string): Promise; + readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }>; +} +export class WorkspaceRegistrySnapshotReader { + constructor(private readonly ops: SnapshotReaderOps) {} + listRetainedSnapshots() { return this.ops.listRetainedSnapshots(); } + read(id: string) { return this.ops.read(id); } + acquireSessionRevision(id: string) { return this.ops.acquireSessionRevision(id); } + readPinned(id: string, commit: string) { return this.ops.readPinned(id, commit); } +} +const snapshotReaders = new WeakMap(); +const addressedClaimers = new WeakMap Promise>(); + +export function workspaceRegistrySnapshotReader(registry: WorkspaceRegistry): WorkspaceRegistrySnapshotReader { + const reader = snapshotReaders.get(registry); + if (!reader) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized"); + return reader; +} + +/** Package-private author/publication coordinator seam: claim before author Git network. */ +export async function claimAddressedPublication(registry: WorkspaceRegistry, request: RegistryAddressedRequestV1): Promise { + const claim = addressedClaimers.get(registry); + if (!claim) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized"); + return claim(request); +} + export class WorkspaceRegistry { private readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory; private readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner; @@ -206,6 +245,14 @@ export class WorkspaceRegistry { this.lifecycleOwner = input.lifecycleOwner; this.participants = input.participants; this.synchronizers = input.synchronizers; + const reader = new WorkspaceRegistrySnapshotReader({ + listRetainedSnapshots: () => this.#listRetainedSnapshots(), + read: (id) => this.#read(id), + acquireSessionRevision: (id) => this.#acquireSessionRevision(id), + readPinned: (id, commit) => this.#readPinned(id, commit), + }); + snapshotReaders.set(this, reader); + addressedClaimers.set(this, (request) => this.#claimAddressed(request)); } async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise { @@ -234,28 +281,41 @@ export class WorkspaceRegistry { installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, expectedBaseCommit: null, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, }; - state = await store.claim(request, request.runId); + state = await store.claim(request, null); } const result = await this.#executeAddressedLocked(store, state, identity); const published = await this.#activeState(); return { kind: "bootstrap_terminal", result: result as Extract, snapshot: this.#addressedSnapshot(published, identity.requestSha256) }; } + async #claimAddressed(request: RegistryAddressedRequestV1): Promise { + if (request.mode !== "create") throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Only create requests may claim a run"); + await registryContext(this).repository.ensureLayout(); + await registryContext(this).lock.run(async () => { + const store = new RegistryAddressedPublicationStore(registryContext(this).repository.root); + let base: ActiveState | undefined; + if (request.operation === "registry_pull") base = await this.#snapshotState(request.expectedBaseCommit); + await store.claim(request, base ? { commit: base.head as Revision40, manifestSha256: registryDigest(base) as never, workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)) } : null); + }); + } + async publishAddressed(request: RegistryAddressedRequestV1): Promise { await registryContext(this).repository.ensureLayout(); return registryContext(this).lock.run(async () => { const store = new RegistryAddressedPublicationStore(registryContext(this).repository.root); let state: RegistryAddressedPublicationStateV1; - try { state = await store.read(request.runId); } - catch { + if (request.mode === "resume") { + try { state = await store.read(request.runId); } + catch { throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed resume requires an existing durable run"); } + } else { let base: ActiveState | undefined; - if (request.operation === "registry_pull" && request.mode === "create") { - base = await this.#snapshotState(request.expectedBaseCommit); - } - state = await store.claim(request, request.runId, base ? { - manifestSha256: registryDigest(base) as never, + if (request.operation === "registry_pull") base = await this.#snapshotState(request.expectedBaseCommit); + // 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, workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)), - } : undefined); + } : null); } const identity: RegistryBootstrapRecoveryIdentityV1 = { operation: "registry_bootstrap", requestSha256: request.requestSha256, @@ -271,8 +331,14 @@ export class WorkspaceRegistry { async #executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise { let state = initial; - if (state.phase === "terminal_durable") return store.readTerminalResult(state.runId, state.terminalResultSha256!); + if (state.phase === "terminal_durable") { + const pinned = await registryContext(this).repository.runRef(state.runId); + if (pinned !== state.advertisedTargetCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed target ref drifted"); + return store.readTerminalResult(state.runId, state.terminalResultSha256!); + } const operation = state.operation; + let publication: "target" | "reconciled_target" | "unchanged" = "target"; + const resumedAtTargetPublished = state.phase === "target_published"; let target = state.advertisedTargetCommit ?? state.fetchedTargetCommit; if (state.phase === "request_claimed") { const status = operation === "registry_bootstrap" ? await registryContext(this).repository.bootstrap() : await registryContext(this).repository.pull(); @@ -290,6 +356,12 @@ export class WorkspaceRegistry { } target = state.fetchedTargetCommit ?? target; if (!target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + // Once advertised, the run-specific ref is immutable evidence. Every resume after + // the fetch barrier revalidates it before reading or publishing any bytes. + if (state.phase !== "target_advertised" && state.phase !== "request_claimed") { + const pinned = await registryContext(this).repository.runRef(state.runId); + if (pinned !== target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed target ref drifted"); + } if (state.phase === "target_fetched") { await this.#materialize(target); const targetState = await this.#snapshotState(target); @@ -311,30 +383,44 @@ export class WorkspaceRegistry { 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 }); } const plan = await this.#planForState(state); - const leases = []; - for (const id of plan.changedWorkspaceIds) leases.push(await registryContext(this).rootLeaseFactory.acquireOrProvision(registryContext(this).rootLeaseFactory.canonicalInput(id))); + const leases: Awaited>[] = []; + try { + for (const id of plan.changedWorkspaceIds) leases.push(await registryContext(this).rootLeaseFactory.acquireOrProvision(registryContext(this).rootLeaseFactory.canonicalInput(id))); + } catch (error) { + for (const lease of [...leases].reverse()) { try { await lease.close(); } catch {} } + throw error; + } let output: RegistryAddressedResultV1 | undefined; - await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => registryContext(this).lifecycleOwner.run({ plan, capabilities, participants: registryContext(this).participants, synchronizers: registryContext(this).synchronizers, action: async () => { - if (state.phase === "planned") { - state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(registryContext(this).participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(registryContext(this).synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never }); - } - if (state.phase === "participants_prepared") { - 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.#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. + await runUnderOrderedWorkspaceWriterLocks(leases, async capabilities => { + await registryContext(this).lifecycleOwner.run({ plan, capabilities, participants: registryContext(this).participants, synchronizers: registryContext(this).synchronizers, action: async () => { + if (state.phase === "planned") { + state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(registryContext(this).participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(registryContext(this).synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never }); + } + if (state.phase === "participants_prepared") { + 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") { + const activeBefore = await this.#tryActiveState(); + if (activeBefore?.head === target) { + 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 }); + } + return undefined; + }}); + // The lifecycle callback has now reconciled target state. Persist terminal while + // ordered writer/root capabilities remain owned, immediately before settlement. 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; + if (resumedAtTargetPublished) { + publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target"; + } + output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable" as const, publication } as RegistryAddressedResultV1; await store.storeTerminalResult(state.runId, output); state = await store.read(state.runId); } - }})); + }); if (!output) output = await store.readTerminalResult(state.runId, state.terminalResultSha256!); return output; } @@ -352,7 +438,7 @@ export class WorkspaceRegistry { 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); }) }; } - async listRetainedSnapshots(): Promise { + async #listRetainedSnapshots(): Promise { await registryContext(this).repository.ensureLayout(); return await registryContext(this).lock.run(async () => { try { @@ -372,7 +458,7 @@ export class WorkspaceRegistry { }); } - async read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> { + async #read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> { const state = await this.#activeState(); const revision = state.revisions.find((candidate) => candidate.id === id); if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); @@ -388,7 +474,7 @@ export class WorkspaceRegistry { * Resolve the active revision and create its cross-process retention lease under the same * repository lock. The lease bridges the interval before `session_manifest.yaml` is durable. */ - async acquireSessionRevision(id: string): Promise { + async #acquireSessionRevision(id: string): Promise { await registryContext(this).repository.ensureLayout(); return await registryContext(this).lock.run(async () => { const state = await this.#activeState(); @@ -437,7 +523,7 @@ export class WorkspaceRegistry { } /** Read a retained immutable snapshot for a session pinned to a historical commit. */ - async readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }> { + async #readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }> { const snapshotPath = workspaceRegistrySnapshotPath(this, safeCommit(commit), id); try { const source = await readFile(snapshotPath, "utf8"); @@ -567,7 +653,7 @@ export class WorkspaceRegistry { const base = await this.#readSnapshotCanonical(request.baseCommit, id); let remote: CanonicalWorkspace | undefined; if (existing) { - remote = (await this.read(id)).workspace; + remote = (await this.#read(id)).workspace; } return new WorkspaceConflictError( this.#changedFields(base, remote), diff --git a/backend/test/registry-addressed-invariants.test.ts b/backend/test/registry-addressed-invariants.test.ts new file mode 100644 index 00000000..a1c6aad0 --- /dev/null +++ b/backend/test/registry-addressed-invariants.test.ts @@ -0,0 +1,29 @@ +import { describe, expect, it } from "vitest"; +import { mkdtemp } from "node:fs/promises"; +import { join } from "node:path"; +import { RegistryAddressedPublicationStore } from "../src/workspaces/registry-publication.js"; + +const request = (runId: string, mode: "create" | "resume" = "create") => ({ + mode, operation: "registry_bootstrap" as const, runId: runId as any, + requestSha256: "1".repeat(64) as any, installationIdentitySha256: "2".repeat(64) as any, + repositoryIdentitySha256: "3".repeat(64) as any, expectedBaseCommit: null, + remoteRefIdentitySha256: "4".repeat(64) as any, +}); + +describe("addressed registry invariants", () => { + it("pins one immutable target OID through fetch and plan", async () => { + const store = new RegistryAddressedPublicationStore(await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-reg-"))); + const runId = "a".repeat(32); + await store.claim(request(runId), null); + await store.transition(runId as any, "target_advertised", { advertisedTargetCommit: "1".repeat(40) as any, immutableTargetRef: `refs/thoth/addressed-runs/${runId}/target` }); + await expect(store.transition(runId as any, "target_fetched", { fetchedTargetCommit: "2".repeat(40) as any })).rejects.toThrow(); + }); + + it("requires resume to find an existing run and rejects create replay", async () => { + const store = new RegistryAddressedPublicationStore(await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-reg-"))); + const runId = "b".repeat(32); + await expect(store.claim(request(runId, "resume"), null)).rejects.toThrow(); + await store.claim(request(runId), null); + await expect(store.claim(request(runId), null)).rejects.toThrow(); + }); +});