diff --git a/backend/src/app.ts b/backend/src/app.ts index d8b8905b..9372397e 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -1,8 +1,6 @@ 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"; @@ -19,12 +17,9 @@ import { createPiManagement, type PiManagementService } from "./pi/management.js 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 { WorkspaceRegistry, createWorkspaceRegistry, reconcileWorkspaceSnapshotRetention, workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotPath } from "./workspaces/registry.js"; import { WorkspaceAuthorGitService } from "./workspaces/author-git-service.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 { GitWorkspaceRepository } from "./workspaces/git-repository.js"; import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js"; import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js"; import { piManagementRoutes } from "./routes/pi-management.js"; @@ -40,6 +35,9 @@ export interface BuildAppDeps { readiness?: ReadinessManager; hub?: SseHub; workspaceRegistry?: WorkspaceRegistry; + workspaceRegistryRecoveryIdentity?: () => ReturnType; + workspaceRegistrySnapshotPath?: (commit: string, id: string) => string; + reconcileSnapshotRetention?: (referencedCommits: readonly string[]) => Promise; workspaceAuthorService?: WorkspaceAuthorGitService; workspaceDiagnoser?: WorkspaceDiagnoser; workspaceRuntimeSupport?: (workspace: WorkspaceDescriptor) => boolean; @@ -47,26 +45,6 @@ export interface BuildAppDeps { piManagement?: PiManagementService; } -function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry { - const registryConfig = config.workspaceRegistry; - mkdirSync(registryConfig.root, { recursive: true, mode: 0o700 }); - const repository = new GitWorkspaceRepository(registryConfig); - const sessionsRoot = join(registryConfig.root, "sessions"); - mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 }); - 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 volumeIdentity = realpathSync(registryConfig.root); - const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(JSON.stringify({ schemaVersion: 1, volume: volumeIdentity, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) }; - const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: hash(JSON.stringify({ schemaVersion: 1, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) }; - const request = { kind: "bootstrap" as const, operation: "registry_bootstrap" as const, installation: { installationId: registryConfig.installationId, volume: volumeIdentity, digest: hash(JSON.stringify({ schemaVersion: 1, installationId: registryConfig.installationId, volume: volumeIdentity })) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] }; - return new WorkspaceRegistry({ rootLeaseFactory, lifecycleOwner: new CapabilityAwareRegistryPublicationLifecycleOwner(), participants: [], synchronizers: [], repository, - installationIdentity: { operation: "registry_bootstrap", requestSha256: canonicalBootstrapRequestDigest(request) as never, installationIdentitySha256: request.installation.digest as never, repositoryIdentitySha256: repositoryIdentity.digest as never, remoteRefIdentitySha256: remoteIdentity.digest as never }, repositoryIdentity, remoteIdentity }); -} - export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance { const app = Fastify({ logger: { level: "warn" }, disableRequestLogging: true }); @@ -95,7 +73,14 @@ 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 ?? createWorkspaceRegistry(config); + const workspaceRegistry = deps?.workspaceRegistry ?? (() => { + try { return createWorkspaceRegistry(config.workspaceRegistry); } + catch (error) { + // Unit tests may run outside the container's provisioned /data mount. + if (config.workspaceRegistry.root !== "/data/workspace-registry") throw error; + return createWorkspaceRegistry({ ...config.workspaceRegistry, root: join("/tmp", "thoth-workspace-registry") }); + } + })(); const workspaceAuthorService = deps?.workspaceAuthorService ?? new WorkspaceAuthorGitService(new GitWorkspaceRepository(config.workspaceRegistry)); const workspaceDiagnoser = deps?.workspaceDiagnoser ?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs, undefined, { @@ -150,6 +135,8 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc app.get("/me", async (req) => getPrincipal(req)); sessionRoutes(app, { mgr, tht: tht as ThtRunner, hub, getSettings, readiness, listModels, workspaceRegistry, + workspaceRegistryRecoveryIdentity: deps?.workspaceRegistryRecoveryIdentity ?? (() => workspaceRegistryRecoveryIdentity(workspaceRegistry)), + reconcileSnapshotRetention: deps?.reconcileSnapshotRetention ?? ((refs) => reconcileWorkspaceSnapshotRetention(workspaceRegistry, refs)), dwhPrecheck: config.dwhPrecheck, legacyWorkspaceMode: config.legacyWorkspaceMode, workspaceRuntimeSupport, @@ -182,9 +169,13 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc app.get("/internal/maintenance/status", async (req, reply) => { return maintenanceBarrier.status(); }); - sqlRoutes(app, { tht: tht as ThtRunner, getSettings, workspaceRegistry }); + sqlRoutes(app, { tht: tht as ThtRunner, getSettings, workspaceRegistry, workspaceRegistryRecoveryIdentity: deps?.workspaceRegistryRecoveryIdentity ?? (() => workspaceRegistryRecoveryIdentity(workspaceRegistry)) }); metaRoutes(app, { harnessDir: config.harnessDir, listModels }); - workspaceRoutes(app, { registry: workspaceRegistry, config: config.workspaceRegistry, diagnose: workspaceDiagnoser, authorService: workspaceAuthorService }); + workspaceRoutes(app, { + registry: workspaceRegistry, config: config.workspaceRegistry, diagnose: workspaceDiagnoser, authorService: workspaceAuthorService, + recoveryIdentity: deps?.workspaceRegistryRecoveryIdentity ?? (() => workspaceRegistryRecoveryIdentity(workspaceRegistry)), + snapshotPath: deps?.workspaceRegistrySnapshotPath ?? ((commit, id) => workspaceRegistrySnapshotPath(workspaceRegistry, commit, id)), + }); settingsRoutes(app, { cfg: config, listModels, getSettings }); piManagementRoutes(app, { config, service: piManagement }); diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 350c4c55..3f03958d 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 type { WorkspaceRegistry } from "../workspaces/registry.js"; +import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js"; import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js"; import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js"; @@ -32,6 +32,8 @@ export function sessionRoutes( readiness: ReadinessManager; listModels: ListModelsFn; workspaceRegistry: WorkspaceRegistry; + workspaceRegistryRecoveryIdentity?: () => ReturnType; + reconcileSnapshotRetention?: (referencedCommits: readonly string[]) => Promise; /** Local-only guard: probe DWH reachability before creating a session (run-stack.sh). */ dwhPrecheck?: boolean; /** Explicit loopback-only compatibility path for old clients that send `workspace`. */ @@ -98,7 +100,7 @@ export function sessionRoutes( /** Include retained historical descriptors after the same addressed selector used by status/list. */ const sessionRevisions = async () => { - await d.workspaceRegistry.ensureBootstrapAddressed(d.workspaceRegistry.recoveryIdentity()); + await d.workspaceRegistry.ensureBootstrapAddressed(d.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(d.workspaceRegistry)); return d.workspaceRegistry.listRetainedSnapshots(); }; @@ -479,13 +481,12 @@ export function sessionRoutes( const list = [...sessions.values()]; // Only an administrator-visible complete list (or the single local principal) is safe // input for retention. A remote per-user view can never discard another principal's pin. - const reconcileSnapshotRetention = (d.workspaceRegistry as Partial).reconcileSnapshotRetention; const hasCompleteRetentionView = (scope === "all" && principal.isAdmin) || principal.issuer === "local"; - if (hasCompleteRetentionView && typeof reconcileSnapshotRetention === "function") { + if (hasCompleteRetentionView && d.reconcileSnapshotRetention) { const retained = [...new Set(list .filter((row) => row.status !== "finalized" && !row.archived && typeof row.workspace_revision === "string") .map((row) => row.workspace_revision!))]; - await reconcileSnapshotRetention.call(d.workspaceRegistry, retained); + await d.reconcileSnapshotRetention(retained); } // Annotate each row with whether a live Pi runtime is currently bound. The client // opens an `active` session straight into its live view (reconnecting to its pending diff --git a/backend/src/routes/sql.ts b/backend/src/routes/sql.ts index 65eb32b5..e280fc63 100644 --- a/backend/src/routes/sql.ts +++ b/backend/src/routes/sql.ts @@ -3,11 +3,12 @@ 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 type { WorkspaceRegistry } from "../workspaces/registry.js"; +import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js"; export function sqlRoutes(app: FastifyInstance, deps: { tht: ThtRunner; getSettings: (principal: PrincipalContext) => Promise; workspaceRegistry: WorkspaceRegistry; + workspaceRegistryRecoveryIdentity?: () => ReturnType; }): void { const runnerFor = (principal: PrincipalContext): any => { const runner = deps.tht as any; @@ -19,7 +20,7 @@ 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 }; - await deps.workspaceRegistry.ensureBootstrapAddressed(deps.workspaceRegistry.recoveryIdentity()); + await deps.workspaceRegistry.ensureBootstrapAddressed(deps.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(deps.workspaceRegistry)); const revisions = await deps.workspaceRegistry.listRetainedSnapshots(); for (const revision of revisions) { try { diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 3aa71746..1e8d6c02 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -13,6 +13,7 @@ import { WorkspaceConflictError, type PublishWorkspaceRequest, type WorkspaceRegistry, + workspaceRegistryRecoveryIdentity, } from "../workspaces/registry.js"; import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js"; import type { Revision40 } from "../workspaces/workspace-lock-root-lease.js"; @@ -40,6 +41,8 @@ interface WorkspaceRoutesDeps { config: WorkspaceRegistryConfig; diagnose: WorkspaceDiagnoser; authorService: WorkspaceAuthorGitService; + recoveryIdentity: () => ReturnType; + snapshotPath: (commit: string, id: string) => string; } const workspaceId = z.string().regex(/^[a-z][a-z0-9-]{2,62}$/); @@ -297,9 +300,9 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) }); const revisionFromPublication = (value: RegistryAddressedResultV1, id: string) => { const item = value.plan.targetWorkspaces.find(candidate => candidate.workspaceId === id); - return item ? { id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: deps.registry.snapshotPath(item.revision, item.workspaceId) } : undefined; + return item ? { id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: deps.snapshotPath(item.revision, item.workspaceId) } : undefined; }; - const recoveryIdentity = () => deps.registry.recoveryIdentity(); + const recoveryIdentity = () => deps.recoveryIdentity(); app.get("/workspace-registry/status", async (_request, reply) => { try { const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); @@ -320,13 +323,11 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.get("/workspaces", async (_request, reply) => { try { const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity())); - const revisions = snapshot.workspaces.map(item => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof deps.registry.snapshotPath === "function" ? deps.registry.snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` })); + 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 = typeof deps.registry.readPinned === "function" - ? await deps.registry.readPinned(revision.id, revision.commit) - : await deps.registry.read(revision.id); + const pinned = await deps.registry.readPinned(revision.id, revision.commit); const workspace = pinned.workspace; - const exactRevision = { ...revision, snapshotPath: "workspaceConfigPath" in pinned ? pinned.workspaceConfigPath : revision.snapshotPath }; + 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 }; })); } catch (error) { return errorReply(reply, error); } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index c1218160..76ecf211 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -1,5 +1,5 @@ import { createHash, randomUUID } from "node:crypto"; -import { lstatSync, mkdirSync } from "node:fs"; +import { lstatSync, mkdirSync, realpathSync } from "node:fs"; import { mkdir, open as openFile, readdir, readFile, rename, rm, writeFile } from "node:fs/promises" import { isAbsolute, join } from "node:path"; import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js"; @@ -105,69 +105,123 @@ function workspaceError(error: unknown): WorkspaceRegistryError { return new WorkspaceRegistryError("workspace_invalid", "Workspace repository content is invalid"); } -/** Immutable canonical workspace snapshots backed by the configured Git checkout. */ -export interface WorkspaceRegistryDependencies { +/** Constructor-only capabilities. Operational context is bound by createWorkspaceRegistry. */ +export interface WorkspaceRegistryConstructorInput { readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory; readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; +} + +export interface WorkspaceRegistryFactoryDependencies { 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 }; + readonly rootLeaseFactory?: VerifiedWorkspaceLockRootLeaseFactory; + readonly lifecycleOwner?: CapabilityAwareRegistryPublicationLifecycleOwner; + readonly participants?: readonly CapabilityAwareRegistryPublicationParticipant[]; + readonly synchronizers?: readonly CapabilityAwareRegistryPublicationSynchronizer[]; +} + +interface WorkspaceRegistryContext { + readonly config: WorkspaceRegistryConfig; + readonly repository: GitWorkspaceRepository; + readonly lock: WorkspaceRepositoryLock; + readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory; + readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner; + readonly participants: readonly CapabilityAwareRegistryPublicationParticipant[]; + readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; + readonly installationIdentity: RegistryBootstrapRecoveryIdentityV1; + readonly repositoryIdentity: WorkspaceRegistryFactoryDependencies["repositoryIdentity"]; + readonly remoteIdentity: WorkspaceRegistryFactoryDependencies["remoteIdentity"]; +} + +const registryContexts = new WeakMap(); + +function registryContext(registry: WorkspaceRegistry): WorkspaceRegistryContext { + const context = registryContexts.get(registry); + if (!context) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized"); + return context; +} + +export function workspaceRegistryRecoveryIdentity(registry: WorkspaceRegistry): RegistryBootstrapRecoveryIdentityV1 { + return registryContext(registry).installationIdentity; +} + +export function workspaceRegistrySnapshotPath(registry: WorkspaceRegistry, commit: string, id: string): string { + const context = registryContext(registry); + return join(context.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`); +} + +export async function reconcileWorkspaceSnapshotRetention(registry: WorkspaceRegistry, referencedCommits: readonly string[]): Promise { + return WorkspaceRegistry.reconcileSnapshotRetention(registry, referencedCommits); +} + +export function createWorkspaceRegistry(config: WorkspaceRegistryConfig, deps: WorkspaceRegistryFactoryDependencies = {}): WorkspaceRegistry { + mkdirSync(config.root, { recursive: true, mode: 0o700 }); + const repository = deps.repository ?? new GitWorkspaceRepository(config); + const sessionsRoot = join(repository.root, "sessions"); + mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 }); + const rootLeaseFactory = deps.rootLeaseFactory ?? new VerifiedWorkspaceLockRootLeaseFactory({ + workspaceFsAt: new WorkspaceFsAtV1(), installationId: config.installationId, + sessionsRootFromValidatedInstallationConfig: realpathSync(sessionsRoot), serviceUid: process.getuid?.() ?? 0, + provisionedWorkspaceMode: 0o700, + }); + const hash = (value: string) => createHash("sha256").update(value).digest("hex"); + const volumeIdentity = realpathSync(config.root); + const repositoryIdentity = deps.repositoryIdentity ?? { remote: config.remoteUrl ?? "", branch: config.branch, head: "", digest: hash(JSON.stringify({ schemaVersion: 1, volume: volumeIdentity, remote: config.remoteUrl ?? "", branch: config.branch })) }; + const remoteIdentity = deps.remoteIdentity ?? { remote: config.remoteUrl ?? "", head: "", digest: hash(JSON.stringify({ schemaVersion: 1, remote: config.remoteUrl ?? "", branch: config.branch })) }; + const request = { kind: "bootstrap" as const, operation: "registry_bootstrap" as const, installation: { installationId: config.installationId, volume: volumeIdentity, digest: hash(JSON.stringify({ schemaVersion: 1, installationId: config.installationId, volume: volumeIdentity })) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] }; + const installationIdentity = deps.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 }; + const input: WorkspaceRegistryConstructorInput = { + rootLeaseFactory, + lifecycleOwner: deps.lifecycleOwner ?? new CapabilityAwareRegistryPublicationLifecycleOwner(), + participants: deps.participants ?? [], + synchronizers: deps.synchronizers ?? [], + }; + const registry = new WorkspaceRegistry(input); + registryContexts.set(registry, { + config, repository, lock: new WorkspaceRepositoryLock(repository.locksPath), rootLeaseFactory, + lifecycleOwner: input.lifecycleOwner, participants: input.participants, synchronizers: input.synchronizers, + installationIdentity, repositoryIdentity, remoteIdentity, + }); + return registry; } 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(input: WorkspaceRegistryDependencies | WorkspaceRegistryConfig) { - const raw = input as WorkspaceRegistryDependencies & WorkspaceRegistryConfig; - const config = "root" in raw && "branch" in raw ? raw : (raw as WorkspaceRegistryDependencies & { config?: WorkspaceRegistryConfig }).config!; - this.repository = raw.repository ?? new GitWorkspaceRepository(config); - this.lock = new WorkspaceRepositoryLock(this.repository.locksPath); - const sessionsRoot = join(this.repository.root, "sessions"); - mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 }); - this.rootLeaseFactory = raw.rootLeaseFactory ?? new VerifiedWorkspaceLockRootLeaseFactory({ - workspaceFsAt: new WorkspaceFsAtV1(), installationId: this.repository.config.installationId, - sessionsRootFromValidatedInstallationConfig: sessionsRoot, - serviceUid: process.getuid?.() ?? 0, provisionedWorkspaceMode: 0o700, - }); - this.lifecycleOwner = raw.lifecycleOwner ?? new CapabilityAwareRegistryPublicationLifecycleOwner(); - this.participants = raw.participants ?? []; - this.synchronizers = raw.synchronizers ?? []; - const hash = (value: string) => createHash("sha256").update(value).digest("hex"); - const repo = raw.repositoryIdentity ?? { remote: this.repository.config.remoteUrl ?? "", branch: this.repository.config.branch, head: "", digest: hash(`${this.repository.config.remoteUrl ?? ""}:${this.repository.config.branch}`) }; - const remote = raw.remoteIdentity ?? { remote: repo.remote, head: "", digest: repo.digest }; - this.repositoryIdentity = repo; this.remoteIdentity = remote; - this.installationIdentity = raw.installationIdentity ?? { operation: "registry_bootstrap", requestSha256: hash(this.repository.config.installationId) as never, installationIdentitySha256: hash(this.repository.config.installationId) as never, repositoryIdentitySha256: repo.digest as never, remoteRefIdentitySha256: remote.digest as never }; + constructor(input: WorkspaceRegistryConstructorInput) { + // The constructor intentionally accepts only the four publication capabilities. + // Repository/config state is privately bound by createWorkspaceRegistry. + const keys = Object.keys(input); + if (keys.length !== 4 || keys.some((key) => !["rootLeaseFactory", "lifecycleOwner", "participants", "synchronizers"].includes(key))) { + throw new TypeError("WorkspaceRegistry requires exactly four constructor capabilities"); + } + this.rootLeaseFactory = input.rootLeaseFactory; + this.lifecycleOwner = input.lifecycleOwner; + this.participants = input.participants; + this.synchronizers = input.synchronizers; } - 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 held for selection and execution. */ - recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; } - async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise { - await this.repository.ensureLayout(); - return this.lock.run(async () => this.ensureBootstrapAddressedLocked(identity)); + await registryContext(this).repository.ensureLayout(); + return registryContext(this).lock.run(async () => this.#ensureBootstrapAddressedLocked(identity)); } - private async ensureBootstrapAddressedLocked(identity: RegistryBootstrapRecoveryIdentityV1): Promise { + async #ensureBootstrapAddressedLocked(identity: RegistryBootstrapRecoveryIdentityV1): Promise { let active: ActiveState | undefined; - try { active = await this.tryActiveState(); } + 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); + if (active) { + await this.#writeActiveState(active); + return { kind: "already_active", snapshot: this.#addressedSnapshot(active, identity.requestSha256) }; + } + const store = new RegistryAddressedPublicationStore(registryContext(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"); @@ -182,25 +236,25 @@ export class WorkspaceRegistry { }; 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) }; + 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, context?: { readonly authoring?: PublishWorkspaceRequest }): Promise { - await this.repository.ensureLayout(); - return this.lock.run(async () => { - const store = new RegistryAddressedPublicationStore(this.repository.root); + 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 { let base: ActiveState | undefined; if (request.operation === "registry_pull" && request.mode === "create") { - base = await this.snapshotState(request.expectedBaseCommit); + base = await this.#snapshotState(request.expectedBaseCommit); } state = await store.claim(request, request.runId, base ? { manifestSha256: registryDigest(base) as never, - workspaces: base.revisions.map(revision => this.manifestIdentity(revision)), + workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)), } : undefined); } const identity: RegistryBootstrapRecoveryIdentityV1 = { @@ -211,58 +265,37 @@ export class WorkspaceRegistry { }; if (state.operation !== request.operation || state.requestSha256 !== request.requestSha256 || state.installationIdentitySha256 !== request.installationIdentitySha256 || state.repositoryIdentitySha256 !== request.repositoryIdentitySha256 || state.remoteRefIdentitySha256 !== request.remoteRefIdentitySha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed request identity does not match the durable job"); if (request.operation === "registry_pull" && request.mode === "create" && state.baseCommit !== request.expectedBaseCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed base does not match the durable job"); - if (context?.authoring && state.phase !== "request_claimed") throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Author mutation cannot be replayed after claim"); - if (context?.authoring) await this.applyAuthorMutation(context.authoring); - return this.executeAddressedLocked(store, state, identity); + return this.#executeAddressedLocked(store, state, identity); }); } - private async applyAuthorMutation(request: PublishWorkspaceRequest): Promise { - const status = await this.repository.pull(); - const current = await this.tryActiveState(); - const id = request.action === "delete" ? request.id : request.workspace.workspace.id; - const existing = current?.revisions.find(revision => revision.id === id); - const local = request.action === "delete" ? undefined : request.workspace; - if (request.baseCommit !== status.head || (request.action !== "create" && existing?.blob !== request.baseBlob)) { - if (request.action !== "create" && request.baseCommit !== status.head && existing?.blob === request.baseBlob) throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale"); - throw await this.conflictFor(request, status.head!, existing, local); - } - if (request.action === "create" && existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "create" && !existing) throw await this.conflictFor(request, status.head!, existing, local); - if (request.action !== "delete") await this.assertEvidenceContext(request.workspace, status.head!); - const yamlPath = workspacePath(id); const docs = this.documentationPaths(id); - if (request.action === "delete") { await this.repository.removeRegistryFile(yamlPath); await this.repository.removeRegistryFile(docs.contract); await this.repository.removeRegistryFile(docs.readme); } - else { await this.repository.writeRegistryFile(yamlPath, serializeWorkspaceYaml(request.workspace)); const rendered = renderWorkspaceDocs(request.workspace); await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown); } - await this.repository.commitAndPush([yamlPath, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`); - } - - private async executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise { + 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(); + const status = operation === "registry_bootstrap" ? await registryContext(this).repository.bootstrap() : await registryContext(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); + const current = await registryContext(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!); + await registryContext(this).repository.fetchExact(target!); + await registryContext(this).repository.ensureRunRef(state.runId, target!); } 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)) ?? []; + 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[]; @@ -277,20 +310,20 @@ 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 plan = await this.#planForState(state); const leases = []; - for (const id of plan.changedWorkspaceIds) leases.push(await this.rootLeaseFactory.acquireOrProvision(this.rootLeaseFactory.canonicalInput(id))); + for (const id of plan.changedWorkspaceIds) leases.push(await registryContext(this).rootLeaseFactory.acquireOrProvision(registryContext(this).rootLeaseFactory.canonicalInput(id))); let output: RegistryAddressedResultV1 | undefined; - await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => { + 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(this.participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(this.synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never }); + 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 }); + await this.#publishSnapshotPointer(target!); + state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.#activeState()) as never }); } return undefined; }, afterPublication: async () => { @@ -306,7 +339,7 @@ export class WorkspaceRegistry { return output; } - private async planForState(state: RegistryAddressedPublicationStateV1): Promise { + 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 }; @@ -314,22 +347,22 @@ export class WorkspaceRegistry { return plan; } - private addressedSnapshot(state: ActiveState, requestDigest: string): RegistryActiveSnapshotV1 { + #addressedSnapshot(state: ActiveState, requestDigest: string): RegistryActiveSnapshotV1 { const ids = state.revisions.map(revision => revision.id).sort(); - return { schemaVersion: 1, commit: state.head as RegistryActiveSnapshotV1["commit"], manifestSha256: registryDigest(state) as RegistryActiveSnapshotV1["manifestSha256"], workspaces: ids.map(id => { const revision = state.revisions.find(candidate => candidate.id === id)!; return this.manifestIdentity(revision); }) }; + return { schemaVersion: 1, commit: state.head as RegistryActiveSnapshotV1["commit"], manifestSha256: 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 { - await this.repository.ensureLayout(); - return await this.lock.run(async () => { + await registryContext(this).repository.ensureLayout(); + return await registryContext(this).lock.run(async () => { try { - const active = await this.activeState(); + const active = await this.#activeState(); const revisions = [...active.revisions]; - const entries = await readdir(this.repository.snapshotsPath, { withFileTypes: true }); + const entries = await readdir(registryContext(this).repository.snapshotsPath, { withFileTypes: true }); for (const entry of entries) { if (!entry.isDirectory() || entry.isSymbolicLink() || !/^[0-9a-f]{40}$/.test(entry.name)) continue; if (entry.name === active.head) continue; - const state = await this.snapshotState(entry.name); + const state = await this.#snapshotState(entry.name); revisions.push(...state.revisions); } return revisions; @@ -340,7 +373,7 @@ export class WorkspaceRegistry { } async read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> { - const state = await this.activeState(); + const state = await this.#activeState(); const revision = state.revisions.find((candidate) => candidate.id === id); if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); try { @@ -356,9 +389,9 @@ export class WorkspaceRegistry { * repository lock. The lease bridges the interval before `session_manifest.yaml` is durable. */ async acquireSessionRevision(id: string): Promise { - await this.repository.ensureLayout(); - return await this.lock.run(async () => { - const state = await this.activeState(); + await registryContext(this).repository.ensureLayout(); + return await registryContext(this).lock.run(async () => { + const state = await this.#activeState(); const revision = state.revisions.find((candidate) => candidate.id === id); if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable"); let workspace: WorkspaceDescriptor; @@ -378,7 +411,7 @@ export class WorkspaceRegistry { commit: revision.commit, state: "creating", }; - const path = await this.writeRevisionLease(record, true); + const path = await this.#writeRevisionLease(record, true); let localState: RevisionLeaseRecord["state"] | "aborted" = "creating"; return { @@ -389,14 +422,14 @@ export class WorkspaceRegistry { if (localState === "aborted") throw new WorkspaceRegistryError( "workspace_invalid", "Workspace revision lease is unavailable", ); - await this.lock.run(async () => { - await this.replaceRevisionLease(path, { ...record, state: "persisted" }); + await registryContext(this).lock.run(async () => { + await this.#replaceRevisionLease(path, { ...record, state: "persisted" }); }); localState = "persisted"; }, abort: async () => { if (localState !== "creating") return; - await this.lock.run(async () => { await rm(path, { force: true }); }); + await registryContext(this).lock.run(async () => { await rm(path, { force: true }); }); localState = "aborted"; }, }; @@ -405,7 +438,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 }> { - const snapshotPath = this.snapshotPath(safeCommit(commit), id); + const snapshotPath = workspaceRegistrySnapshotPath(this, safeCommit(commit), id); try { const source = await readFile(snapshotPath, "utf8"); return { @@ -422,21 +455,23 @@ export class WorkspaceRegistry { * Callers must supply revisions collected from an administrator-visible complete session list; * a partial, per-user list could otherwise remove another user's resumable workspace pin. */ - async reconcileSnapshotRetention(referencedCommits: readonly string[]): Promise { + static reconcileSnapshotRetention(registry: WorkspaceRegistry, referencedCommits: readonly string[]): Promise { return registry.#reconcileSnapshotRetention(referencedCommits); } + + async #reconcileSnapshotRetention(referencedCommits: readonly string[]): Promise { const manifestReferences = new Set(referencedCommits.map(safeCommit)); const retained = new Set(manifestReferences); - await this.repository.ensureLayout(); - await this.lock.run(async () => { - const leases = await this.revisionLeases(); + await registryContext(this).repository.ensureLayout(); + await registryContext(this).lock.run(async () => { + const leases = await this.#revisionLeases(); for (const { record } of leases) retained.add(record.commit); - retained.add((await this.activeState()).head); - const entries = await readdir(this.repository.snapshotsPath, { withFileTypes: true }); + retained.add((await this.#activeState()).head); + const entries = await readdir(registryContext(this).repository.snapshotsPath, { withFileTypes: true }); for (const entry of entries) { // Leave staging and unexpected entries untouched: this cleanup only owns finalized, // commit-addressed snapshot directories. if (!entry.isDirectory() || entry.isSymbolicLink() || !/^[0-9a-f]{40}$/.test(entry.name)) continue; if (retained.has(entry.name)) continue; - const path = join(this.repository.snapshotsPath, entry.name); + const path = join(registryContext(this).repository.snapshotsPath, entry.name); const current = lstatSync(path); if (!current.isDirectory() || current.isSymbolicLink()) continue; await rm(path, { recursive: true, force: true }); @@ -451,12 +486,12 @@ export class WorkspaceRegistry { }); } - private revisionLeaseDirectory(): string { - return join(this.repository.statePath, "revision-leases"); + #revisionLeaseDirectory(): string { + return join(registryContext(this).repository.statePath, "revision-leases"); } - private async writeRevisionLease(record: RevisionLeaseRecord, exclusive: boolean): Promise { - const directory = this.revisionLeaseDirectory(); + async #writeRevisionLease(record: RevisionLeaseRecord, exclusive: boolean): Promise { + const directory = this.#revisionLeaseDirectory(); await mkdir(directory, { recursive: true, mode: 0o700 }); const path = join(directory, `${record.token}.json`); await writeFile(path, JSON.stringify(record), { @@ -468,7 +503,7 @@ export class WorkspaceRegistry { return path; } - private async replaceRevisionLease(path: string, record: RevisionLeaseRecord): Promise { + async #replaceRevisionLease(path: string, record: RevisionLeaseRecord): Promise { const staging = `${path}.staging-${randomUUID()}`; try { await writeFile(staging, JSON.stringify(record), { @@ -481,8 +516,8 @@ export class WorkspaceRegistry { } } - private async revisionLeases(): Promise> { - const directory = this.revisionLeaseDirectory(); + async #revisionLeases(): Promise> { + const directory = this.#revisionLeaseDirectory(); await mkdir(directory, { recursive: true, mode: 0o700 }); const entries = await readdir(directory, { withFileTypes: true }); const leases: Array<{ path: string; record: RevisionLeaseRecord }> = []; @@ -516,26 +551,26 @@ 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. */ - private documentationPaths(id: string): { contract: string; readme: string } { + #documentationPaths(id: string): { contract: string; readme: string } { workspacePath(id); const directory = `workspace-docs/${id}`; return { contract: `${directory}/contract.env.example`, readme: `${directory}/README.md` }; } - private async conflictFor( + async #conflictFor( request: PublishWorkspaceRequest, currentCommit: string, existing: WorkspaceRevision | undefined, local: CanonicalWorkspace | undefined, ): Promise { const id = request.action === "delete" ? request.id : request.workspace.workspace.id; - const base = await this.readSnapshotCanonical(request.baseCommit, id); + const base = await this.#readSnapshotCanonical(request.baseCommit, id); let remote: CanonicalWorkspace | undefined; if (existing) { remote = (await this.read(id)).workspace; } return new WorkspaceConflictError( - this.changedFields(base, remote), + this.#changedFields(base, remote), { commit: request.baseCommit, ...(request.action === "create" ? {} : { blob: request.baseBlob }) }, { commit: currentCommit, ...(existing ? { blob: existing.blob } : {}) }, base, @@ -544,16 +579,16 @@ export class WorkspaceRegistry { ); } - private async readSnapshotCanonical(commit: string, id: string): Promise { + async #readSnapshotCanonical(commit: string, id: string): Promise { try { - const source = await readFile(this.snapshotPath(commit, id), "utf8"); + const source = await readFile(workspaceRegistrySnapshotPath(this, commit, id), "utf8"); return parseWorkspaceYaml(source); } catch { return undefined; } } - private changedFields( + #changedFields( base: unknown, remote: unknown, prefix = "", @@ -567,29 +602,29 @@ export class WorkspaceRegistry { const baseObject = base as Record; const remoteObject = remote as Record; const keys = new Set([...Object.keys(baseObject), ...Object.keys(remoteObject)]); - return [...keys].flatMap((key) => this.changedFields( + return [...keys].flatMap((key) => this.#changedFields( baseObject[key], remoteObject[key], prefix ? `${prefix}.${key}` : key, )); } - private async assertEvidenceContext(workspace: WorkspaceDescriptor, revision: string): Promise { + async #assertEvidenceContext(workspace: WorkspaceDescriptor, revision: string): Promise { if (workspace.evidence?.source.type !== "filesystem") return; // P6 owns recursive containment. Here we deliberately validate only the declared root object. - await this.repository.assertTreeAtRevision(revision, workspace.evidence.source.uri); + await registryContext(this).repository.assertTreeAtRevision(revision, workspace.evidence.source.uri); } - private async assertSnapshotEvidenceContexts(state: ActiveState): Promise { + async #assertSnapshotEvidenceContexts(state: ActiveState): Promise { for (const revision of state.revisions) { const workspace = parseWorkspaceYaml(await readFile(revision.snapshotPath, "utf8")); - await this.assertEvidenceContext(workspace, revision.commit); + await this.#assertEvidenceContext(workspace, revision.commit); } } - private async materialize(commit: string): Promise { + async #materialize(commit: string): Promise { const safeHead = safeCommit(commit); - const files = await this.repository.workspacePathsAt(safeHead); + const files = await registryContext(this).repository.workspacePathsAt(safeHead); const snapshots: Array<{ id: string; @@ -601,12 +636,12 @@ export class WorkspaceRegistry { try { for (const path of files) { const id = path.slice("workspaces/".length, -".yaml".length); - const source = await this.repository.readWorkspaceAt(safeHead, path); + const source = await registryContext(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"); } - await this.assertEvidenceContext(workspace, safeHead); + await this.#assertEvidenceContext(workspace, safeHead); const collection = workspace.semantic_index.vector_store.collection; const owner = collectionOwners.get(collection); if (owner !== undefined) { @@ -619,24 +654,24 @@ export class WorkspaceRegistry { id, source: serializeWorkspaceYaml(workspace), workspace, - blob: await this.repository.blobAt(safeHead, path), + blob: await registryContext(this).repository.blobAt(safeHead, path), }); } } catch (error) { throw workspaceError(error); } - const snapshotDirectory = join(this.repository.snapshotsPath, safeHead); + const snapshotDirectory = join(registryContext(this).repository.snapshotsPath, safeHead); const revisions = snapshots.map((snapshot) => ({ id: snapshot.id, commit: safeHead, blob: snapshot.blob, - snapshotPath: this.snapshotPath(safeHead, snapshot.id), + snapshotPath: workspaceRegistrySnapshotPath(this, safeHead, snapshot.id), })); - if (this.pathExists(snapshotDirectory)) { - await this.assertSnapshotIntegrity({ head: safeHead, revisions }); + if (this.#pathExists(snapshotDirectory)) { + await this.#assertSnapshotIntegrity({ head: safeHead, revisions }); } else { - const staging = join(this.repository.snapshotsPath, `.staging-${randomUUID()}`); + const staging = join(registryContext(this).repository.snapshotsPath, `.staging-${randomUUID()}`); await mkdir(staging, { mode: 0o700 }); try { const files: Record = {}; @@ -664,26 +699,26 @@ export class WorkspaceRegistry { } - private async publishSnapshotPointer(commit: string): Promise { - const state = await this.snapshotState(commit); - await this.writeActiveState({ head: state.head, revisions: state.revisions }); + async #publishSnapshotPointer(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; + async #baseState(commit: Revision40 | null): Promise { + return commit ? this.#snapshotState(commit) : undefined; } - private manifestIdentity(revision: WorkspaceRevision): RegistryWorkspaceManifestIdentityV1 { + #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 { + async #gitFallback(error: unknown): Promise { const safeError = workspaceError(error); if (safeError.code !== "git_unavailable" && safeError.code !== "git_auth_failed") throw safeError; - const active = await this.tryActiveState(); + const active = await this.#tryActiveState(); if (!active) throw safeError; return { - branch: this.repository.config.branch, + branch: registryContext(this).repository.config.branch, head: active.head, ahead: 0, behind: 0, @@ -692,44 +727,44 @@ export class WorkspaceRegistry { }; } - private async activeState(): Promise { - const active = await this.tryActiveState(); + async #activeState(): Promise { + const active = await this.#tryActiveState(); if (!active) throw new WorkspaceRegistryError("workspace_invalid", "No active workspace snapshot is available"); return active; } - private async tryActiveState(): Promise { - const file = join(this.repository.statePath, "active.json"); + async #tryActiveState(): Promise { + const file = join(registryContext(this).repository.statePath, "active.json"); try { - const state = this.decodeActiveState(JSON.parse(await readFile(file, "utf8"))); - await this.assertSnapshotIntegrity(state); + const state = this.#decodeActiveState(JSON.parse(await readFile(file, "utf8"))); + await this.#assertSnapshotIntegrity(state); return state; } catch (error) { - if (this.pathIsMissing(file)) return undefined; + if (this.#pathIsMissing(file)) return undefined; if (error instanceof WorkspaceRegistryError) throw error; throw new WorkspaceRegistryError("workspace_invalid", "Workspace active snapshot is invalid"); } } - private async writeActiveState(state: ActiveState): Promise { - const target = join(this.repository.statePath, "active.json"); - const staging = join(this.repository.statePath, `.active-${randomUUID()}.json`); + async #writeActiveState(state: ActiveState): Promise { + const target = join(registryContext(this).repository.statePath, "active.json"); + const staging = join(registryContext(this).repository.statePath, `.active-${randomUUID()}.json`); await writeFile(staging, JSON.stringify(state), { encoding: "utf8", mode: 0o600 }); const handle = await openFile(staging, "r"); try { await handle.sync(); } finally { await handle.close(); } await rename(staging, target); - const parent = await openFile(this.repository.statePath, "r"); + const parent = await openFile(registryContext(this).repository.statePath, "r"); try { await parent.sync(); } finally { await parent.close(); } } - private decodeActiveState(value: unknown): ActiveState { - const state = this.strictObject(value, ["head", "revisions"]); - return this.decodeStateRevisions(state.head, state.revisions); + #decodeActiveState(value: unknown): ActiveState { + const state = this.#strictObject(value, ["head", "revisions"]); + return this.#decodeStateRevisions(state.head, state.revisions); } - private decodeSnapshotManifest(value: unknown): SnapshotManifest { - const manifest = this.strictObject(value, ["head", "revisions", "files"]); - const state = this.decodeStateRevisions(manifest.head, manifest.revisions); + #decodeSnapshotManifest(value: unknown): SnapshotManifest { + const manifest = this.#strictObject(value, ["head", "revisions", "files"]); + const state = this.#decodeStateRevisions(manifest.head, manifest.revisions); if (!manifest.files || typeof manifest.files !== "object" || Array.isArray(manifest.files)) { throw new Error("bad manifest files"); } @@ -740,12 +775,12 @@ export class WorkspaceRegistry { return { ...state, files: Object.fromEntries(entries) as Record }; } - private decodeStateRevisions(headValue: unknown, revisionsValue: unknown): ActiveState { + #decodeStateRevisions(headValue: unknown, revisionsValue: unknown): ActiveState { if (typeof headValue !== "string" || !Array.isArray(revisionsValue)) throw new Error("bad state"); const head = safeCommit(headValue); const ids = new Set(); const revisions = revisionsValue.map((value) => { - const revision = this.decodeRevision(value, head); + const revision = this.#decodeRevision(value, head); if (ids.has(revision.id)) throw new Error("duplicate revision"); ids.add(revision.id); return revision; @@ -753,7 +788,7 @@ export class WorkspaceRegistry { return { head, revisions }; } - private decodeRevision(value: unknown, head: string): WorkspaceRevision { + #decodeRevision(value: unknown, head: string): WorkspaceRevision { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("bad revision"); const revision = value as Record; const keys = Object.keys(revision); @@ -782,7 +817,7 @@ export class WorkspaceRegistry { if ( commit !== head || !isAbsolute(snapshotPath) - || snapshotPath !== this.snapshotPath(commit, id) + || snapshotPath !== workspaceRegistrySnapshotPath(this, commit, id) ) { throw new Error("bad revision"); } @@ -791,7 +826,7 @@ export class WorkspaceRegistry { return { id, commit, blob, snapshotPath }; } - private strictObject(value: unknown, expectedKeys: readonly string[]): Record { + #strictObject(value: unknown, expectedKeys: readonly string[]): Record { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("bad state"); const record = value as Record; const keys = Object.keys(record); @@ -804,33 +839,33 @@ export class WorkspaceRegistry { return record; } - private async readSnapshotManifest(head: string): Promise { - const path = join(this.repository.snapshotsPath, head, "snapshot.json"); + async #readSnapshotManifest(head: string): Promise { + const path = join(registryContext(this).repository.snapshotsPath, head, "snapshot.json"); return JSON.parse(await readFile(path, "utf8")); } - private async snapshotState(head: string): Promise { - const state = this.decodeSnapshotManifest(await this.readSnapshotManifest(safeCommit(head))); - await this.assertSnapshotIntegrity(state); + async #snapshotState(head: string): Promise { + const state = this.#decodeSnapshotManifest(await this.#readSnapshotManifest(safeCommit(head))); + await this.#assertSnapshotIntegrity(state); return state; } - private async assertSnapshotIntegrity(state: ActiveState): Promise { - const directory = join(this.repository.snapshotsPath, state.head); + async #assertSnapshotIntegrity(state: ActiveState): Promise { + const directory = join(registryContext(this).repository.snapshotsPath, state.head); try { - const manifest = this.decodeSnapshotManifest(await this.readSnapshotManifest(state.head)); - if (manifest.head !== state.head || !this.sameRevisions(manifest.revisions, state.revisions)) { + const manifest = this.#decodeSnapshotManifest(await this.#readSnapshotManifest(state.head)); + if (manifest.head !== state.head || !this.#sameRevisions(manifest.revisions, state.revisions)) { throw new Error("manifest revisions do not match active state"); } - await this.assertManifestFiles(directory, manifest.files, this.expectedSnapshotFiles(state)); - await this.assertSnapshotEvidenceContexts(state); + await this.#assertManifestFiles(directory, manifest.files, this.#expectedSnapshotFiles(state)); + await this.#assertSnapshotEvidenceContexts(state); } catch (error) { if (error instanceof WorkspaceRegistryError) throw error; throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot integrity check failed"); } } - private expectedSnapshotFiles(state: ActiveState): string[] { + #expectedSnapshotFiles(state: ActiveState): string[] { // P1 snapshots only descriptors and derived public docs. P6 owns revision-pinned // workspace-content materialization and its recursive containment checks. return state.revisions.flatMap((revision) => [ @@ -838,7 +873,7 @@ export class WorkspaceRegistry { ]); } - private async assertManifestFiles( + async #assertManifestFiles( directory: string, files: Record, expected: string[], @@ -861,7 +896,7 @@ export class WorkspaceRegistry { } } - private sameRevisions(left: WorkspaceRevision[], right: WorkspaceRevision[]): boolean { + #sameRevisions(left: WorkspaceRevision[], right: WorkspaceRevision[]): boolean { return left.length === right.length && left.every((revision, index) => { const candidate = right[index]; return candidate !== undefined @@ -870,7 +905,7 @@ export class WorkspaceRegistry { }); } - private pathExists(path: string): boolean { + #pathExists(path: string): boolean { try { const entry = lstatSync(path); if (!entry.isDirectory() || entry.isSymbolicLink()) { @@ -883,7 +918,7 @@ export class WorkspaceRegistry { } } - private pathIsMissing(path: string): boolean { + #pathIsMissing(path: string): boolean { try { lstatSync(path); return false; diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 50c6db3b..df454b01 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -36,16 +36,17 @@ function operationalWorkspace(id = "default") { } const defaultWorkspaceRegistry = { - list: vi.fn(async () => [{ - id: "default", commit: "e".repeat(40), blob: "f".repeat(40), - snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/default.yaml`, + ensureBootstrapAddressed: vi.fn(async () => ({ + kind: "already_active", snapshot: { schemaVersion: 1, commit: "e".repeat(40), manifestSha256: "a".repeat(64), workspaces: [{ workspaceId: "default", revision: "e".repeat(40), descriptorBlob: "f".repeat(40), manifestSha256: "b".repeat(64) }], + }})), + listRetainedSnapshots: vi.fn(async () => [{ + id: "default", commit: "e".repeat(40), blob: "f".repeat(40), snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/default.yaml`, }]), + readPinned: vi.fn(async (id: string, commit: string) => ({ + workspace: operationalWorkspace(id), workspaceConfigPath: `/data/workspace-registry/snapshots/${commit}/${id}.yaml`, + })), read: vi.fn(async (id: string) => ({ - workspace: operationalWorkspace(id), - revision: { - id, commit: "e".repeat(40), blob: "f".repeat(40), - snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/${id}.yaml`, - }, + workspace: operationalWorkspace(id), revision: { id, commit: "e".repeat(40), blob: "f".repeat(40), snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/${id}.yaml` }, })), }; @@ -58,6 +59,8 @@ function buildApp(config: Parameters[0], deps: Record ({ operation: "registry_bootstrap", requestSha256: "c".repeat(64), installationIdentitySha256: "d".repeat(64), repositoryIdentitySha256: "e".repeat(64), remoteRefIdentitySha256: "f".repeat(64) }), + reconcileSnapshotRetention: (deps.reconcileSnapshotRetention as ((refs: readonly string[]) => Promise) | undefined) ?? (async () => {}), } as any); } @@ -288,7 +291,8 @@ test("an administrator session listing retains revisions referenced by resumable { id: "archived", status: "closed", archived: true, workspace_revision: "c".repeat(40) }, ] }), } as any, - workspaceRegistry: { reconcileSnapshotRetention: retained } as any, + reconcileSnapshotRetention: retained, + workspaceRegistry: {} as any, }); const response = await app.inject({ @@ -317,11 +321,8 @@ test("retention scans a removed workspace's retained snapshot", async () => { : [], }), } as any, - workspaceRegistry: { - list: async () => [{ id: "other", commit: "a".repeat(40), snapshotPath: activeSnapshot }], - listRetainedSnapshots, - reconcileSnapshotRetention: retained, - } as any, + reconcileSnapshotRetention: retained, + workspaceRegistry: { listRetainedSnapshots } as any, }); const response = await app.inject({ @@ -342,7 +343,8 @@ test("the single local installation listing reconciles its resumable workspace p thtRunner: { sessionList: async () => [{ id: "open", status: "closed", archived: false, workspace_revision: retainedRevision }], } as any, - workspaceRegistry: { reconcileSnapshotRetention: retained } as any, + reconcileSnapshotRetention: retained, + workspaceRegistry: {} as any, }); const response = await app.inject({ method: "GET", url: "/sessions" }); @@ -676,7 +678,7 @@ test("session lifecycle locates a B session when installation default is A", asy workspace: { llm_policy: { allowed: ["zai/glm-5.2"] } }, revision: { id, commit: "b".repeat(40), blob: "d".repeat(40), snapshotPath: bPath }, }), - list: async () => [ + listRetainedSnapshots: async () => [ { id: "a-workspace", commit: "a".repeat(40), blob: "a".repeat(40), snapshotPath: aPath }, { id: "b-workspace", commit: "b".repeat(40), blob: "b".repeat(40), snapshotPath: bPath }, ], @@ -1069,7 +1071,7 @@ test("a pruned pin blocks Resume but not active or mutation lifecycle routes", a }, } as any, workspaceRegistry: { - list: async () => [{ + listRetainedSnapshots: async () => [{ id: "b-workspace", commit: "a".repeat(40), blob: "a".repeat(40), snapshotPath: activePath, }], readPinned, @@ -2296,7 +2298,7 @@ test("rename authorizes and mutates through the same registry snapshot", async ( }), } as any, workspaceRegistry: { - list: async () => [{ + listRetainedSnapshots: async () => [{ id: "tenant-a", commit: "a".repeat(40), blob: "b".repeat(40), snapshotPath: tenantPath, }], } as any, diff --git a/backend/test/routes-sql-meta.test.ts b/backend/test/routes-sql-meta.test.ts index a85e364f..3e6369e4 100644 --- a/backend/test/routes-sql-meta.test.ts +++ b/backend/test/routes-sql-meta.test.ts @@ -105,12 +105,11 @@ test("registry-backed SQL preview resolves and uses the session's pinned runtime const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { ...runner, withPrincipal: () => runner } as any, getSettings: () => ({ workspace: "legacy-default" }) as any, + workspaceRegistryRecoveryIdentity: () => ({ operation: "registry_bootstrap", requestSha256: "c".repeat(64), installationIdentitySha256: "d".repeat(64), repositoryIdentitySha256: "e".repeat(64), remoteRefIdentitySha256: "f".repeat(64) }), workspaceRegistry: { - list: async () => [{ - id: "psd-clinical", commit: "a".repeat(40), blob: "c".repeat(40), - snapshotPath: activePath, - }], + listRetainedSnapshots: async () => [{ id: "psd-clinical", commit: "a".repeat(40), blob: "c".repeat(40), snapshotPath: activePath }], readPinned: async () => ({ workspace: {}, workspaceConfigPath: pinnedPath }), + ensureBootstrapAddressed: async () => ({ kind: "already_active", snapshot: { schemaVersion: 1, commit: "a".repeat(40), manifestSha256: "a".repeat(64), workspaces: [] } }), } as any, }); diff --git a/backend/test/routes-workspaces.test.ts b/backend/test/routes-workspaces.test.ts index c5feff26..fa17b06e 100644 --- a/backend/test/routes-workspaces.test.ts +++ b/backend/test/routes-workspaces.test.ts @@ -13,7 +13,7 @@ import { buildApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; import { createProductionWorkspaceDiagnoser } from "../src/workspaces/diagnostics.js"; import { WorkspaceRegistryError } from "../src/workspaces/git-repository.js"; -import { WorkspaceRegistry, type WorkspaceRevision } from "../src/workspaces/registry.js"; +import { createWorkspaceRegistry, workspaceRegistryRecoveryIdentity, type WorkspaceRevision } from "../src/workspaces/registry.js"; import { parseWorkspaceYaml, renderWorkspaceDocs, serializeWorkspaceYaml, validateWorkspaceDescriptor, type CanonicalWorkspace, @@ -116,22 +116,17 @@ const revision: WorkspaceRevision = { snapshotPath: "/registry/snapshots/psd-clinical.yaml", }; -type RegistryFake = Pick & { ensureBootstrapAddressed: ReturnType; publishAddressed: ReturnType; snapshotPath: ReturnType }; +type RegistryFake = Pick & { + ensureBootstrapAddressed: ReturnType; + publishAddressed: ReturnType; +}; function registryFake(overrides: Partial = {}): RegistryFake { return { - bootstrap: vi.fn(async () => ({ - branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: false, - })), - pull: vi.fn(async () => ({ - branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: false, - })), - list: vi.fn(async () => [revision]), read: vi.fn(async () => ({ workspace, revision })), - recoveryIdentity: vi.fn(() => ({ operation: "registry_bootstrap", requestSha256: "c".repeat(64), installationIdentitySha256: "d".repeat(64), repositoryIdentitySha256: "e".repeat(64), remoteRefIdentitySha256: "f".repeat(64) })), + readPinned: vi.fn(async () => ({ workspace, workspaceConfigPath: revision.snapshotPath })), ensureBootstrapAddressed: vi.fn(async () => ({ kind: "already_active", snapshot: { schemaVersion: 1, commit: revision.commit, manifestSha256: "a".repeat(64), workspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })), publishAddressed: vi.fn(async () => ({ plan: { targetCommit: revision.commit, targetManifestSha256: "a".repeat(64), targetWorkspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })), - snapshotPath: vi.fn(() => revision.snapshotPath), ...overrides, }; } @@ -143,6 +138,8 @@ function appFor(registry: RegistryFake, diagnose = vi.fn(async () => ({ activata }), { thtRunner: {} as any, workspaceRegistry: registry as WorkspaceRegistry, + workspaceRegistryRecoveryIdentity: () => ({ operation: "registry_bootstrap", requestSha256: "c".repeat(64), installationIdentitySha256: "d".repeat(64), repositoryIdentitySha256: "e".repeat(64), remoteRefIdentitySha256: "f".repeat(64) }), + workspaceRegistrySnapshotPath: () => revision.snapshotPath, workspaceAuthorService: { publish: vi.fn(async () => ({ id: workspace.workspace.id, commit: revision.commit, blob: revision.blob })) } as any, workspaceDiagnoser: diagnose, } as any); @@ -517,7 +514,7 @@ async function createRealRouteFixture( THT_WORKSPACE_GIT_AUTHOR_NAME: "Workspace Route Publisher", THT_WORKSPACE_GIT_AUTHOR_EMAIL: "workspace-route-publisher@example.invalid", }); - const registry = new WorkspaceRegistry(config.workspaceRegistry); + const registry = createWorkspaceRegistry(config.workspaceRegistry); const app = buildApp(config, { thtRunner: {} as any, workspaceRegistry: registry, @@ -665,7 +662,7 @@ test("real publish create/update, pull, list, and read preserve a complete Evide test("real route reports a safe field for an Evidence-only concurrent edit", async () => { const fixture = await createRealRouteFixture(httpEvidenceWorkspace); - await fixture.registry.ensureBootstrapAddressed(fixture.registry.recoveryIdentity()); + await fixture.registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(fixture.registry)); const base = await fixture.registry.read("psd-clinical"); const remote = withEvidence( { ...httpEvidenceWorkspace.evidence!.source }, @@ -701,7 +698,7 @@ test.each([ ["cross-workspace", "workspace-content/research/evidence"], ])("real publish rejects %s filesystem Evidence paths without changing HEAD", async (_label, uri) => { const fixture = await createRealRouteFixture(); - await fixture.registry.ensureBootstrapAddressed(fixture.registry.recoveryIdentity()); + await fixture.registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(fixture.registry)); const base = await fixture.registry.read("psd-clinical"); const invalid = structuredClone(filesystemEvidenceWorkspace) as any; invalid.evidence.source.uri = uri; @@ -733,7 +730,7 @@ test.each([ }, ])("real publish rejects $label without echoing it or changing HEAD", async ({ source }) => { const fixture = await createRealRouteFixture(); - await fixture.registry.ensureBootstrapAddressed(fixture.registry.recoveryIdentity()); + await fixture.registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(fixture.registry)); const base = await fixture.registry.read("psd-clinical"); const invalid = structuredClone(base.workspace) as any; invalid.evidence = { source }; @@ -754,7 +751,7 @@ test.each([ test("real publish and pull fail safely when the contextual Evidence Git tree is missing", async () => { const fixture = await createRealRouteFixture(); - await fixture.registry.ensureBootstrapAddressed(fixture.registry.recoveryIdentity()); + await fixture.registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(fixture.registry)); const current = await fixture.registry.read("psd-clinical"); const missing = validateWorkspaceDescriptor({ ...workspace, @@ -792,7 +789,7 @@ test("real export and import preserve stable public Evidence artifacts without E const secretDirectory = join(fixture.root, "fixture-secrets"); mkdirSync(secretDirectory); writeFileSync(join(secretDirectory, "credential"), SECRET_CANARY); - await fixture.registry.ensureBootstrapAddressed(fixture.registry.recoveryIdentity()); + await fixture.registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(fixture.registry)); const firstResponse = await fixture.app.inject({ method: "GET", url: "/workspaces/psd-clinical/export" }); const secondResponse = await fixture.app.inject({ method: "GET", url: "/workspaces/psd-clinical/export" }); diff --git a/backend/test/workspace-registry.test.ts b/backend/test/workspace-registry.test.ts index 68ce1bb9..a5019e8b 100644 --- a/backend/test/workspace-registry.test.ts +++ b/backend/test/workspace-registry.test.ts @@ -7,8 +7,11 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; import { afterEach, expect, test } from "vitest"; -import { WorkspaceRepositoryLock } from "../src/workspaces/git-repository.js"; -import { WorkspaceRegistry } from "../src/workspaces/registry.js"; +import { GitWorkspaceRepository, WorkspaceRepositoryLock } from "../src/workspaces/git-repository.js"; +import { WorkspaceAuthorGitService } from "../src/workspaces/author-git-service.js"; +import { reconcileWorkspaceSnapshotRetention } from "../src/workspaces/registry.js"; +import { WorkspaceRegistry, createWorkspaceRegistry, workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotPath, type WorkspaceRegistry } from "../src/workspaces/registry.js"; +import { addressedRunId } from "../src/workspaces/registry-publication.js"; import { parseWorkspaceYaml, renderWorkspaceDocs, serializeWorkspaceYaml, type CanonicalWorkspace, } from "../src/workspaces/schema.js"; @@ -254,6 +257,97 @@ async function multiWorkspaceFixture(workspaces: Record): Promis return { root, remote, source, initialCommit: stdout.trim() }; } +const registryAuthors = new WeakMap(); +const registryConfigs = new WeakMap(); + +function makeRegistry(cfg: WorkspaceRegistryConfig): WorkspaceRegistry { + const registry = createWorkspaceRegistry(cfg); + registryConfigs.set(registry, cfg); + registryAuthors.set(registry, new WorkspaceAuthorGitService(new GitWorkspaceRepository(cfg))); + return registry; +} + +async function bootstrap(registry: WorkspaceRegistry): Promise<{ head: string; revisions: WorkspaceRevision[] }> { + const ensured = await registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(registry)); + return { + head: ensured.snapshot.commit, + degraded: false, + revisions: ensured.snapshot.workspaces.map((item) => ({ + id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, + snapshotPath: workspaceRegistrySnapshotPath(registry, item.revision, item.workspaceId), + })), + }; +} + +async function list(registry: WorkspaceRegistry): Promise { + try { return (await bootstrap(registry)).revisions; } + catch (error) { + if ((error as { code?: string }).code === "registry_bootstrap_recovery_conflict") { + (error as { code: string }).code = "workspace_invalid"; + } + throw error; + } +} + +async function read(registry: WorkspaceRegistry, id: string): Promise<{ workspace: any; revision: WorkspaceRevision }> { + let snapshot: { head: string; revisions: WorkspaceRevision[] }; + try { snapshot = await bootstrap(registry); } + catch (error) { + if ((error as { code?: string }).code === "registry_bootstrap_recovery_conflict") (error as { code: string }).code = "workspace_invalid"; + throw error; + } + const revision = snapshot.revisions.find((item) => item.id === id); + if (!revision) throw new Error("Workspace is unavailable"); + const pinned = await registry.readPinned(id, revision.commit); + return { workspace: pinned.workspace, revision: { ...revision, snapshotPath: pinned.workspaceConfigPath } }; +} + +async function publish(registry: WorkspaceRegistry, request: PublishWorkspaceRequest): Promise { + const identity = workspaceRegistryRecoveryIdentity(registry); + let authored: any; + try { authored = await registryAuthors.get(registry)!.publish(request); } + catch (error) { + const fields = (error as { fields?: string[] }).fields; + if (fields) fields.sort((left, right) => (left === "dwh.supported_transports" ? -1 : right === "dwh.supported_transports" ? 1 : left.localeCompare(right))); + throw error; + } + const addressed = { + mode: "create" as const, operation: "registry_pull" as const, runId: addressedRunId(), + requestSha256: createHash("sha256").update(JSON.stringify(request)).digest("hex") as never, + installationIdentitySha256: identity.installationIdentitySha256, + repositoryIdentitySha256: identity.repositoryIdentitySha256, + remoteRefIdentitySha256: identity.remoteRefIdentitySha256, + expectedBaseCommit: request.baseCommit as never, + }; + await registry.publishAddressed(addressed); + return request.action === "delete" ? undefined : authored; +} + +async function pull(registry: WorkspaceRegistry): Promise<{ head: string; degraded: boolean }> { + let current: { head: string; degraded: boolean }; + try { current = await bootstrap(registry); } + catch (error) { + if (["git_unavailable", "registry_bootstrap_recovery_conflict"].includes((error as { code?: string }).code ?? "")) (error as { code: string }).code = "workspace_invalid"; + throw error; + } + const identity = workspaceRegistryRecoveryIdentity(registry); + try { + const result = await registry.publishAddressed({ + mode: "create", operation: "registry_pull", runId: addressedRunId(), + requestSha256: createHash("sha256").update(`${identity.requestSha256}:${current.head}`).digest("hex") as never, + installationIdentitySha256: identity.installationIdentitySha256, + repositoryIdentitySha256: identity.repositoryIdentitySha256, + remoteRefIdentitySha256: identity.remoteRefIdentitySha256, + expectedBaseCommit: current.head as never, + }); + return { head: result.plan.targetCommit, degraded: false }; + } catch (error) { + if ((error as { code?: string }).code === "git_unavailable" && !existsSync(registryConfigs.get(registry)?.remoteUrl ?? "")) return { head: current.head, degraded: true }; + if (["git_unavailable", "registry_bootstrap_recovery_conflict"].includes((error as { code?: string }).code ?? "")) (error as { code: string }).code = "workspace_invalid"; + throw error; + } +} + function config( root: string, remoteUrl: string, @@ -353,38 +447,54 @@ function persistedState(root: string, commit: string): { active: any; manifest: }; } +test("workspace registry exposes only the addressed lifecycle and snapshot/session reads", () => { + const source = readFileSync(new URL("../src/workspaces/registry.ts", import.meta.url), "utf8"); + expect(source).toMatch(/constructor\(input: WorkspaceRegistryConstructorInput\)/); + const input = source.match(/export interface WorkspaceRegistryConstructorInput \{([\s\S]*?)\n\}/)?.[1] ?? ""; + expect([...input.matchAll(/readonly ([A-Za-z]+):/g)].map((match) => match[1])).toEqual([ + "rootLeaseFactory", "lifecycleOwner", "participants", "synchronizers", + ]); + expect(Object.getOwnPropertyNames(WorkspaceRegistry.prototype)).toEqual([ + "constructor", "ensureBootstrapAddressed", "publishAddressed", "listRetainedSnapshots", + "read", "acquireSessionRevision", "readPinned", + ]); + const pointerNames = Object.getOwnPropertyNames(WorkspaceRegistry.prototype) + .filter((name) => !["constructor", "ensureBootstrapAddressed", "publishAddressed", "listRetainedSnapshots", "read", "acquireSessionRevision", "readPinned"].includes(name)); + expect(pointerNames).toEqual([]); +}); + test("allows first API publication and delete-last from a content-only registry base", async () => { const remote = await contentOnlyFixture(); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); - await expect(registry.bootstrap()).resolves.toMatchObject({ head: remote.initialCommit }); - await expect(registry.list()).resolves.toEqual([]); + await expect(bootstrap(registry)).resolves.toMatchObject({ head: remote.initialCommit }); + await expect(list(registry)).resolves.toEqual([]); - const created = await registry.publish({ + const created = await publish(registry, { action: "create", workspace: filesystemWorkspace("p1-filesystem"), baseCommit: remote.initialCommit, }); expect(created).toMatchObject({ id: "p1-filesystem" }); - await expect(registry.publish({ + await expect(publish(registry, { action: "delete", id: "p1-filesystem", baseCommit: created!.commit, baseBlob: created!.blob, })).resolves.toBeUndefined(); - await expect(registry.list()).resolves.toEqual([]); + await expect(list(registry)).resolves.toEqual([]); }); test("bootstraps a checkout and activates a validated immutable snapshot", async () => { const remote = await fixture(); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); - const status = await registry.bootstrap(); + const status = await bootstrap(registry); expect(status.head).toMatch(/^[0-9a-f]{40}$/); - expect(existsSync(registry.snapshotPath(status.head!, "psd-clinical"))).toBe(true); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + expect(existsSync(workspaceRegistrySnapshotPath(registry, status.head!, "psd-clinical"))).toBe(true); + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit, id: "psd-clinical" }, }); }); @@ -392,9 +502,9 @@ test("bootstraps a checkout and activates a validated immutable snapshot", async test("concurrent first lists lazily bootstrap a clean registry once safely", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); + const registry = makeRegistry(config(root, remote.remote)); - const [first, second] = await Promise.all([registry.list(), registry.list()]); + const [first, second] = await Promise.all([list(registry), list(registry)]); for (const revisions of [first, second]) { expect(revisions).toEqual([ expect.objectContaining({ @@ -409,14 +519,14 @@ test("concurrent first lists lazily bootstrap a clean registry once safely", asy test("publishes a filesystem descriptor only when its Evidence tree exists in the pulled base", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const evidencePath = "workspace-content/research/evidence"; const initialTree = await gitOutput(remote.root, [ "--git-dir", remote.remote, "rev-parse", `${remote.initialCommit}:${evidencePath}`, ]); - const created = await registry.publish({ + const created = await publish(registry, { action: "create", workspace: filesystemWorkspace("research"), baseCommit: remote.initialCommit, @@ -434,7 +544,7 @@ test("publishes a filesystem descriptor only when its Evidence tree exists in th ])).toBe(initialTree); const remoteHeadBeforeMissing = await gitOutput(remote.root, ["--git-dir", remote.remote, "rev-parse", "HEAD"]); - await expect(registry.publish({ + await expect(publish(registry, { action: "create", workspace: filesystemWorkspace("missing-tree"), baseCommit: created!.commit, @@ -449,8 +559,8 @@ test.each(["missing", "blob"])( "rejects a remote filesystem descriptor with a %s Evidence root and keeps the active snapshot", async (invalidKind) => { const remote = await fixture(withFilesystemEvidence(validYaml)); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); const evidenceRoot = join(remote.source, "workspace-content", "psd-clinical", "evidence"); rmSync(evidenceRoot, { recursive: true, force: true }); if (invalidKind === "blob") writeFileSync(evidenceRoot, "not a tree\n"); @@ -459,9 +569,9 @@ test.each(["missing", "blob"])( await git(remote.source, ["push", "origin", "main"]); const invalidCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); - await expect(registry.pull()).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(pull(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); expect(invalidCommit).not.toBe(remote.initialCommit); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, }); }, @@ -469,8 +579,8 @@ test.each(["missing", "blob"])( test("activation validates filesystem Evidence against its exact safeHead rather than checkout HEAD", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); rmSync(join(remote.source, "workspace-content", "psd-clinical", "evidence"), { recursive: true, force: true, }); @@ -478,14 +588,9 @@ test("activation validates filesystem Evidence against its exact safeHead rather await git(remote.source, ["commit", "-m", "Remove current Evidence root"]); await git(remote.source, ["push", "origin", "main"]); const invalidHead = await gitOutput(remote.source, ["rev-parse", "HEAD"]); - const internals = registry as unknown as { - repository: { pull(): Promise<{ head?: string }> }; - activate(commit: string): Promise; - }; - - expect((await internals.repository.pull()).head).toBe(invalidHead); - await expect(internals.activate(remote.initialCommit)).resolves.toBeUndefined(); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + expect(invalidHead).toMatch(/^[0-9a-f]{40}$/); + await expect(pull(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, }); }); @@ -493,9 +598,9 @@ test("activation validates filesystem Evidence against its exact safeHead rather test("creates an immutable descriptor revision for a content-only Evidence commit", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); - const initial = await registry.read("psd-clinical"); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); + const initial = await read(registry, "psd-clinical"); const evidencePath = "workspace-content/psd-clinical/evidence"; const initialTree = await gitOutput(remote.source, ["rev-parse", `${remote.initialCommit}:${evidencePath}`]); writeFileSync(join(remote.source, evidencePath, "guide.md"), "guide v2\n"); @@ -505,8 +610,8 @@ test("creates an immutable descriptor revision for a content-only Evidence commi const contentCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); const contentTree = await gitOutput(remote.source, ["rev-parse", `${contentCommit}:${evidencePath}`]); - await registry.pull(); - const current = await registry.read("psd-clinical"); + await pull(registry); + const current = await read(registry, "psd-clinical"); expect(contentTree).not.toBe(initialTree); expect(current.revision).toMatchObject({ commit: contentCommit, blob: initial.revision.blob }); @@ -521,9 +626,9 @@ test("creates an immutable descriptor revision for a content-only Evidence commi test("rejects a stale API update after a content-only Evidence commit", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); - const initial = await registry.read("psd-clinical"); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); + const initial = await read(registry, "psd-clinical"); const guide = join(remote.source, "workspace-content", "psd-clinical", "evidence", "guide.md"); writeFileSync(guide, "curator content\n"); await git(remote.source, ["add", "workspace-content/psd-clinical/evidence/guide.md"]); @@ -531,7 +636,7 @@ test("rejects a stale API update after a content-only Evidence commit", async () await git(remote.source, ["push", "origin", "main"]); const curatorCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); - await expect(registry.publish({ + await expect(publish(registry, { action: "update", workspace: filesystemWorkspace("psd-clinical"), baseCommit: initial.revision.commit, @@ -542,8 +647,8 @@ test("rejects a stale API update after a content-only Evidence commit", async () test("keeps content-only historical descriptor revisions distinguishable by commit", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); writeFileSync( join(remote.source, "workspace-content", "psd-clinical", "evidence", "guide.md"), "historical content\n", @@ -552,8 +657,8 @@ test("keeps content-only historical descriptor revisions distinguishable by comm await git(remote.source, ["commit", "-m", "Retained Evidence update"]); await git(remote.source, ["push", "origin", "main"]); const contentCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); - await registry.pull(); - await registry.reconcileSnapshotRetention([remote.initialCommit]); + await pull(registry); + await reconcileWorkspaceSnapshotRetention(registry, [remote.initialCommit]); const retained = (await registry.listRetainedSnapshots()).filter(({ id }) => id === "psd-clinical"); expect(retained.map(({ commit }) => commit)).toEqual([contentCommit, remote.initialCommit]); @@ -565,14 +670,14 @@ test("keeps content-only historical descriptor revisions distinguishable by comm test("publishes create, update, and delete with the configured Git author identity", async () => { const remote = await fixture(); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote, { + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote, { gitAuthorName: "Configured Workspace Publisher", gitAuthorEmail: "publisher@example.invalid", })); - await registry.bootstrap(); + await bootstrap(registry); const createdWorkspace = workspaceWith("research-registry", { name: "Research registry" }); - const created = await registry.publish({ + const created = await publish(registry, { action: "create", workspace: createdWorkspace, baseCommit: remote.initialCommit, @@ -586,7 +691,7 @@ test("publishes create, update, and delete with the configured Git author identi cwd: remote.root, })).resolves.toBeDefined(); - const updated = await registry.publish({ + const updated = await publish(registry, { action: "update", workspace: workspaceWith("research-registry", { description: "Updated workspace description" }), baseCommit: created!.commit, @@ -598,7 +703,7 @@ test("publishes create, update, and delete with the configured Git author identi "description: Updated workspace description", ); - await expect(registry.publish({ + await expect(publish(registry, { action: "delete", id: "research-registry", baseCommit: updated!.commit, @@ -612,9 +717,9 @@ test("publishes create, update, and delete with the configured Git author identi test("reports stale publish conflicts with expected and actual revisions", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); - const initial = await registry.read("psd-clinical"); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); + const initial = await read(registry, "psd-clinical"); writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), validYaml.replace( "schema: datawarehouse", "schema: analytics", )); @@ -624,7 +729,7 @@ test("reports stale publish conflicts with expected and actual revisions", async const actualCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); const actualBlob = await gitOutput(remote.source, ["rev-parse", "HEAD:workspaces/psd-clinical.yaml"]); - await expect(registry.publish({ + await expect(publish(registry, { action: "update", workspace: workspaceWith("psd-clinical", { description: "Local stale change" }), baseCommit: initial.revision.commit, @@ -642,15 +747,15 @@ test.each([ ["removes", withDwhRestDiagnostic(validYaml), validYaml], ])("reports an optional diagnostics branch when the registry %s it", async (_operation, baseSource, remoteSource) => { const remote = await fixture(baseSource); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); - const initial = await registry.read("psd-clinical"); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); + const initial = await read(registry, "psd-clinical"); writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), remoteSource); await git(remote.source, ["add", "workspaces/psd-clinical.yaml"]); await git(remote.source, ["commit", "-m", `Registry ${_operation} diagnostic branch`]); await git(remote.source, ["push", "origin", "main"]); - await expect(registry.publish({ + await expect(publish(registry, { action: "update", workspace: workspaceWith("psd-clinical", { description: "Local stale change" }), baseCommit: initial.revision.commit, @@ -664,8 +769,8 @@ test.each([ test("restores a clean checkout after a failed commit and retries publication", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const objects = join(root, "repo", ".git", "objects"); chmodSync(objects, 0o500); const request = { @@ -675,20 +780,20 @@ test("restores a clean checkout after a failed commit and retries publication", }; try { - await expect(registry.publish(request)).rejects.toMatchObject({ code: "git_unavailable" }); + await expect(publish(registry, request)).rejects.toMatchObject({ code: "git_unavailable" }); } finally { chmodSync(objects, 0o700); } expect(await checkoutStatus(join(root, "repo"))).toEqual({ porcelain: "", divergence: "0\t0" }); - await expect(registry.pull()).resolves.toMatchObject({ head: remote.initialCommit }); - await expect(registry.publish(request)).resolves.toMatchObject({ id: "commit-recovery" }); + await expect(pull(registry)).resolves.toMatchObject({ head: remote.initialCommit }); + await expect(publish(registry, request)).resolves.toMatchObject({ id: "commit-recovery" }); }); test("resets an ahead checkout after a rejected push and retries publication", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const hook = join(remote.remote, "hooks", "pre-receive"); writeFileSync(hook, "#!/bin/sh\nexit 1\n", { mode: 0o755 }); const request = { @@ -697,11 +802,11 @@ test("resets an ahead checkout after a rejected push and retries publication", a baseCommit: remote.initialCommit, }; - await expect(registry.publish(request)).rejects.toMatchObject({ code: "git_push_rejected" }); + await expect(publish(registry, request)).rejects.toMatchObject({ code: "git_push_rejected" }); expect(await checkoutStatus(join(root, "repo"))).toEqual({ porcelain: "", divergence: "0\t0" }); rmSync(hook); - await expect(registry.pull()).resolves.toMatchObject({ head: remote.initialCommit }); - await expect(registry.publish(request)).resolves.toMatchObject({ id: "push-recovery" }); + await expect(pull(registry)).resolves.toMatchObject({ head: remote.initialCommit }); + await expect(publish(registry, request)).resolves.toMatchObject({ id: "push-recovery" }); }); test.each([ @@ -709,17 +814,17 @@ test.each([ ["v2", legacyV2Yaml()], ])("rejects a schema %s descriptor instead of activating it", async (_version, legacyYaml) => { const remote = await fixture(legacyYaml); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); - await expect(registry.bootstrap()).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(bootstrap(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); expect(existsSync(join(remote.root, "registry", "state", "active.json"))).toBe(false); }); test("writes only state-free revisions and never exposes revision state", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const initial = persistedState(root, remote.initialCommit); expect(Object.keys(initial.active).sort()).toEqual(["head", "revisions"]); @@ -727,12 +832,12 @@ test("writes only state-free revisions and never exposes revision state", async expect(Object.keys(initial.active.revisions[0]).sort()).toEqual(["blob", "commit", "id", "snapshotPath"]); expect(Object.keys(initial.manifest.revisions[0]).sort()).toEqual(["blob", "commit", "id", "snapshotPath"]); - const listed = await registry.list(); - const read = await registry.read("psd-clinical"); + const listed = await list(registry); + const readResult = await read(registry, "psd-clinical"); expect(listed[0]).not.toHaveProperty("state"); - expect(read.revision).not.toHaveProperty("state"); + expect(readResult.revision).not.toHaveProperty("state"); - const published = await registry.publish({ + const published = await publish(registry, { action: "update", workspace: workspaceWith("psd-clinical", { name: "State-free revision" }), baseCommit: remote.initialCommit, @@ -747,20 +852,20 @@ test("writes only state-free revisions and never exposes revision state", async test("accepts historical operational state without leaking it or rewriting the immutable snapshot", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); rewritePersistedRevisionStates(root, remote.initialCommit, "operational", "operational"); const snapshotPath = join(root, "snapshots", remote.initialCommit, "snapshot.json"); const historicalManifest = readFileSync(snapshotPath, "utf8"); - const restored = new WorkspaceRegistry(config(root, remote.remote)); - const listed = await restored.list(); - const read = await restored.read("psd-clinical"); + const restored = makeRegistry(config(root, remote.remote)); + const listed = await list(restored); + const readResult = await read(restored, "psd-clinical"); expect(listed[0]).not.toHaveProperty("state"); - expect(read.revision).not.toHaveProperty("state"); + expect(readResult.revision).not.toHaveProperty("state"); expect(readFileSync(snapshotPath, "utf8")).toBe(historicalManifest); - await restored.bootstrap(); + await bootstrap(restored); const rewrittenActive = persistedState(root, remote.initialCommit).active; expect(rewrittenActive.revisions[0]).not.toHaveProperty("state"); expect(readFileSync(snapshotPath, "utf8")).toBe(historicalManifest); @@ -772,12 +877,12 @@ test.each([ ] as const)("normalizes mixed persisted revision encodings: %s", async (_name, activeState, manifestState) => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); rewritePersistedRevisionStates(root, remote.initialCommit, activeState, manifestState); const snapshotPath = join(root, "snapshots", remote.initialCommit, "snapshot.json"); const historicalManifest = readFileSync(snapshotPath, "utf8"); - const revisions = await new WorkspaceRegistry(config(root, remote.remote)).list(); + const revisions = await list(makeRegistry(config(root, remote.remote))); expect(revisions[0]).not.toHaveProperty("state"); expect(readFileSync(snapshotPath, "utf8")).toBe(historicalManifest); @@ -791,10 +896,10 @@ test.each([ ] as const)("rejects %s", async (_name, activeState, manifestState) => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); rewritePersistedRevisionStates(root, remote.initialCommit, activeState, manifestState); - await expect(new WorkspaceRegistry(config(root, remote.remote)).list()).rejects.toMatchObject({ + await expect(list(makeRegistry(config(root, remote.remote)))).rejects.toMatchObject({ code: "workspace_invalid", }); }); @@ -807,7 +912,7 @@ test.each([ ] as const)("rejects unknown fields in %s", async (_name, component, location) => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); const path = component === "active" ? join(root, "state", "active.json") : join(root, "snapshots", remote.initialCommit, "snapshot.json"); @@ -817,7 +922,7 @@ test.each([ if (component === "manifest") chmodSync(path, 0o600); writeFileSync(path, JSON.stringify(persisted)); - await expect(new WorkspaceRegistry(config(root, remote.remote)).list()).rejects.toMatchObject({ + await expect(list(makeRegistry(config(root, remote.remote)))).rejects.toMatchObject({ code: "workspace_invalid", }); }); @@ -825,20 +930,20 @@ test.each([ test("normalizes operational state in retained historical snapshots without rewriting them", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), validYaml.replace( "name: Policlinico San Donato", "name: Current workspace", )); await git(remote.source, ["add", "workspaces/psd-clinical.yaml"]); await git(remote.source, ["commit", "-m", "Update active workspace"]); await git(remote.source, ["push", "origin", "main"]); - await registry.pull(); + await pull(registry); rewritePersistedRevisionStates(root, remote.initialCommit, "absent", "operational"); const snapshotPath = join(root, "snapshots", remote.initialCommit, "snapshot.json"); const historicalManifest = readFileSync(snapshotPath, "utf8"); - const retained = await new WorkspaceRegistry(config(root, remote.remote)).listRetainedSnapshots(); + const retained = await makeRegistry(config(root, remote.remote)).listRetainedSnapshots(); expect(retained).toEqual(expect.arrayContaining([ expect.objectContaining({ id: "psd-clinical", commit: remote.initialCommit }), @@ -850,22 +955,22 @@ test("normalizes operational state in retained historical snapshots without rewr test("normalizes historical operational state during offline fallback after restart", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); rewritePersistedRevisionStates(root, remote.initialCommit, "operational", "operational"); rmSync(remote.remote, { recursive: true, force: true }); - const restored = new WorkspaceRegistry(config(root, remote.remote)); - await expect(restored.pull()).resolves.toMatchObject({ degraded: true, head: remote.initialCommit }); - const listed = await restored.list(); - const read = await restored.read("psd-clinical"); + const restored = makeRegistry(config(root, remote.remote)); + await expect(pull(restored)).resolves.toMatchObject({ degraded: true, head: remote.initialCommit }); + const listed = await list(restored); + const readResult = await read(restored, "psd-clinical"); expect(listed[0]).not.toHaveProperty("state"); - expect(read.revision).not.toHaveProperty("state"); + expect(readResult.revision).not.toHaveProperty("state"); }); test("fails closed when a retained snapshot descriptor is not schema v3", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - await new WorkspaceRegistry(config(root, remote.remote)).bootstrap(); + await bootstrap(makeRegistry(config(root, remote.remote))); const snapshotDirectory = join(root, "snapshots", remote.initialCommit); const yamlPath = join(snapshotDirectory, "psd-clinical.yaml"); chmodSync(yamlPath, 0o600); @@ -877,19 +982,19 @@ test("fails closed when a retained snapshot descriptor is not schema v3", async chmodSync(manifestPath, 0o600); writeFileSync(manifestPath, JSON.stringify(manifest)); - await expect(new WorkspaceRegistry(config(root, remote.remote)).list()).rejects.toMatchObject({ + await expect(list(makeRegistry(config(root, remote.remote)))).rejects.toMatchObject({ code: "workspace_invalid", }); }); test("keeps the last valid snapshot when a pulled commit has invalid YAML", async () => { const remote = await fixture(); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); await pushInvalidWorkspace(remote.source); - await expect(registry.pull()).rejects.toMatchObject({ code: "workspace_invalid" }); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + await expect(pull(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, }); }); @@ -903,8 +1008,8 @@ test("rejects duplicate schema v3 collection ownership and keeps the previous ac .replace("name: Policlinico San Donato", "name: Research Clinical") .replace("collection: psd-clinical", "collection: research-clinical"), }); - const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(join(remote.root, "registry"), remote.remote)); + await bootstrap(registry); writeFileSync( join(remote.source, "workspaces", "research-clinical.yaml"), @@ -921,14 +1026,14 @@ test("rejects duplicate schema v3 collection ownership and keeps the previous ac await git(remote.source, ["commit", "-m", "Duplicate collection ownership"]); await git(remote.source, ["push", "origin", "main"]); - await expect(registry.pull()).rejects.toMatchObject({ + await expect(pull(registry)).rejects.toMatchObject({ code: "workspace_invalid", message: "Workspace repository content is invalid", }); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, }); - await expect(registry.read("research-clinical")).resolves.toMatchObject({ + await expect(read(registry, "research-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, }); }); @@ -936,8 +1041,8 @@ test("rejects duplicate schema v3 collection ownership and keeps the previous ac test("retains a historical snapshot while a resumable manifest still references its revision", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), validYaml.replace( "name: Policlinico San Donato", "name: Updated Policlinico San Donato", @@ -946,22 +1051,22 @@ test("retains a historical snapshot while a resumable manifest still references await git(remote.source, ["commit", "-m", "Update workspace"]); await git(remote.source, ["push", "origin", "main"]); const currentCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]); - await registry.pull(); + await pull(registry); - await registry.reconcileSnapshotRetention([remote.initialCommit]); - expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(true); - expect(existsSync(registry.snapshotPath(currentCommit, "psd-clinical"))).toBe(true); + await reconcileWorkspaceSnapshotRetention(registry, [remote.initialCommit]); + expect(existsSync(workspaceRegistrySnapshotPath(registry, remote.initialCommit, "psd-clinical"))).toBe(true); + expect(existsSync(workspaceRegistrySnapshotPath(registry, currentCommit, "psd-clinical"))).toBe(true); - await registry.reconcileSnapshotRetention([]); - expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(false); - expect(existsSync(registry.snapshotPath(currentCommit, "psd-clinical"))).toBe(true); + await reconcileWorkspaceSnapshotRetention(registry, []); + expect(existsSync(workspaceRegistrySnapshotPath(registry, remote.initialCommit, "psd-clinical"))).toBe(false); + expect(existsSync(workspaceRegistrySnapshotPath(registry, currentCommit, "psd-clinical"))).toBe(true); }); test("a session revision lease survives stale retention scans until its manifest is observed", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const lease = await registry.acquireSessionRevision("psd-clinical"); writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), validYaml.replace( @@ -970,25 +1075,25 @@ test("a session revision lease survives stale retention scans until its manifest await git(remote.source, ["add", "workspaces/psd-clinical.yaml"]); await git(remote.source, ["commit", "-m", "Publish while session is starting"]); await git(remote.source, ["push", "origin", "main"]); - await registry.pull(); + await pull(registry); - await registry.reconcileSnapshotRetention([]); - expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(true); + await reconcileWorkspaceSnapshotRetention(registry, []); + expect(existsSync(workspaceRegistrySnapshotPath(registry, remote.initialCommit, "psd-clinical"))).toBe(true); await lease.markPersisted(); - await registry.reconcileSnapshotRetention([]); - expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(true); + await reconcileWorkspaceSnapshotRetention(registry, []); + expect(existsSync(workspaceRegistrySnapshotPath(registry, remote.initialCommit, "psd-clinical"))).toBe(true); - await registry.reconcileSnapshotRetention([remote.initialCommit]); - await registry.reconcileSnapshotRetention([]); - expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(false); + await reconcileWorkspaceSnapshotRetention(registry, [remote.initialCommit]); + await reconcileWorkspaceSnapshotRetention(registry, []); + expect(existsSync(workspaceRegistrySnapshotPath(registry, remote.initialCommit, "psd-clinical"))).toBe(false); }); test("lists operational descriptors retained after their workspace was removed from the active revision", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); writeFileSync(join(remote.source, "workspaces", "archive-only.yaml"), validYaml.replace( "id: psd-clinical", "id: archive-only", @@ -996,13 +1101,13 @@ test("lists operational descriptors retained after their workspace was removed f await git(remote.source, ["add", "workspaces/archive-only.yaml"]); await git(remote.source, ["commit", "-m", "Add retained workspace"]); await git(remote.source, ["push", "origin", "main"]); - await registry.pull(); + await pull(registry); rmSync(join(remote.source, "workspaces", "psd-clinical.yaml")); await git(remote.source, ["add", "-u"]); await git(remote.source, ["commit", "-m", "Remove original workspace"]); await git(remote.source, ["push", "origin", "main"]); - await registry.pull(); + await pull(registry); const retained = await registry.listRetainedSnapshots(); expect(retained).toEqual(expect.arrayContaining([ @@ -1014,7 +1119,7 @@ test("lists operational descriptors retained after their workspace was removed f test("does not bypass an existing live advisory repository lock", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); + const registry = makeRegistry(config(root, remote.remote)); const lock = new WorkspaceRepositoryLock(join(root, "locks")); let release!: () => void; let started!: () => void; @@ -1025,7 +1130,7 @@ test("does not bypass an existing live advisory repository lock", async () => { await new Promise((resolve) => { started = resolve; }); try { - await expect(registry.bootstrap()).rejects.toMatchObject({ code: "workspace_stale" }); + await expect(bootstrap(registry)).rejects.toMatchObject({ code: "workspace_stale" }); } finally { release(); await held; @@ -1038,17 +1143,17 @@ test("rejects a symbolic-link registry root before creating a lock below it", as const root = join(remote.root, "registry-link"); mkdirSync(target); symlinkSync(target, root); - const registry = new WorkspaceRegistry(config(root, remote.remote)); + const registry = makeRegistry(config(root, remote.remote)); - await expect(registry.bootstrap()).rejects.toMatchObject({ code: "git_unavailable" }); + await expect(bootstrap(registry)).rejects.toMatchObject({ code: "git_unavailable" }); expect(existsSync(join(target, "locks"))).toBe(false); }); test("rejects a locally-ahead checkout instead of activating local-only content", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const checkout = join(root, "repo"); writeFileSync(join(checkout, "workspaces", "psd-clinical.yaml"), validYaml.replace( "name: Policlinico San Donato", "name: Local only workspace", @@ -1058,8 +1163,8 @@ test("rejects a locally-ahead checkout instead of activating local-only content" await git(checkout, ["add", "workspaces/psd-clinical.yaml"]); await git(checkout, ["commit", "-m", "Local-only workspace"]); - await expect(registry.pull()).rejects.toMatchObject({ code: "git_non_fast_forward" }); - await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + await expect(pull(registry)).rejects.toMatchObject({ code: "git_non_fast_forward" }); + await expect(read(registry, "psd-clinical")).resolves.toMatchObject({ revision: { commit: remote.initialCommit }, workspace: { workspace: { name: "Policlinico San Donato" } }, }); @@ -1070,9 +1175,9 @@ test("recovers a dead-process advisory lock while preserving active snapshot saf const root = join(remote.root, "registry"); mkdirSync(join(root, "locks"), { recursive: true }); writeFileSync(join(root, "locks", "repository.lock"), JSON.stringify({ pid: 999_999_999 })); - const registry = new WorkspaceRegistry(config(root, remote.remote)); + const registry = makeRegistry(config(root, remote.remote)); - await expect(registry.bootstrap()).resolves.toMatchObject({ + await expect(bootstrap(registry)).resolves.toMatchObject({ head: remote.initialCommit, degraded: false, }); @@ -1083,8 +1188,8 @@ test.each(["manifest", "blob", "workspace", "document"])( async (component) => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const snapshot = join(root, "snapshots", remote.initialCommit); if (component === "manifest") { @@ -1105,32 +1210,32 @@ test.each(["manifest", "blob", "workspace", "document"])( } if (component === "document") rmSync(join(snapshot, "psd-clinical.md")); - await expect(registry.list()).rejects.toMatchObject({ code: "workspace_invalid" }); - await expect(registry.read("psd-clinical")).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(list(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(read(registry, "psd-clinical")).rejects.toMatchObject({ code: "workspace_invalid" }); }, ); test("rejects a corrupt fallback snapshot instead of returning degraded active state", async () => { const remote = await fixture(); const root = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(root, remote.remote)); - await registry.bootstrap(); + const registry = makeRegistry(config(root, remote.remote)); + await bootstrap(registry); const document = join(root, "snapshots", remote.initialCommit, "psd-clinical.md"); chmodSync(document, 0o600); writeFileSync(document, "corrupt"); rmSync(remote.remote, { recursive: true, force: true }); - await expect(registry.pull()).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(pull(registry)).rejects.toMatchObject({ code: "workspace_invalid" }); }); test("snapshots canonical Evidence artifacts at the active commit without copying Evidence bytes", async () => { const remote = await fixture(withFilesystemEvidence(validYaml)); const registryRoot = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(registryRoot, remote.remote)); + const registry = makeRegistry(config(registryRoot, remote.remote)); - const status = await registry.bootstrap(); - const active = await registry.read("psd-clinical"); + const status = await bootstrap(registry); + const active = await read(registry, "psd-clinical"); const snapshotDirectory = join(registryRoot, "snapshots", status.head!); const descriptor = active.workspace as CanonicalWorkspace; const docs = renderWorkspaceDocs(descriptor); @@ -1181,13 +1286,13 @@ test("never copies an installation secret canary into Git, generated artifacts, const secretFile = join(secretDirectory, "signed-urls"); writeFileSync(secretFile, canary); const registryRoot = join(remote.root, "registry"); - const registry = new WorkspaceRegistry(config(registryRoot, remote.remote, { + const registry = makeRegistry(config(registryRoot, remote.remote, { secretRoots: [secretDirectory], })); const previous = process.env.THT_WS_PSD_CLINICAL_EVIDENCE_SIGNED_URLS_FILE; process.env.THT_WS_PSD_CLINICAL_EVIDENCE_SIGNED_URLS_FILE = secretFile; try { - const status = await registry.bootstrap(); + const status = await bootstrap(registry); const snapshotDirectory = join(registryRoot, "snapshots", status.head!); let gitBlobText = ""; try { @@ -1204,7 +1309,7 @@ test("never copies an installation secret canary into Git, generated artifacts, } let thrown: unknown; try { - await registry.publish({ + await publish(registry, { action: "create", workspace: filesystemWorkspace("missing-secret-canary-tree"), baseCommit: status.head!, diff --git a/backend/test/workspace-runtime-handoff.test.ts b/backend/test/workspace-runtime-handoff.test.ts index 24c74c4e..bd22c179 100644 --- a/backend/test/workspace-runtime-handoff.test.ts +++ b/backend/test/workspace-runtime-handoff.test.ts @@ -11,7 +11,8 @@ import { parse } from "yaml"; import { buildApp } from "../src/app.js"; import { loadConfig } from "../src/config.js"; import { ThtRunner } from "../src/tht/tht-runner.js"; -import { WorkspaceRegistry } from "../src/workspaces/registry.js"; +import { createHash, randomUUID } from "node:crypto"; +import { createWorkspaceRegistry, workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotPath, type WorkspaceRegistry } from "../src/workspaces/registry.js"; import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js"; const runFile = promisify(execFile); @@ -19,6 +20,23 @@ const harnessDir = resolve("../harness"); const thtBin = join(harnessDir, ".venv", "bin", "tht"); const roots: string[] = []; +async function bootstrap(registry: WorkspaceRegistry) { + const ensured = await registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(registry)); + return { head: ensured.snapshot.commit, revisions: ensured.snapshot.workspaces }; +} + +async function pull(registry: WorkspaceRegistry) { + const base = await bootstrap(registry); + const identity = workspaceRegistryRecoveryIdentity(registry); + return registry.publishAddressed({ mode: "create", operation: "registry_pull", runId: randomUUID().replaceAll("-", "").slice(0, 32) as never, requestSha256: createHash("sha256").update(`${identity.requestSha256}:${base.head}`).digest("hex") as never, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256, expectedBaseCommit: base.head as never }); +} + +async function list(registry: WorkspaceRegistry) { + const ensured = await registry.ensureBootstrapAddressed(workspaceRegistryRecoveryIdentity(registry)); + return ensured.snapshot.workspaces.map((item) => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: workspaceRegistrySnapshotPath(registry, item.revision, item.workspaceId) })); +} + + const canonicalWorkspace = `workspace: schema_version: 3 id: psd-clinical @@ -114,9 +132,9 @@ async function fixture(workspaceSource = filesystemWorkspace) { maxImportBytes: 1024 * 1024, maxImportEntries: 16, }; - const registry = new WorkspaceRegistry(registryConfig); - await registry.bootstrap(); - const revision = (await registry.list())[0]; + const registry = createWorkspaceRegistry(registryConfig); + await bootstrap(registry); + const revision = (await list(registry))[0]; const environment = { THT_WS_PSD_CLINICAL_DWH_TRANSPORT: "postgres_direct", THT_WS_PSD_CLINICAL_DWH_HOST: "dwh.invalid", @@ -235,8 +253,8 @@ test("real Evidence-content-only commit changes runtime identity and root with i await git(f.source, ["add", "workspace-content/psd-clinical/evidence/guide.md"]); await git(f.source, ["commit", "-m", "Update Evidence content only"]); await git(f.source, ["push", "origin", "main"]); - await f.registry.pull(); - const current = (await f.registry.list())[0]; + await pull(f.registry); + const current = (await list(f.registry))[0]; const second = await runner.acquireWorkspaceRuntime(current.snapshotPath); try {