diff --git a/backend/src/app.ts b/backend/src/app.ts index f1210f2a..6fd937c9 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -1,6 +1,8 @@ import Fastify, { type FastifyInstance } from "fastify"; import cors from "@fastify/cors"; import { join } from "node:path"; +import { mkdirSync, realpathSync } from "node:fs"; +import { createHash } from "node:crypto"; import type { AppConfig } from "./config.js"; import { ThtRunner } from "./tht/tht-runner.js"; import { PiProcessManager } from "./pi/pi-process-manager.js"; @@ -18,6 +20,10 @@ import { loadSettings, type Settings } from "./settings/settings-store.js"; import { ReadinessManager } from "./runtime/readiness-manager.js"; import { MaintenanceBarrier } from "./runtime/maintenance-gate.js"; import { WorkspaceRegistry } from "./workspaces/registry.js"; +import { GitWorkspaceRepository } from "./workspaces/git-repository.js"; +import { WorkspaceFsAtV1 } from "./workspaces/workspace-fs-at.js"; +import { VerifiedWorkspaceLockRootLeaseFactory } from "./workspaces/workspace-lock-root-lease.js"; +import { CapabilityAwareRegistryPublicationLifecycleOwner, canonicalBootstrapRequestDigest } from "./workspaces/registry-publication.js"; import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js"; import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js"; import { piManagementRoutes } from "./routes/pi-management.js"; @@ -39,6 +45,24 @@ export interface BuildAppDeps { piManagement?: PiManagementService; } +function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry { + const registryConfig = config.workspaceRegistry; + const repository = new GitWorkspaceRepository(registryConfig); + const sessionsRoot = join(registryConfig.root, "sessions"); + mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 }); + const rootLeaseFactory = new VerifiedWorkspaceLockRootLeaseFactory({ + workspaceFsAt: new WorkspaceFsAtV1(), installationId: registryConfig.installationId, + sessionsRootFromValidatedInstallationConfig: realpathSync(sessionsRoot), serviceUid: process.getuid?.() ?? 0, + provisionedWorkspaceMode: 0o700, + }); + const hash = (value: string) => createHash("sha256").update(value).digest("hex"); + const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(`${registryConfig.remoteUrl ?? ""}:${registryConfig.branch}`) }; + const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: repositoryIdentity.digest }; + const request = { kind: "bootstrap" as const, installation: { installationId: registryConfig.installationId, digest: hash(registryConfig.installationId) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] }; + return new WorkspaceRegistry({ rootLeaseFactory, lifecycleOwner: new CapabilityAwareRegistryPublicationLifecycleOwner(), participants: [], synchronizers: [], repository, + installationIdentity: { operation: "registry_bootstrap", requestSha256: canonicalBootstrapRequestDigest(request) as never, installationIdentitySha256: request.installation.digest as never, repositoryIdentitySha256: repositoryIdentity.digest as never, remoteRefIdentitySha256: remoteIdentity.digest as never }, repositoryIdentity, remoteIdentity }); +} + export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance { const app = Fastify({ logger: { level: "warn" }, disableRequestLogging: true }); @@ -67,7 +91,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc }); const mgr = deps?.mgr ?? new PiProcessManager(config, deps?.spawnFn ? { spawnFn: deps.spawnFn } : undefined); const hub = deps?.hub ?? new SseHub(); - const workspaceRegistry = deps?.workspaceRegistry ?? new WorkspaceRegistry(config.workspaceRegistry); + const workspaceRegistry = deps?.workspaceRegistry ?? createWorkspaceRegistry(config); const workspaceDiagnoser = deps?.workspaceDiagnoser ?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs, undefined, { internalQdrantUrl: config.internalQdrantUrl, diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 1338af53..350c4c55 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -96,13 +96,10 @@ export function sessionRoutes( }); app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); }); - /** Include retained historical descriptors so removed workspaces remain resumable. */ + /** Include retained historical descriptors after the same addressed selector used by status/list. */ const sessionRevisions = async () => { - const registry = d.workspaceRegistry as Partial; - if (typeof registry.listRetainedSnapshots === "function") { - return await registry.listRetainedSnapshots(); - } - return await d.workspaceRegistry.list(); + await d.workspaceRegistry.ensureBootstrapAddressed(d.workspaceRegistry.recoveryIdentity()); + return d.workspaceRegistry.listRetainedSnapshots(); }; const isNotFound = (error: unknown) => @@ -142,7 +139,7 @@ export function sessionRoutes( throw error; } }; - let revisions: Awaited>; + let revisions: Awaited>; try { revisions = await sessionRevisions(); } catch (registryError) { diff --git a/backend/src/routes/sql.ts b/backend/src/routes/sql.ts index a5e48554..65eb32b5 100644 --- a/backend/src/routes/sql.ts +++ b/backend/src/routes/sql.ts @@ -19,10 +19,8 @@ export function sqlRoutes(app: FastifyInstance, deps: { const locate = async (principal: PrincipalContext, id: string, legacyWorkspace?: string) => { const runner = runnerFor(principal); if (typeof runner.sessionShow !== "function") return { manifest: {}, workspace: legacyWorkspace }; - const registry = deps.workspaceRegistry as Partial; - const revisions = typeof registry.listRetainedSnapshots === "function" - ? await registry.listRetainedSnapshots.call(deps.workspaceRegistry) - : await deps.workspaceRegistry.list(); + await deps.workspaceRegistry.ensureBootstrapAddressed(deps.workspaceRegistry.recoveryIdentity()); + const revisions = await deps.workspaceRegistry.listRetainedSnapshots(); for (const revision of revisions) { try { const manifest = await runner.sessionShow(id, revision.snapshotPath); diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 89e9ab47..1318e7c5 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -7,6 +7,7 @@ import yazl from "yazl"; import { z } from "zod"; import type { WorkspaceRegistryConfig } from "../workspaces/types.js"; import { WorkspaceRegistryError } from "../workspaces/git-repository.js"; +import { addressedRunId } from "../workspaces/registry-publication.js"; import { WorkspaceConflictError, type PublishWorkspaceRequest, @@ -365,8 +366,11 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.post("/workspaces/publish", async (request, reply) => { try { - const result = await deps.registry.publish(publishRequest(request.body)); - return result ? { revision: result } : reply.code(204).send(); + const identity = deps.registry.recoveryIdentity(); + const body = request.body as { baseCommit?: string }; + if (!body.baseCommit || !/^[0-9a-f]{40}$/.test(body.baseCommit)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid"); + const result = await deps.registry.publishAddressed({ mode: "create", operation: "registry_pull", runId: addressedRunId(), requestSha256: identity.requestSha256 as never, installationIdentitySha256: identity.installationIdentitySha256 as never, repositoryIdentitySha256: identity.repositoryIdentitySha256 as never, remoteRefIdentitySha256: identity.remoteRefIdentitySha256 as never, expectedBaseCommit: body.baseCommit as never }); + return { revision: result }; } catch (error) { return errorReply(reply, error); } diff --git a/backend/src/workspaces/git-repository.ts b/backend/src/workspaces/git-repository.ts index c146523f..66f2cede 100644 --- a/backend/src/workspaces/git-repository.ts +++ b/backend/src/workspaces/git-repository.ts @@ -70,7 +70,7 @@ export class GitWorkspaceRepository { readonly locksPath: string; private readonly hooksPath: string; - constructor(private readonly config: WorkspaceRegistryConfig) { + constructor(readonly config: WorkspaceRegistryConfig) { if (!isAbsolute(config.root)) { throw new WorkspaceRegistryError("git_unavailable", "Workspace registry root is unavailable"); } @@ -129,6 +129,44 @@ export class GitWorkspaceRepository { }; } + private validateRevision(revision: string): void { + if (!/^[0-9a-f]{40}$/.test(revision)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid"); + } + async workspacePathsAt(revision: string): Promise { + this.validateRevision(revision); + const output = await this.git(["ls-tree", "-r", "--name-only", revision, "--", "workspaces"]); + const paths = output.trim() === "" ? [] : output.trim().split("\n"); + for (const path of paths) if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository contains an invalid path"); + return paths; + } + async readWorkspaceAt(revision: string, path: string): Promise { + this.validateRevision(revision); + if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid"); + return this.git(["show", `${revision}:${path}`]); + } + async blobAt(revision: string, path: string): Promise { + this.validateRevision(revision); + if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid"); + return (await this.git(["rev-parse", `${revision}:${path}`])).trim(); + } + async fetchExact(revision: string): Promise { + this.validateRevision(revision); + if (!this.config.remoteUrl) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry remote is unavailable"); + await this.git(["fetch", "--no-tags", "origin", revision]); + } + async ensureRunRef(runId: string, revision: string): Promise { + this.validateRevision(revision); + if (!/^[0-9a-f]{32}$/.test(runId)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace publication run is invalid"); + const ref = `refs/thoth/addressed-runs/${runId}/target`; + const current = await this.gitOptional(["rev-parse", "--verify", ref]); + if (current !== undefined && current.trim() !== revision) throw new WorkspaceRegistryError("workspace_conflict", "Workspace publication target conflicts"); + if (current === undefined) await this.git(["update-ref", ref, revision]); + } + async runRef(runId: string): Promise { + if (!/^[0-9a-f]{32}$/.test(runId)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace publication run is invalid"); + return (await this.gitOptional(["rev-parse", "--verify", `refs/thoth/addressed-runs/${runId}/target`]))?.trim(); + } + async workspacePaths(): Promise { const output = await this.git(["ls-tree", "-r", "--name-only", "HEAD", "--", "workspaces"]); const paths = output.trim() === "" ? [] : output.trim().split("\n"); diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index bf377aca..d94e4d17 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -30,27 +30,43 @@ export interface CapabilityAwareRegistryPublicationParticipant { readonly par 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 { + constructor() {} 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 }> = []; + // 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. + const prepared: Array<{ readonly participant: CapabilityAwareRegistryPublicationParticipant; readonly value: unknown; readonly 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); - })); + if (index >= input.capabilities.workspaceIds.length) { + for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "planned"); + for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "participants_prepared"); + 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"); + return result; } - 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; + const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId; + return input.capabilities.forWorkspace(id, ({ rootLease, writerCapability }) => + writerCapability.runUnderSessionReadersExclusive(async readers => { + const lease: AddressedWorkspacePublicationLeaseV1 = { + workspaceId: id, + rootLease, + writerCapability, + quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id), + readers, + }; + for (const participant of input.participants) { + const value = await participant.prepare(input.plan, lease); + prepared.push({ participant, value, lease }); + } + return enter(index + 1); + }), + ); }; return enter(0); } } + export type RegistryAddressedRequestV1 = | { readonly mode: "create"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: null; readonly remoteRefIdentitySha256: Sha256Hex } | { readonly mode: "resume"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex } @@ -83,6 +99,7 @@ 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 resultPath(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.result.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 { @@ -92,6 +109,28 @@ export class RegistryAddressedPublicationStore { 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: null, baseWorkspaces: [], changedSetRule: null }; await this.durable(this.path(runId), state, true); return state as unknown as RegistryAddressedPublicationStateV1; } + async setBase(runId: RegistryRunId32, baseManifestSha256: Sha256Hex, baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]): Promise { + const old = await this.read(runId); + if (old.phase !== "request_claimed" || old.operation !== "registry_pull") throw CONFLICT(); + const next = { ...old, baseManifestSha256, baseWorkspaces, priorStateSha256: registryDigest(old) as Sha256Hex }; + await this.durable(this.path(runId), next); + return next; + } + + async storeTerminalResult(runId: RegistryRunId32, result: RegistryAddressedResultV1): Promise { + await this.dirs(); + const path = this.resultPath(runId); + await this.durable(path, { schemaVersion: 1, runId, result }, true); + } + async readTerminalResult(runId: RegistryRunId32, expectedDigest: Sha256Hex): Promise { + let parsed: unknown; + try { parsed = JSON.parse((await strictRead(this.resultPath(runId))).text); } catch { throw CONFLICT(); } + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw CONFLICT(); + const value = parsed as { schemaVersion?: unknown; runId?: unknown; result?: unknown }; + if (value.schemaVersion !== 1 || value.runId !== runId || !value.result || registryDigest(value.result) !== expectedDigest) throw CONFLICT(); + return value.result as RegistryAddressedResultV1; + } + async read(runId: RegistryRunId32): Promise { const path = this.path(runId); let parsed: unknown; try { parsed = JSON.parse((await strictRead(path)).text); } catch { throw CONFLICT(); } @@ -109,12 +148,69 @@ export class RegistryAddressedPublicationStore { 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(); + const required: Record = { + request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], target_fetched: ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit"], + planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"], + }; + for (const key of required[s.phase]) if (s[key] === null || s[key] === undefined) 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; } + async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial = {}): Promise { + const old = await this.read(runId); + 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"]); + 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[key] !== old[key]) throw CONFLICT(); + } + const next = { ...old, ...patch, phase, priorStateSha256: registryDigest(old) as Sha256Hex } as RegistryAddressedPublicationStateV1; + const required: Record = { + request_claimed: [], + target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], + target_fetched: ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit"], + planned: ["fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "changedSetSha256", "planSha256", "changedSetRule"], + participants_prepared: ["participantsSha256"], + publication_intent_durable: ["publicationIntentSha256"], + target_published: ["publishedActiveStateSha256"], + terminal_durable: ["terminalResultSha256"], + }; + for (const key of required[phase]) if (next[key] === null || next[key] === undefined) throw CONFLICT(); + await this.durable(this.path(runId), next); + return next; + } snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryActiveSnapshotV1 { if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces) throw CONFLICT(); return { schemaVersion: 1, commit: state.targetCommit, manifestSha256: state.targetManifestSha256, workspaces: state.targetWorkspaces }; } 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; } + 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)); + const resultNames = entries.filter(x => /^[0-9a-f]{32}\.result\.json$/.test(x)); + if (names.length + resultNames.length !== entries.length) throw CONFLICT(); + let total = 0; + const identities = new Map(); + const out: RegistryAddressedPublicationStateV1[] = []; + for (const name of entries) { + 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)}`); + } + for (const name of names) { + const state = await this.read(name.slice(0, -5) as RegistryRunId32); + if (state.phase === "terminal_durable") { + await this.readTerminalResult(state.runId, state.terminalResultSha256 as Sha256Hex); + } + out.push(state); + } + const verify = (await readdir(this.jobsDirectory)).sort(); + if (verify.length !== entries.length || verify.some((name, i) => name !== entries[i])) throw CONFLICT(); + for (const name of entries) { + 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 9dbac9f3..f8b784bf 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -17,16 +17,9 @@ import { type WorkspaceDescriptor, } from "./schema.js"; import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js"; -import { - RegistryAddressedPublicationStore, - addressedRunId, - canonicalBootstrapRequestDigest, - type RegistryAddressedSnapshotV1, - type RegistryEnsureBootstrapAddressedResultV1, - type RegistryAddressedRequestV1, - type RegistryAddressedPublicationResultV1, -} from "./registry-publication.js"; - +import type { VerifiedWorkspaceLockRootLeaseFactory, Revision40, CanonicalWorkspaceId } from "./workspace-lock-root-lease.js"; +import { runUnderOrderedWorkspaceWriterLocks, type OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js"; +import { RegistryAddressedPublicationStore, addressedRunId, canonicalBootstrapRequestDigest, registryDigest, CapabilityAwareRegistryPublicationLifecycleOwner, type CapabilityAwareRegistryPublicationParticipant, type CapabilityAwareRegistryPublicationSynchronizer, type RegistryAddressedPlanV1, type RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.js"; export type { GitStatus } from "./git-repository.js"; export interface WorkspaceRevision { @@ -112,131 +105,178 @@ function workspaceError(error: unknown): WorkspaceRegistryError { } /** Immutable canonical workspace snapshots backed by the configured Git checkout. */ +export interface WorkspaceRegistryDependencies { + readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory; + readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner; + readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; + readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; + readonly repository?: GitWorkspaceRepository; + readonly installationIdentity?: RegistryBootstrapRecoveryIdentityV1; + readonly repositoryIdentity?: { readonly digest: string; readonly remote: string; readonly branch: string; readonly head: string }; + readonly remoteIdentity?: { readonly digest: string; readonly remote: string; readonly head: string }; +} + export class WorkspaceRegistry { private readonly repository: GitWorkspaceRepository; private readonly lock: WorkspaceRepositoryLock; + private readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory; + private readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner; + private readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; + private readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; + private readonly installationIdentity: RegistryBootstrapRecoveryIdentityV1; + private readonly repositoryIdentity: WorkspaceRegistryDependencies["repositoryIdentity"]; + private readonly remoteIdentity: WorkspaceRegistryDependencies["remoteIdentity"]; - constructor(private readonly config: WorkspaceRegistryConfig) { - this.repository = new GitWorkspaceRepository(config); + constructor(input: WorkspaceRegistryDependencies) { + if (!input.repository || !input.installationIdentity || !input.repositoryIdentity || !input.remoteIdentity) throw new Error("WorkspaceRegistry requires bound repository identities"); + this.repository = input.repository; this.lock = new WorkspaceRepositoryLock(this.repository.locksPath); + this.rootLeaseFactory = input.rootLeaseFactory; + this.lifecycleOwner = input.lifecycleOwner; + this.participants = input.participants; + this.synchronizers = input.synchronizers; + this.installationIdentity = input.installationIdentity; + this.repositoryIdentity = input.repositoryIdentity; + this.remoteIdentity = input.remoteIdentity; } snapshotPath(commit: string, id: string): string { return join(this.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`); } - /** Automatic addressed recovery. The repository lock is acquired before active-state inspection. */ - async ensureBootstrapAddressed(request?: any): Promise { + /** Automatic addressed recovery. The repository lock is held for selection and execution. */ + recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; } + + async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise { await this.repository.ensureLayout(); - 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) }; + return this.lock.run(async () => this.ensureBootstrapAddressedLocked(identity)); + } + + private async ensureBootstrapAddressedLocked(identity: RegistryBootstrapRecoveryIdentityV1): Promise { + let active: ActiveState | undefined; + try { active = await this.tryActiveState(); } + catch { throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); } + if (active) return { kind: "already_active", snapshot: this.addressedSnapshot(active, identity.requestSha256) }; + const store = new RegistryAddressedPublicationStore(this.repository.root); + let jobs: readonly RegistryAddressedPublicationStateV1[]; + try { jobs = await store.scan(); } catch { throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); } + const nonterminal = jobs.filter(job => job.phase !== "terminal_durable"); + const matching = nonterminal.filter(job => job.operation === identity.operation && job.requestSha256 === identity.requestSha256 && job.installationIdentitySha256 === identity.installationIdentitySha256 && job.repositoryIdentitySha256 === identity.repositoryIdentitySha256 && job.remoteRefIdentitySha256 === identity.remoteRefIdentitySha256); + 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 = matching[0]; + if (!state) { + const request: RegistryAddressedRequestV1 = { + mode: "create", operation: "registry_bootstrap", runId: addressedRunId(), requestSha256: identity.requestSha256, + installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, + expectedBaseCommit: null, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, + }; + state = await store.claim(request, request.runId); + } + 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 publishAddressed(request: RegistryAddressedRequestV1): Promise { + await this.repository.ensureLayout(); + return this.lock.run(async () => { 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); + let state: RegistryAddressedPublicationStateV1; + try { state = await store.read(request.runId); } + catch { + state = await store.claim(request, request.runId); + if (request.operation === "registry_pull" && request.mode === "create") { + const base = await this.snapshotState(request.expectedBaseCommit); + const baseWorkspaces = base.revisions.map(revision => this.manifestIdentity(revision)); + state = await store.setBase(state.runId, registryDigest(base) as never, baseWorkspaces); + } } - 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) }; + const identity: RegistryBootstrapRecoveryIdentityV1 = { + operation: "registry_bootstrap", requestSha256: request.requestSha256, + installationIdentitySha256: request.installationIdentitySha256, + repositoryIdentitySha256: request.repositoryIdentitySha256, + remoteRefIdentitySha256: request.remoteRefIdentitySha256, + }; + return this.executeAddressedLocked(store, state, identity); }); } - async publishAddressed(request: 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 }; + private async executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise { + let state = initial; + if (state.phase === "terminal_durable") return store.readTerminalResult(state.runId, state.terminalResultSha256!); + const operation = state.operation; + let target = state.advertisedTargetCommit ?? state.fetchedTargetCommit; + if (state.phase === "request_claimed") { + const status = operation === "registry_bootstrap" ? await this.repository.bootstrap() : await this.repository.pull(); + target = safeCommit(status.head ?? "") as Revision40; + state = await store.transition(state.runId, "target_advertised", { advertisedTargetCommit: target, immutableTargetRef: `refs/thoth/addressed-runs/${state.runId}/target` }); + } else if (!target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + if (state.phase === "target_advertised") { + const current = await this.repository.runRef(state.runId); + if (current !== undefined && current !== target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + if (current === undefined) { + await this.repository.fetchExact(target!); + await this.repository.ensureRunRef(state.runId, 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 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 }); + state = await store.transition(state.runId, "target_fetched", { fetchedTargetCommit: target }); + } + target = state.fetchedTargetCommit ?? target; + if (!target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + if (state.phase === "target_fetched") { + await this.materialize(target); + const targetState = await this.snapshotState(target); + const base = operation === "registry_pull" ? await this.baseState(state.baseCommit) : undefined; + const targetWorkspaces = targetState.revisions.map(revision => this.manifestIdentity(revision)); + const baseWorkspaces = base?.revisions.map(revision => this.manifestIdentity(revision)) ?? []; + const baseIds = new Set(baseWorkspaces.map(workspace => workspace.workspaceId)); + const targetIds = new Set(targetWorkspaces.map(workspace => workspace.workspaceId)); + const changedWorkspaceIds = [...new Set([...baseIds, ...targetIds])].filter(id => !baseIds.has(id) || !targetIds.has(id) || registryDigest(baseWorkspaces.find(x => x.workspaceId === id)) !== registryDigest(targetWorkspaces.find(x => x.workspaceId === id))).sort() as CanonicalWorkspaceId[]; + const plan = operation === "registry_bootstrap" ? { + schemaVersion: 1, operation, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, + jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit!, immutableTargetRef: state.immutableTargetRef!, fetchedTargetCommit: target as Revision40, targetCommit: target as Revision40, + targetManifestSha256: registryDigest(targetState) as never, targetWorkspaces, changedWorkspaceIds: targetWorkspaces.map(x => x.workspaceId).sort() as CanonicalWorkspaceId[], changedSetSha256: registryDigest(targetWorkspaces.map(x => x.workspaceId).sort()) as never, changedSetRule: "all_target_workspace_ids" as const, baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] as const, + } : { + schemaVersion: 1, operation, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, + jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit!, immutableTargetRef: state.immutableTargetRef!, fetchedTargetCommit: target as Revision40, targetCommit: target as Revision40, + targetManifestSha256: registryDigest(targetWorkspaces), targetWorkspaces, changedWorkspaceIds, changedSetSha256: registryDigest(changedWorkspaceIds) as never, changedSetRule: "symmetric_base_target_workspace_difference" as const, baseCommit: state.baseCommit!, baseManifestSha256: registryDigest(base) as never, baseWorkspaces, + }; + state = await store.transition(state.runId, "planned", { targetCommit: target as Revision40, targetManifestSha256: plan.targetManifestSha256 as never, targetWorkspaces: plan.targetWorkspaces, changedWorkspaceIds: plan.changedWorkspaceIds, changedSetSha256: plan.changedSetSha256 as never, planSha256: registryDigest(plan) as never, changedSetRule: plan.changedSetRule }); + } + const plan = await this.planForState(state); + const leases = []; + for (const id of plan.changedWorkspaceIds) leases.push(await this.rootLeaseFactory.acquireOrProvision(this.rootLeaseFactory.canonicalInput(id))); + const result = await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => { + if (state.phase === "planned") { + state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(this.participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(this.synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never }); + } + 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.publishMaterialized(target!); + state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.activeState()) as never }); + } + const output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan), phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1; + await store.storeTerminalResult(state.runId, output); + state = await store.transition(state.runId, "terminal_durable", { terminalResultSha256: registryDigest(output) as never }); + return output; + }})); return result; } - 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 async planForState(state: RegistryAddressedPublicationStateV1): Promise { + if (!state.targetCommit || !state.targetWorkspaces || !state.changedWorkspaceIds || !state.planSha256 || !state.changedSetSha256 || !state.targetManifestSha256 || !state.advertisedTargetCommit || !state.immutableTargetRef || !state.fetchedTargetCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + const common = { schemaVersion: 1 as const, installationIdentitySha256: state.installationIdentitySha256, repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit, immutableTargetRef: state.immutableTargetRef, fetchedTargetCommit: state.fetchedTargetCommit, targetCommit: state.targetCommit, targetManifestSha256: state.targetManifestSha256, targetWorkspaces: state.targetWorkspaces, changedWorkspaceIds: state.changedWorkspaceIds, changedSetSha256: state.changedSetSha256 }; + const plan: RegistryAddressedPlanV1 = state.operation === "registry_bootstrap" ? { ...common, operation: "registry_bootstrap", changedSetRule: "all_target_workspace_ids", baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] } : { ...common, operation: "registry_pull", changedSetRule: "symmetric_base_target_workspace_difference", baseCommit: state.baseCommit!, baseManifestSha256: state.baseManifestSha256!, baseWorkspaces: state.baseWorkspaces }; + if (registryDigest(plan) !== state.planSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt"); + return plan; } - private addressedSnapshot(state: ActiveState, request: any): RegistryAddressedSnapshotV1 { + + private addressedSnapshot(state: ActiveState, requestDigest: string): RegistryActiveSnapshotV1 { const ids = state.revisions.map(revision => revision.id).sort(); - 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 }; + 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 bootstrap(): Promise { - await this.repository.ensureLayout(); - return await this.lock.run(async () => { - try { - const status = await this.repository.bootstrap(); - await this.activate(status.head!); - return status; - } catch (error) { - return await this.gitFallback(error); - } - }); - } - - async pull(): Promise { - await this.repository.ensureLayout(); - return await this.lock.run(async () => { - try { - const status = await this.repository.pull(); - await this.activate(status.head!); - return status; - } catch (error) { - return await this.gitFallback(error); - } - }); - } - - async list(): Promise { - const active = await this.tryActiveState(); - if (active) return active.revisions; - // A clean installation has no active snapshot until the first registry operation. Keep - // this lazy so health/startup remain available when Git is temporarily unreachable, while - // still refusing corrupted existing state (tryActiveState throws instead of returning none). - await this.bootstrap(); - return (await this.activeState()).revisions; - } - - /** - * List every intact retained snapshot, current snapshots first. Session discovery and - * retention use this rather than only the active revision so removing a workspace from - * Git cannot strand a resumable session that still pins one of its older descriptors. - */ async listRetainedSnapshots(): Promise { await this.repository.ensureLayout(); return await this.lock.run(async () => { @@ -434,57 +474,6 @@ export class WorkspaceRegistry { * Publish canonical YAML and derived public documentation as one optimistic Git revision. * The browser never provides paths or generated artifacts; those are derived server-side. */ - async publish(request: PublishWorkspaceRequest): Promise { - await this.repository.ensureLayout(); - return await this.lock.run(async () => { - const status = await this.repository.pull(); - await this.activate(status.head!); - const current = await this.activeState(); - const id = request.action === "delete" ? request.id : request.workspace.workspace.id; - const existing = current.revisions.find((revision) => revision.id === id); - const local = request.action === "delete" ? undefined : request.workspace; - - if (request.baseCommit !== status.head || ( - request.action !== "create" && existing?.blob !== request.baseBlob - )) { - const contentOnlyStale = request.action !== "create" - && request.baseCommit !== status.head - && existing?.blob === request.baseBlob; - if (contentOnlyStale) { - throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale"); - } - throw await this.conflictFor(request, status.head!, existing, local); - } - if (request.action === "create" && existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "create" && !existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "delete") { - await this.assertEvidenceContext(request.workspace, status.head!); - } - - const yamlPath = workspacePath(id); - const docPaths = this.documentationPaths(id); - if (request.action === "delete") { - await this.repository.removeRegistryFile(yamlPath); - await this.repository.removeRegistryFile(docPaths.contract); - await this.repository.removeRegistryFile(docPaths.readme); - } else { - const canonical = request.workspace; - const source = serializeWorkspaceYaml(canonical); - const docs = renderWorkspaceDocs(canonical); - await this.repository.writeRegistryFile(yamlPath, source); - await this.repository.writeRegistryFile(docPaths.contract, docs.envExample); - await this.repository.writeRegistryFile(docPaths.readme, docs.markdown); - } - - const next = await this.repository.commitAndPush( - [yamlPath, docPaths.contract, docPaths.readme], - request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`, - ); - await this.activate(next.head!); - return (await this.activeState()).revisions.find((revision) => revision.id === id); - }); - } - private documentationPaths(id: string): { contract: string; readme: string } { workspacePath(id); const directory = `workspace-docs/${id}`; @@ -556,9 +545,9 @@ export class WorkspaceRegistry { } } - private async activate(commit: string): Promise { + private async materialize(commit: string): Promise { const safeHead = safeCommit(commit); - const files = await this.repository.workspacePaths(); + const files = await this.repository.workspacePathsAt(safeHead); const snapshots: Array<{ id: string; @@ -570,7 +559,7 @@ export class WorkspaceRegistry { try { for (const path of files) { const id = path.slice("workspaces/".length, -".yaml".length); - const source = await this.repository.readWorkspace(path); + const source = await this.repository.readWorkspaceAt(safeHead, path); const workspace = parseWorkspaceYaml(source); if (workspace.workspace.id !== id) { throw new WorkspaceRegistryError("workspace_invalid", "Workspace ID does not match its repository path"); @@ -588,7 +577,7 @@ export class WorkspaceRegistry { id, source: serializeWorkspaceYaml(workspace), workspace, - blob: await this.repository.blob(path), + blob: await this.repository.blobAt(safeHead, path), }); } } catch (error) { @@ -631,7 +620,19 @@ export class WorkspaceRegistry { } } - await this.writeActiveState({ head: safeHead, revisions }); + } + + private async publishMaterialized(commit: string): Promise { + const state = await this.snapshotState(commit); + await this.writeActiveState({ head: state.head, revisions: state.revisions }); + } + + private async baseState(commit: Revision40 | null): Promise { + return commit ? this.snapshotState(commit) : undefined; + } + + private manifestIdentity(revision: WorkspaceRevision): RegistryWorkspaceManifestIdentityV1 { + return { workspaceId: revision.id as RegistryWorkspaceManifestIdentityV1["workspaceId"], revision: revision.commit as RegistryWorkspaceManifestIdentityV1["revision"], descriptorBlob: revision.blob as RegistryWorkspaceManifestIdentityV1["descriptorBlob"], manifestSha256: registryDigest(revision) as RegistryWorkspaceManifestIdentityV1["manifestSha256"] }; } private async gitFallback(error: unknown): Promise { @@ -640,7 +641,7 @@ export class WorkspaceRegistry { const active = await this.tryActiveState(); if (!active) throw safeError; return { - branch: this.config.branch, + branch: this.repository.config.branch, head: active.head, ahead: 0, behind: 0,