refactor workspace registry addressed API

This commit is contained in:
2026-08-11 21:46:01 +02:00
parent 1cc0b50446
commit e2fecf9953
10 changed files with 582 additions and 432 deletions
+21 -30
View File
@@ -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<typeof workspaceRegistryRecoveryIdentity>;
workspaceRegistrySnapshotPath?: (commit: string, id: string) => string;
reconcileSnapshotRetention?: (referencedCommits: readonly string[]) => Promise<void>;
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 });
+6 -5
View File
@@ -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<typeof workspaceRegistryRecoveryIdentity>;
reconcileSnapshotRetention?: (referencedCommits: readonly string[]) => Promise<void>;
/** 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<WorkspaceRegistry>).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
+3 -2
View File
@@ -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<Settings>;
workspaceRegistry: WorkspaceRegistry;
workspaceRegistryRecoveryIdentity?: () => ReturnType<typeof workspaceRegistryRecoveryIdentity>;
}): 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 {
+8 -7
View File
@@ -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<typeof workspaceRegistryRecoveryIdentity>;
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); }
+223 -188
View File
@@ -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<unknown>[];
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<unknown>[];
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<unknown>[];
readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[];
readonly installationIdentity: RegistryBootstrapRecoveryIdentityV1;
readonly repositoryIdentity: WorkspaceRegistryFactoryDependencies["repositoryIdentity"];
readonly remoteIdentity: WorkspaceRegistryFactoryDependencies["remoteIdentity"];
}
const registryContexts = new WeakMap<WorkspaceRegistry, WorkspaceRegistryContext>();
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<void> {
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<unknown>[];
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<RegistryEnsureBootstrapAddressedResultV1> {
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<RegistryEnsureBootstrapAddressedResultV1> {
async #ensureBootstrapAddressedLocked(identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryEnsureBootstrapAddressedResultV1> {
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<RegistryAddressedResultV1, { operation: "registry_bootstrap" }>, 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<RegistryAddressedResultV1, { operation: "registry_bootstrap" }>, snapshot: this.#addressedSnapshot(published, identity.requestSha256) };
}
async publishAddressed(request: RegistryAddressedRequestV1, context?: { readonly authoring?: PublishWorkspaceRequest }): Promise<RegistryAddressedResultV1> {
await this.repository.ensureLayout();
return this.lock.run(async () => {
const store = new RegistryAddressedPublicationStore(this.repository.root);
async publishAddressed(request: RegistryAddressedRequestV1): Promise<RegistryAddressedResultV1> {
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<void> {
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<RegistryAddressedResultV1> {
async #executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryAddressedResultV1> {
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<RegistryAddressedPlanV1> {
async #planForState(state: RegistryAddressedPublicationStateV1): Promise<RegistryAddressedPlanV1> {
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<WorkspaceRevision[]> {
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<SessionRevisionLease> {
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<void> {
static reconcileSnapshotRetention(registry: WorkspaceRegistry, referencedCommits: readonly string[]): Promise<void> { return registry.#reconcileSnapshotRetention(referencedCommits); }
async #reconcileSnapshotRetention(referencedCommits: readonly string[]): Promise<void> {
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<string> {
const directory = this.revisionLeaseDirectory();
async #writeRevisionLease(record: RevisionLeaseRecord, exclusive: boolean): Promise<string> {
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<void> {
async #replaceRevisionLease(path: string, record: RevisionLeaseRecord): Promise<void> {
const staging = `${path}.staging-${randomUUID()}`;
try {
await writeFile(staging, JSON.stringify(record), {
@@ -481,8 +516,8 @@ export class WorkspaceRegistry {
}
}
private async revisionLeases(): Promise<Array<{ path: string; record: RevisionLeaseRecord }>> {
const directory = this.revisionLeaseDirectory();
async #revisionLeases(): Promise<Array<{ path: string; record: RevisionLeaseRecord }>> {
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<WorkspaceConflictError> {
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<CanonicalWorkspace | undefined> {
async #readSnapshotCanonical(commit: string, id: string): Promise<CanonicalWorkspace | undefined> {
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<string, unknown>;
const remoteObject = remote as Record<string, unknown>;
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<void> {
async #assertEvidenceContext(workspace: WorkspaceDescriptor, revision: string): Promise<void> {
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<void> {
async #assertSnapshotEvidenceContexts(state: ActiveState): Promise<void> {
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<void> {
async #materialize(commit: string): Promise<void> {
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<string, string> = {};
@@ -664,26 +699,26 @@ export class WorkspaceRegistry {
}
private async publishSnapshotPointer(commit: string): Promise<void> {
const state = await this.snapshotState(commit);
await this.writeActiveState({ head: state.head, revisions: state.revisions });
async #publishSnapshotPointer(commit: string): Promise<void> {
const state = await this.#snapshotState(commit);
await this.#writeActiveState({ head: state.head, revisions: state.revisions });
}
private async baseState(commit: Revision40 | null): Promise<ActiveState | undefined> {
return commit ? this.snapshotState(commit) : undefined;
async #baseState(commit: Revision40 | null): Promise<ActiveState | undefined> {
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<GitStatus> {
async #gitFallback(error: unknown): Promise<GitStatus> {
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<ActiveState> {
const active = await this.tryActiveState();
async #activeState(): Promise<ActiveState> {
const active = await this.#tryActiveState();
if (!active) throw new WorkspaceRegistryError("workspace_invalid", "No active workspace snapshot is available");
return active;
}
private async tryActiveState(): Promise<ActiveState | undefined> {
const file = join(this.repository.statePath, "active.json");
async #tryActiveState(): Promise<ActiveState | undefined> {
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<void> {
const target = join(this.repository.statePath, "active.json");
const staging = join(this.repository.statePath, `.active-${randomUUID()}.json`);
async #writeActiveState(state: ActiveState): Promise<void> {
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<string, string> };
}
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<string>();
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<string, unknown>;
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<string, unknown> {
#strictObject(value: unknown, expectedKeys: readonly string[]): Record<string, unknown> {
if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("bad state");
const record = value as Record<string, unknown>;
const keys = Object.keys(record);
@@ -804,33 +839,33 @@ export class WorkspaceRegistry {
return record;
}
private async readSnapshotManifest(head: string): Promise<unknown> {
const path = join(this.repository.snapshotsPath, head, "snapshot.json");
async #readSnapshotManifest(head: string): Promise<unknown> {
const path = join(registryContext(this).repository.snapshotsPath, head, "snapshot.json");
return JSON.parse(await readFile(path, "utf8"));
}
private async snapshotState(head: string): Promise<ActiveState> {
const state = this.decodeSnapshotManifest(await this.readSnapshotManifest(safeCommit(head)));
await this.assertSnapshotIntegrity(state);
async #snapshotState(head: string): Promise<ActiveState> {
const state = this.#decodeSnapshotManifest(await this.#readSnapshotManifest(safeCommit(head)));
await this.#assertSnapshotIntegrity(state);
return state;
}
private async assertSnapshotIntegrity(state: ActiveState): Promise<void> {
const directory = join(this.repository.snapshotsPath, state.head);
async #assertSnapshotIntegrity(state: ActiveState): Promise<void> {
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<string, string>,
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;
+20 -18
View File
@@ -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<typeof buildRealApp>[0], deps: Record<strin
...deps,
...(thtRunner ? { thtRunner } : {}),
workspaceRegistry: { ...defaultWorkspaceRegistry, ...(deps.workspaceRegistry as object | undefined) },
workspaceRegistryRecoveryIdentity: () => ({ 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<void>) | 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,
+3 -4
View File
@@ -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,
});
+14 -17
View File
@@ -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<WorkspaceRegistry, "read" | "recoveryIdentity"> & { ensureBootstrapAddressed: ReturnType<typeof vi.fn>; publishAddressed: ReturnType<typeof vi.fn>; snapshotPath: ReturnType<typeof vi.fn> };
type RegistryFake = Pick<WorkspaceRegistry, "read" | "readPinned"> & {
ensureBootstrapAddressed: ReturnType<typeof vi.fn>;
publishAddressed: ReturnType<typeof vi.fn>;
};
function registryFake(overrides: Partial<RegistryFake> = {}): 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" });
+260 -155
View File
@@ -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<string, string>): Promis
return { root, remote, source, initialCommit: stdout.trim() };
}
const registryAuthors = new WeakMap<WorkspaceRegistry, WorkspaceAuthorGitService>();
const registryConfigs = new WeakMap<WorkspaceRegistry, WorkspaceRegistryConfig>();
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<WorkspaceRevision[]> {
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<any> {
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<void>;
};
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<void>((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!,
+24 -6
View File
@@ -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 {