Files
ThothII/backend/src/workspaces/registry.ts
T

614 lines
24 KiB
TypeScript

import { createHash, randomUUID } from "node:crypto";
import { lstatSync } from "node:fs";
import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises";
import { isAbsolute, join } from "node:path";
import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js";
import {
GitWorkspaceRepository,
WorkspaceRegistryError,
WorkspaceRepositoryLock,
type GitStatus,
} from "./git-repository.js";
import {
isCanonicalWorkspace,
parseWorkspaceYaml,
serializeWorkspaceYaml,
type CanonicalWorkspace,
type WorkspaceDescriptor,
} from "./schema.js";
import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js";
export type { GitStatus } from "./git-repository.js";
export interface WorkspaceRevision {
id: string;
commit: string;
blob: string;
snapshotPath: string;
state: "operational" | "migration_required";
}
export type PublishWorkspaceRequest =
| { action: "create"; workspace: CanonicalWorkspace; baseCommit: string }
| { action: "update"; workspace: CanonicalWorkspace; baseCommit: string; baseBlob: string }
| { action: "delete"; id: string; baseCommit: string; baseBlob: string };
export class WorkspaceConflictError extends WorkspaceRegistryError {
constructor(
readonly fields: string[],
readonly expected: { commit: string; blob?: string },
readonly actual: { commit: string; blob?: string },
readonly base?: CanonicalWorkspace,
readonly local?: CanonicalWorkspace,
readonly remote?: CanonicalWorkspace,
) {
super("workspace_conflict", "Workspace revision conflicts with the active registry");
this.name = "WorkspaceConflictError";
}
}
interface ActiveState {
head: string;
revisions: WorkspaceRevision[];
}
interface SnapshotManifest extends ActiveState {
files: Record<string, string>;
}
type LegacyWorkspaceRevision = Omit<WorkspaceRevision, "state">;
interface LegacyActiveState {
head: string;
revisions: LegacyWorkspaceRevision[];
}
interface LegacySnapshotManifest extends LegacyActiveState {
files: Record<string, string>;
}
function workspacePath(id: string): string {
if (!/^[a-z][a-z0-9-]{2,62}$/.test(id)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace ID is invalid");
}
return `workspaces/${id}.yaml`;
}
function safeCommit(commit: string): string {
if (!/^[0-9a-f]{40}$/.test(commit)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
}
return commit;
}
function safeBlob(blob: string): string {
if (!/^[0-9a-f]{40}$/.test(blob)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot blob is invalid");
}
return blob;
}
function digest(contents: string | Buffer): string {
return createHash("sha256").update(contents).digest("hex");
}
function workspaceError(error: unknown): WorkspaceRegistryError {
if (error instanceof WorkspaceRegistryError) return error;
return new WorkspaceRegistryError("workspace_invalid", "Workspace repository content is invalid");
}
/** Immutable canonical workspace snapshots backed by the configured Git checkout. */
export class WorkspaceRegistry {
private readonly repository: GitWorkspaceRepository;
private readonly lock: WorkspaceRepositoryLock;
constructor(private readonly config: WorkspaceRegistryConfig) {
this.repository = new GitWorkspaceRepository(config);
this.lock = new WorkspaceRepositoryLock(this.repository.locksPath);
}
snapshotPath(commit: string, id: string): string {
return join(this.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`);
}
async bootstrap(): Promise<GitStatus> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
try {
const status = await this.repository.bootstrap();
await this.activate(status.head!);
return status;
} catch (error) {
return await this.gitFallback(error);
}
});
}
async pull(): Promise<GitStatus> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
try {
const status = await this.repository.pull();
await this.activate(status.head!);
return status;
} catch (error) {
return await this.gitFallback(error);
}
});
}
async list(): Promise<WorkspaceRevision[]> {
return (await this.activeState()).revisions;
}
async read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> {
const state = await this.activeState();
const revision = state.revisions.find((candidate) => candidate.id === id);
if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
try {
const source = await readFile(revision.snapshotPath, "utf8");
return { workspace: parseWorkspaceYaml(source), revision };
} catch (error) {
throw workspaceError(error);
}
}
/** 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);
try {
const source = await readFile(snapshotPath, "utf8");
return { workspace: parseWorkspaceYaml(source), workspaceConfigPath: snapshotPath };
} catch (error) {
throw workspaceError(error);
}
}
/**
* Publish canonical YAML and derived public documentation as one optimistic Git revision.
* The browser never provides paths or generated artifacts; those are derived server-side.
*/
async publish(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
const status = await this.repository.pull();
await this.activate(status.head!);
const current = await this.activeState();
const id = request.action === "delete" ? request.id : request.workspace.workspace.id;
const existing = current.revisions.find((revision) => revision.id === id);
const local = request.action === "delete" ? undefined : request.workspace;
if (request.baseCommit !== status.head || (
request.action !== "create" && existing?.blob !== request.baseBlob
)) {
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);
const yamlPath = workspacePath(id);
const docPaths = this.documentationPaths(id);
if (request.action === "delete") {
await this.repository.removeRegistryFile(yamlPath);
await this.repository.removeRegistryFile(docPaths.contract);
await this.repository.removeRegistryFile(docPaths.readme);
} else {
const canonical = request.workspace;
const source = serializeWorkspaceYaml(canonical);
const docs = renderWorkspaceDocs(canonical);
await this.repository.writeRegistryFile(yamlPath, source);
await this.repository.writeRegistryFile(docPaths.contract, docs.envExample);
await this.repository.writeRegistryFile(docPaths.readme, docs.markdown);
}
const next = await this.repository.commitAndPush(
[yamlPath, docPaths.contract, docPaths.readme],
request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`,
);
await this.activate(next.head!);
return (await this.activeState()).revisions.find((revision) => revision.id === id);
});
}
private documentationPaths(id: string): { contract: string; readme: string } {
workspacePath(id);
const directory = `workspace-docs/${id}`;
return { contract: `${directory}/contract.env.example`, readme: `${directory}/README.md` };
}
private 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);
let remote: CanonicalWorkspace | undefined;
if (existing) {
const read = await this.read(id);
remote = isCanonicalWorkspace(read.workspace) ? read.workspace : undefined;
}
return new WorkspaceConflictError(
this.changedFields(base, remote),
{ commit: request.baseCommit, ...(request.action === "create" ? {} : { blob: request.baseBlob }) },
{ commit: currentCommit, ...(existing ? { blob: existing.blob } : {}) },
base,
local,
remote,
);
}
private async readSnapshotCanonical(commit: string, id: string): Promise<CanonicalWorkspace | undefined> {
try {
const source = await readFile(this.snapshotPath(commit, id), "utf8");
const workspace = parseWorkspaceYaml(source);
return isCanonicalWorkspace(workspace) ? workspace : undefined;
} catch {
return undefined;
}
}
private changedFields(
base: unknown,
remote: unknown,
prefix = "",
): string[] {
if (base === undefined || remote === undefined) {
return base === remote ? [] : [prefix || "workspace.id"];
}
if (Array.isArray(base) || Array.isArray(remote) || typeof base !== "object" || typeof remote !== "object") {
return JSON.stringify(base) === JSON.stringify(remote) ? [] : [prefix];
}
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(
baseObject[key],
remoteObject[key],
prefix ? `${prefix}.${key}` : key,
));
}
private async activate(commit: string): Promise<void> {
const safeHead = safeCommit(commit);
const files = await this.repository.workspacePaths();
if (files.length === 0) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository contains no workspaces");
}
const snapshots: Array<{
id: string;
source: string;
workspace: WorkspaceDescriptor;
blob: string;
state: WorkspaceRevision["state"];
}> = [];
try {
for (const path of files) {
const id = path.slice("workspaces/".length, -".yaml".length);
const source = await this.repository.readWorkspace(path);
const workspace = parseWorkspaceYaml(source);
if (workspace.workspace.id !== id) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace ID does not match its repository path");
}
let snapshotSource = source;
const state: WorkspaceRevision["state"] = isCanonicalWorkspace(workspace)
? "operational"
: "migration_required";
if (isCanonicalWorkspace(workspace)) {
buildInstallationContract(workspace);
renderWorkspaceDocs(workspace);
snapshotSource = serializeWorkspaceYaml(workspace);
}
snapshots.push({
id,
source: snapshotSource,
workspace,
blob: await this.repository.blob(path),
state,
});
}
} catch (error) {
throw workspaceError(error);
}
const snapshotDirectory = join(this.repository.snapshotsPath, safeHead);
const revisions = snapshots.map((snapshot) => ({
id: snapshot.id,
commit: safeHead,
blob: snapshot.blob,
snapshotPath: this.snapshotPath(safeHead, snapshot.id),
state: snapshot.state,
}));
if (this.pathExists(snapshotDirectory)) {
await this.assertOrMigrateSnapshotIntegrity({ head: safeHead, revisions });
} else {
const staging = join(this.repository.snapshotsPath, `.staging-${randomUUID()}`);
await mkdir(staging, { mode: 0o700 });
try {
const files: Record<string, string> = {};
for (const snapshot of snapshots) {
const yamlName = `${snapshot.id}.yaml`;
const envName = `${snapshot.id}.env.example`;
const docsName = `${snapshot.id}.md`;
await writeFile(join(staging, yamlName), snapshot.source, { encoding: "utf8", mode: 0o400 });
files[yamlName] = digest(snapshot.source);
if (snapshot.state === "operational") {
const docs = renderWorkspaceDocs(snapshot.workspace);
await writeFile(join(staging, envName), docs.envExample, { encoding: "utf8", mode: 0o400 });
await writeFile(join(staging, docsName), docs.markdown, { encoding: "utf8", mode: 0o400 });
files[envName] = digest(docs.envExample);
files[docsName] = digest(docs.markdown);
}
}
await writeFile(join(staging, "snapshot.json"), JSON.stringify({ head: safeHead, revisions, files }), {
encoding: "utf8", mode: 0o400,
});
await rename(staging, snapshotDirectory);
} catch (error) {
await rm(staging, { recursive: true, force: true });
throw error;
}
}
await this.writeActiveState({ head: safeHead, revisions });
}
private 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();
if (!active) throw safeError;
return {
branch: this.config.branch,
head: active.head,
ahead: 0,
behind: 0,
degraded: true,
lastError: safeError.code,
};
}
private 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");
try {
const parsed: unknown = JSON.parse(await readFile(file, "utf8"));
if (this.isLegacyActiveState(parsed)) {
return await this.migrateLegacyActiveState(parsed);
}
const state = parsed as ActiveState;
this.assertActiveState(state);
await this.assertSnapshotIntegrity(state);
return state;
} catch (error) {
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`);
await writeFile(staging, JSON.stringify(state), { encoding: "utf8", mode: 0o600 });
await rename(staging, target);
}
private async writeSnapshotManifest(directory: string, manifest: SnapshotManifest): Promise<void> {
const target = join(directory, "snapshot.json");
const staging = join(directory, `.snapshot-${randomUUID()}.json`);
await writeFile(staging, JSON.stringify(manifest), { encoding: "utf8", mode: 0o400 });
await rename(staging, target);
}
private assertActiveState(state: ActiveState): void {
safeCommit(state.head);
if (!Array.isArray(state.revisions) || state.revisions.length === 0) throw new Error("bad state");
const ids = new Set<string>();
for (const revision of state.revisions) {
safeCommit(revision.commit);
safeBlob(revision.blob);
if (revision.commit !== state.head || ids.has(revision.id)) throw new Error("bad revision");
if (revision.state !== "operational" && revision.state !== "migration_required") throw new Error("bad revision");
ids.add(revision.id);
workspacePath(revision.id);
if (!isAbsolute(revision.snapshotPath) || revision.snapshotPath !== this.snapshotPath(revision.commit, revision.id)) {
throw new Error("bad snapshot path");
}
}
}
private isLegacyActiveState(value: unknown): value is LegacyActiveState {
if (!value || typeof value !== "object") return false;
const revisions = (value as { revisions?: unknown }).revisions;
return Array.isArray(revisions) && revisions.length > 0 && revisions.every((revision) => (
revision && typeof revision === "object" && !("state" in revision)
));
}
private isLegacySnapshotManifest(value: unknown): value is LegacySnapshotManifest {
return this.isLegacyActiveState(value)
&& !!(value as { files?: unknown }).files
&& typeof (value as { files?: unknown }).files === "object"
&& !Array.isArray((value as { files?: unknown }).files);
}
private assertLegacyActiveState(state: LegacyActiveState): void {
safeCommit(state.head);
if (!Array.isArray(state.revisions) || state.revisions.length === 0) throw new Error("bad legacy state");
const ids = new Set<string>();
for (const revision of state.revisions) {
safeCommit(revision.commit);
safeBlob(revision.blob);
if (revision.commit !== state.head || ids.has(revision.id)) throw new Error("bad legacy revision");
ids.add(revision.id);
workspacePath(revision.id);
if (!isAbsolute(revision.snapshotPath) || revision.snapshotPath !== this.snapshotPath(revision.commit, revision.id)) {
throw new Error("bad legacy snapshot path");
}
}
}
private async migrateLegacyActiveState(legacy: LegacyActiveState): Promise<ActiveState> {
this.assertLegacyActiveState(legacy);
const state = await this.deriveStateFromLegacyRevisions(legacy);
const manifest = await this.readSnapshotManifest(state.head);
if (this.isLegacySnapshotManifest(manifest)) {
await this.migrateLegacySnapshotManifest(state, manifest);
} else {
await this.assertSnapshotIntegrity(state);
}
await this.writeActiveState(state);
return state;
}
private async deriveStateFromLegacyRevisions(legacy: LegacyActiveState): Promise<ActiveState> {
const revisions: WorkspaceRevision[] = [];
for (const revision of legacy.revisions) {
const source = await readFile(revision.snapshotPath, "utf8");
const workspace = parseWorkspaceYaml(source);
if (workspace.workspace.id !== revision.id) throw new Error("legacy snapshot workspace is invalid");
revisions.push({
...revision,
state: isCanonicalWorkspace(workspace) ? "operational" : "migration_required",
});
}
return { head: legacy.head, revisions };
}
private async assertOrMigrateSnapshotIntegrity(state: ActiveState): Promise<void> {
const manifest = await this.readSnapshotManifest(state.head);
if (this.isLegacySnapshotManifest(manifest)) {
await this.migrateLegacySnapshotManifest(state, manifest);
return;
}
await this.assertSnapshotIntegrity(state);
}
private async migrateLegacySnapshotManifest(
state: ActiveState,
suppliedManifest?: LegacySnapshotManifest,
): Promise<void> {
const manifest = suppliedManifest ?? await this.readSnapshotManifest(state.head);
try {
if (!this.isLegacySnapshotManifest(manifest)) throw new Error("snapshot is not pre-state");
this.assertLegacyActiveState(manifest);
if (manifest.head !== state.head || !this.sameLegacyRevisions(manifest.revisions, state.revisions)) {
throw new Error("legacy manifest revisions do not match active state");
}
const derived = await this.deriveStateFromLegacyRevisions(manifest);
if (!this.sameRevisions(derived.revisions, state.revisions)) {
throw new Error("legacy manifest state does not match workspace snapshots");
}
const directory = join(this.repository.snapshotsPath, state.head);
const legacyExpected = state.revisions.flatMap((revision) => [
`${revision.id}.yaml`, `${revision.id}.env.example`, `${revision.id}.md`,
]);
await this.assertManifestFiles(directory, manifest.files, legacyExpected);
const expected = this.expectedSnapshotFiles(state);
const files = Object.fromEntries(expected.map((name) => [name, manifest.files[name]]));
await this.writeSnapshotManifest(directory, { ...state, files });
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot integrity check failed");
}
}
private async readSnapshotManifest(head: string): Promise<unknown> {
const path = join(this.repository.snapshotsPath, head, "snapshot.json");
return JSON.parse(await readFile(path, "utf8"));
}
private async assertSnapshotIntegrity(state: ActiveState): Promise<void> {
const directory = join(this.repository.snapshotsPath, state.head);
try {
const manifest = await this.readSnapshotManifest(state.head) as SnapshotManifest;
this.assertActiveState(manifest);
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));
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot integrity check failed");
}
}
private expectedSnapshotFiles(state: ActiveState): string[] {
return state.revisions.flatMap((revision) => revision.state === "operational"
? [`${revision.id}.yaml`, `${revision.id}.env.example`, `${revision.id}.md`]
: [`${revision.id}.yaml`]);
}
private async assertManifestFiles(
directory: string,
files: Record<string, string>,
expected: string[],
): Promise<void> {
if (!files || typeof files !== "object" || Object.keys(files).length !== expected.length || !expected.every((name) => (
/^[0-9a-f]{64}$/.test(files[name] ?? "")
))) throw new Error("manifest files are invalid");
for (const name of expected) {
const path = join(directory, name);
const entry = lstatSync(path);
if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("snapshot file is invalid");
const contents = await readFile(path);
if (digest(contents) !== files[name]) throw new Error("snapshot file does not match manifest");
if (name.endsWith(".yaml")) {
const workspace = parseWorkspaceYaml(contents.toString("utf8"));
if (workspace.workspace.id !== name.slice(0, -".yaml".length)) {
throw new Error("snapshot workspace is invalid");
}
}
}
}
private sameRevisions(left: WorkspaceRevision[], right: WorkspaceRevision[]): boolean {
return left.length === right.length && left.every((revision, index) => {
const candidate = right[index];
return candidate !== undefined
&& candidate.id === revision.id && candidate.commit === revision.commit
&& candidate.blob === revision.blob && candidate.snapshotPath === revision.snapshotPath
&& candidate.state === revision.state;
});
}
private sameLegacyRevisions(left: LegacyWorkspaceRevision[], right: WorkspaceRevision[]): boolean {
return left.length === right.length && left.every((revision, index) => {
const candidate = right[index];
return candidate !== undefined
&& candidate.id === revision.id && candidate.commit === revision.commit
&& candidate.blob === revision.blob && candidate.snapshotPath === revision.snapshotPath;
});
}
private pathExists(path: string): boolean {
try {
const entry = lstatSync(path);
if (!entry.isDirectory() || entry.isSymbolicLink()) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot path is invalid");
}
return true;
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
return false;
}
}
private pathIsMissing(path: string): boolean {
try {
lstatSync(path);
return false;
} catch {
return true;
}
}
}