915 lines
38 KiB
TypeScript
915 lines
38 KiB
TypeScript
import { createHash, randomUUID } from "node:crypto";
|
|
import { lstatSync, readdirSync, readFileSync } from "node:fs";
|
|
import { mkdir, readdir, readFile, rename, rm, writeFile } from "node:fs/promises";
|
|
import { isAbsolute, join } from "node:path";
|
|
import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js";
|
|
import { parseAnnotationsYaml } from "./annotations.js";
|
|
import { syncAnnotations } from "./annotations-sync.js";
|
|
import { materializeEvidenceTree } from "./evidence/materialization.js";
|
|
import { assertCatalogMatchesDescriptor, parseWorkspaceCatalogYaml, type WorkspaceCatalog, type WorkspaceCatalogEntry } from "./catalog.js";
|
|
import {
|
|
GitWorkspaceRepository,
|
|
WorkspaceRegistryError,
|
|
WorkspaceRepositoryLock,
|
|
normalizeRepositoryIdentity,
|
|
type GitStatus,
|
|
} from "./git-repository.js";
|
|
import {
|
|
parseWorkspaceYaml,
|
|
serializeWorkspaceYaml,
|
|
validateOperationalWorkspace,
|
|
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;
|
|
}
|
|
|
|
export interface StoredWorkspaceIntegrity {
|
|
state: "uninitialized" | "active";
|
|
workspaces: number;
|
|
fingerprint: string;
|
|
}
|
|
|
|
export interface SessionRevisionLease {
|
|
workspace: WorkspaceDescriptor;
|
|
revision: WorkspaceRevision;
|
|
/** Mark the manifest durable; retention removes the lease only after observing that manifest. */
|
|
markPersisted(): Promise<void>;
|
|
/** Remove a lease for a session that failed before its manifest was durable. */
|
|
abort(): Promise<void>;
|
|
}
|
|
|
|
interface ActiveState {
|
|
head: string;
|
|
revisions: WorkspaceRevision[];
|
|
catalog?: WorkspaceCatalog;
|
|
}
|
|
|
|
interface SnapshotManifest extends ActiveState {
|
|
files: Record<string, string>;
|
|
}
|
|
|
|
interface RevisionLeaseRecord {
|
|
version: 1;
|
|
token: string;
|
|
workspaceId: string;
|
|
commit: string;
|
|
state: "creating" | "persisted";
|
|
}
|
|
|
|
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 `${id}/workspace.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), `${id}.yaml`);
|
|
}
|
|
|
|
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 listCatalog(): Promise<Array<WorkspaceCatalogEntry & {
|
|
configurationState: "ready";
|
|
revision: WorkspaceRevision;
|
|
}>> {
|
|
const active = await this.tryActiveState();
|
|
if (!active) {
|
|
await this.bootstrap();
|
|
return await this.listCatalog();
|
|
}
|
|
const catalog = active.catalog ?? { schema_version: 1 as const, workspaces: [] };
|
|
return catalog.workspaces.map((entry) => ({
|
|
...entry,
|
|
configurationState: "ready" as const,
|
|
revision: active.revisions.find((revision) => revision.id === entry.id)!,
|
|
}));
|
|
}
|
|
|
|
async list(): Promise<WorkspaceRevision[]> {
|
|
const active = await this.tryActiveState();
|
|
if (active) return active.revisions;
|
|
// A clean installation has no active snapshot until the first registry operation. Keep
|
|
// this lazy so health/startup remain available when Git is temporarily unreachable, while
|
|
// still refusing corrupted existing state (tryActiveState throws instead of returning none).
|
|
await this.bootstrap();
|
|
return (await this.activeState()).revisions;
|
|
}
|
|
|
|
/**
|
|
* List every intact retained snapshot, current snapshots first. Session discovery and
|
|
* retention use this rather than only the active revision so removing a workspace from
|
|
* Git cannot strand a resumable session that still pins one of its older descriptors.
|
|
*/
|
|
async listRetainedSnapshots(): Promise<WorkspaceRevision[]> {
|
|
await this.repository.ensureLayout();
|
|
return await this.lock.run(async () => {
|
|
try {
|
|
const active = await this.activeState();
|
|
const revisions = [...active.revisions];
|
|
const entries = await readdir(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);
|
|
revisions.push(...state.revisions);
|
|
}
|
|
return revisions;
|
|
} catch (error) {
|
|
throw workspaceError(error);
|
|
}
|
|
});
|
|
}
|
|
|
|
/**
|
|
* Validate only persisted local registry state. Restore uses this path while the installation
|
|
* is stopped: it must neither contact Git nor turn a never-used registry into initialized state.
|
|
* The only never-initialized shape is an existing empty root. An initialized root is exactly
|
|
* repo/, snapshots/, state/, and locks/: locks contains repository.lock plus an empty
|
|
* empty-hooks/, state contains active.json plus an optional empty revision-leases/, and
|
|
* snapshots contains exact immutable commit snapshots plus an optional empty runtime/.
|
|
* Descendants may contain only ordinary directories and regular files; links and special files
|
|
* are rejected by the stable whole-tree fingerprint before any shape is accepted.
|
|
*/
|
|
async verifyStoredState(): Promise<StoredWorkspaceIntegrity> {
|
|
try {
|
|
const before = await this.storedStateFingerprint();
|
|
const rootEntries = await readdir(this.repository.root, { withFileTypes: true });
|
|
if (rootEntries.length === 0) {
|
|
const after = await this.storedStateFingerprint();
|
|
if (after !== before) throw new Error("workspace registry changed during inspection");
|
|
return { state: "uninitialized", workspaces: 0, fingerprint: `sha256:${before}` };
|
|
}
|
|
this.assertExactDirectoryEntries(rootEntries, {
|
|
locks: "directory", repo: "directory", snapshots: "directory", state: "directory",
|
|
});
|
|
|
|
this.assertExactDirectoryEntries(
|
|
await readdir(this.repository.locksPath, { withFileTypes: true }),
|
|
{ "empty-hooks": "directory", "repository.lock": "file" },
|
|
);
|
|
this.assertExactDirectoryEntries(
|
|
await readdir(join(this.repository.locksPath, "empty-hooks"), { withFileTypes: true }),
|
|
{},
|
|
);
|
|
|
|
const stateEntries = await readdir(this.repository.statePath, { withFileTypes: true });
|
|
const stateShape: Record<string, "file" | "directory"> = { "active.json": "file" };
|
|
if (stateEntries.some((entry) => entry.name === "revision-leases")) {
|
|
stateShape["revision-leases"] = "directory";
|
|
}
|
|
this.assertExactDirectoryEntries(stateEntries, stateShape);
|
|
if (stateShape["revision-leases"] !== undefined) {
|
|
this.assertExactDirectoryEntries(
|
|
await readdir(join(this.repository.statePath, "revision-leases"), { withFileTypes: true }),
|
|
{},
|
|
);
|
|
}
|
|
|
|
const active = this.decodeActiveState(JSON.parse(await readFile(
|
|
join(this.repository.statePath, "active.json"), "utf8",
|
|
)));
|
|
const snapshotEntries = await readdir(this.repository.snapshotsPath, { withFileTypes: true });
|
|
if (snapshotEntries.length === 0) throw new Error("workspace snapshots are unavailable");
|
|
let activeSnapshotFound = false;
|
|
for (const entry of snapshotEntries) {
|
|
if (entry.name === "runtime") {
|
|
if (!entry.isDirectory() || entry.isSymbolicLink()) throw new Error("workspace runtime path is invalid");
|
|
this.assertExactDirectoryEntries(
|
|
await readdir(join(this.repository.snapshotsPath, "runtime"), { withFileTypes: true }),
|
|
{},
|
|
);
|
|
continue;
|
|
}
|
|
if (!/^[0-9a-f]{40}$/.test(entry.name) || !entry.isDirectory() || entry.isSymbolicLink()) {
|
|
throw new Error("workspace snapshot path is invalid");
|
|
}
|
|
const state = entry.name === active.head ? active : await this.readStoredSnapshotState(entry.name);
|
|
await this.assertSnapshotIntegrity(state, false, true);
|
|
if (entry.name === active.head) activeSnapshotFound = true;
|
|
}
|
|
if (!activeSnapshotFound) throw new Error("active workspace snapshot is unavailable");
|
|
|
|
const after = await this.storedStateFingerprint();
|
|
if (after !== before) throw new Error("workspace registry changed during inspection");
|
|
return { state: "active", workspaces: active.revisions.length, fingerprint: `sha256:${before}` };
|
|
} catch (error) {
|
|
throw workspaceError(error);
|
|
}
|
|
}
|
|
|
|
private assertExactDirectoryEntries(
|
|
entries: Array<{ name: string; isFile(): boolean; isDirectory(): boolean; isSymbolicLink(): boolean }>,
|
|
expected: Record<string, "file" | "directory">,
|
|
): void {
|
|
if (entries.length !== Object.keys(expected).length) throw new Error("workspace registry shape is invalid");
|
|
for (const entry of entries) {
|
|
const kind = expected[entry.name];
|
|
if (kind === undefined || entry.isSymbolicLink()
|
|
|| (kind === "file" && !entry.isFile())
|
|
|| (kind === "directory" && !entry.isDirectory())) {
|
|
throw new Error("workspace registry shape is invalid");
|
|
}
|
|
}
|
|
}
|
|
|
|
private async storedStateFingerprint(): Promise<string> {
|
|
const records: string[] = [];
|
|
let entries = 0;
|
|
let totalBytes = 0;
|
|
const visit = async (path: string, relative: string): Promise<void> => {
|
|
const before = lstatSync(path);
|
|
if (before.isSymbolicLink()) throw new Error("workspace registry link is invalid");
|
|
const metadata = [
|
|
before.dev, before.ino, before.mode, before.uid, before.gid,
|
|
before.size, before.mtimeMs, before.ctimeMs,
|
|
].join(":");
|
|
entries += 1;
|
|
if (entries > 65_536) throw new Error("workspace registry contains too many entries");
|
|
if (before.isDirectory()) {
|
|
records.push(`directory:${relative}:${metadata}`);
|
|
const children = await readdir(path);
|
|
children.sort();
|
|
for (const name of children) {
|
|
await visit(join(path, name), relative === "." ? name : `${relative}/${name}`);
|
|
}
|
|
} else if (before.isFile()) {
|
|
if (before.size > 256 * 1024 * 1024) throw new Error("workspace registry file is too large");
|
|
totalBytes += before.size;
|
|
if (totalBytes > 2 * 1024 * 1024 * 1024) throw new Error("workspace registry is too large");
|
|
const contents = await readFile(path);
|
|
records.push(`file:${relative}:${metadata}:${contents.length}:${digest(contents)}`);
|
|
} else {
|
|
throw new Error("workspace registry entry is invalid");
|
|
}
|
|
const after = lstatSync(path);
|
|
if (before.dev !== after.dev || before.ino !== after.ino || before.mode !== after.mode
|
|
|| before.uid !== after.uid || before.gid !== after.gid || before.size !== after.size
|
|
|| before.mtimeMs !== after.mtimeMs || before.ctimeMs !== after.ctimeMs) {
|
|
throw new Error("workspace registry changed during inspection");
|
|
}
|
|
};
|
|
await visit(this.repository.root, ".");
|
|
return digest(records.join("\n"));
|
|
}
|
|
|
|
private async readStoredSnapshotState(head: string): Promise<ActiveState> {
|
|
return this.decodeSnapshotManifest(await this.readSnapshotManifest(safeCommit(head)));
|
|
}
|
|
|
|
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);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Resolve the active revision and create its cross-process retention lease under the same
|
|
* repository lock. The lease bridges the interval before `session_manifest.yaml` is durable.
|
|
*/
|
|
async acquireSessionRevision(id: string): Promise<SessionRevisionLease> {
|
|
await this.repository.ensureLayout();
|
|
return await 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;
|
|
try {
|
|
workspace = validateOperationalWorkspace(
|
|
parseWorkspaceYaml(await readFile(revision.snapshotPath, "utf8")),
|
|
);
|
|
} catch (error) {
|
|
throw workspaceError(error);
|
|
}
|
|
|
|
const token = randomUUID();
|
|
const record: RevisionLeaseRecord = {
|
|
version: 1,
|
|
token,
|
|
workspaceId: id,
|
|
commit: revision.commit,
|
|
state: "creating",
|
|
};
|
|
const path = await this.writeRevisionLease(record, true);
|
|
let localState: RevisionLeaseRecord["state"] | "aborted" = "creating";
|
|
|
|
return {
|
|
workspace,
|
|
revision,
|
|
markPersisted: async () => {
|
|
if (localState === "persisted") return;
|
|
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" });
|
|
});
|
|
localState = "persisted";
|
|
},
|
|
abort: async () => {
|
|
if (localState !== "creating") return;
|
|
await this.lock.run(async () => { await rm(path, { force: true }); });
|
|
localState = "aborted";
|
|
},
|
|
};
|
|
});
|
|
}
|
|
|
|
/** 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: validateOperationalWorkspace(parseWorkspaceYaml(source)),
|
|
workspaceConfigPath: snapshotPath,
|
|
};
|
|
} catch (error) {
|
|
throw workspaceError(error);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Garbage-collect obsolete immutable snapshots without breaking cold Resume.
|
|
* 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> {
|
|
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();
|
|
for (const { record } of leases) retained.add(record.commit);
|
|
retained.add((await this.activeState()).head);
|
|
const entries = await readdir(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 current = lstatSync(path);
|
|
if (!current.isDirectory() || current.isSymbolicLink()) continue;
|
|
await rm(path, { recursive: true, force: true });
|
|
}
|
|
// A persisted lease is handed off only when this exact authoritative scan has observed a
|
|
// manifest pin for its commit. A stale scan therefore keeps the lease and cannot prune it.
|
|
for (const { path, record } of leases) {
|
|
if (record.state === "persisted" && manifestReferences.has(record.commit)) {
|
|
await rm(path, { force: true });
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
private revisionLeaseDirectory(): string {
|
|
return join(this.repository.statePath, "revision-leases");
|
|
}
|
|
|
|
private 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), {
|
|
encoding: "utf8",
|
|
mode: 0o600,
|
|
flush: true,
|
|
...(exclusive ? { flag: "wx" } : {}),
|
|
});
|
|
return path;
|
|
}
|
|
|
|
private async replaceRevisionLease(path: string, record: RevisionLeaseRecord): Promise<void> {
|
|
const staging = `${path}.staging-${randomUUID()}`;
|
|
try {
|
|
await writeFile(staging, JSON.stringify(record), {
|
|
encoding: "utf8", mode: 0o600, flag: "wx", flush: true,
|
|
});
|
|
await rename(staging, path);
|
|
} catch (error) {
|
|
await rm(staging, { force: true });
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
private 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 }> = [];
|
|
for (const entry of entries) {
|
|
if (!entry.isFile() || entry.isSymbolicLink() || !/^[0-9a-f-]{36}\.json$/.test(entry.name)) {
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
|
|
}
|
|
const path = join(directory, entry.name);
|
|
let record: RevisionLeaseRecord;
|
|
try {
|
|
record = JSON.parse(await readFile(path, "utf8")) as RevisionLeaseRecord;
|
|
} catch {
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
|
|
}
|
|
if (
|
|
record.version !== 1
|
|
|| `${record.token}.json` !== entry.name
|
|
|| !/^[0-9a-f-]{36}$/.test(record.token)
|
|
|| !/^[a-z][a-z0-9-]{2,62}$/.test(record.workspaceId)
|
|
|| !/^[0-9a-f]{40}$/.test(record.commit)
|
|
|| (record.state !== "creating" && record.state !== "persisted")
|
|
) {
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision lease is invalid");
|
|
}
|
|
leases.push({ path, record });
|
|
}
|
|
return leases;
|
|
}
|
|
|
|
private 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);
|
|
}
|
|
|
|
private 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);
|
|
}
|
|
}
|
|
|
|
private async activate(commit: string): Promise<void> {
|
|
const safeHead = safeCommit(commit);
|
|
const catalog = parseWorkspaceCatalogYaml(await this.repository.readCatalog(safeHead));
|
|
const catalogById = new Map(catalog.workspaces.map((entry) => [entry.id, entry]));
|
|
for (const id of await this.repository.workspaceDirectories()) {
|
|
if (!catalogById.has(id)) {
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace directory is not listed in the catalog");
|
|
}
|
|
}
|
|
const files = await this.repository.workspacePaths();
|
|
const descriptorIds = new Set(files.map((path) => path.slice(0, -"/workspace.yaml".length)));
|
|
for (const id of catalogById.keys()) {
|
|
if (!descriptorIds.has(id)) {
|
|
throw new WorkspaceRegistryError(
|
|
"workspace_invalid",
|
|
"Every catalog workspace must have a published descriptor",
|
|
);
|
|
}
|
|
}
|
|
|
|
const snapshots: Array<{
|
|
id: string;
|
|
source: string;
|
|
workspace: WorkspaceDescriptor;
|
|
blob: string;
|
|
}> = [];
|
|
const collectionOwners = new Map<string, string>();
|
|
try {
|
|
for (const path of files) {
|
|
const id = path.slice(0, -"/workspace.yaml".length);
|
|
const entry = catalogById.get(id);
|
|
if (!entry) throw new WorkspaceRegistryError("workspace_invalid", "Workspace descriptor is not listed in the catalog");
|
|
const source = await this.repository.readWorkspace(path, safeHead);
|
|
const workspace = parseWorkspaceYaml(source);
|
|
if (workspace.workspace.id !== id) {
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace ID does not match its repository path");
|
|
}
|
|
assertCatalogMatchesDescriptor(entry, workspace);
|
|
await this.assertEvidenceContext(workspace, safeHead);
|
|
const annotations = await this.repository.annotationsObject(safeHead, id);
|
|
if (annotations !== undefined) {
|
|
parseAnnotationsYaml(annotations.contents.toString("utf8"));
|
|
if (this.config.dataRoot !== undefined) {
|
|
syncAnnotations({
|
|
dataRoot: this.config.dataRoot,
|
|
workspaceId: id,
|
|
commit: safeHead,
|
|
blobId: annotations.blobId,
|
|
contents: annotations.contents,
|
|
});
|
|
}
|
|
}
|
|
const collection = workspace.semantic_index.vector_store.collection;
|
|
const owner = collectionOwners.get(collection);
|
|
if (owner !== undefined) {
|
|
throw new Error(`duplicate qdrant collection ownership: ${collection} (${owner}, ${id})`);
|
|
}
|
|
collectionOwners.set(collection, id);
|
|
buildInstallationContract(workspace);
|
|
renderWorkspaceDocs(workspace);
|
|
snapshots.push({
|
|
id,
|
|
source: serializeWorkspaceYaml(workspace),
|
|
workspace,
|
|
blob: await this.repository.blob(path, safeHead),
|
|
});
|
|
}
|
|
} 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),
|
|
}));
|
|
if (this.pathExists(snapshotDirectory)) {
|
|
await this.assertSnapshotIntegrity({ head: safeHead, revisions, catalog });
|
|
} 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);
|
|
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);
|
|
if (snapshot.workspace.evidence?.source.type === "filesystem") {
|
|
const materialized = await materializeEvidenceTree({
|
|
repository: this.repository,
|
|
revision: safeHead,
|
|
id: snapshot.id,
|
|
targetDirectory: join(staging, snapshot.id),
|
|
limits: this.evidenceMaterializationLimits(),
|
|
});
|
|
files[`${snapshot.id}/evidence.manifest.json`] = materialized.manifestDigest;
|
|
}
|
|
}
|
|
await writeFile(join(staging, "snapshot.json"), JSON.stringify({ head: safeHead, revisions, catalog, 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, catalog });
|
|
}
|
|
|
|
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,
|
|
...(this.config.remoteUrl
|
|
? { repository: normalizeRepositoryIdentity(this.config.remoteUrl) }
|
|
: {}),
|
|
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 state = this.decodeActiveState(JSON.parse(await readFile(file, "utf8")));
|
|
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 decodeActiveState(value: unknown): ActiveState {
|
|
const record = this.optionalKeyObject(value, ["head", "revisions"], ["catalog"]);
|
|
const state = this.decodeStateRevisions(record.head, record.revisions);
|
|
return record.catalog === undefined ? state : { ...state, catalog: this.decodeCatalog(record.catalog) };
|
|
}
|
|
|
|
private decodeSnapshotManifest(value: unknown): SnapshotManifest {
|
|
const manifest = this.optionalKeyObject(value, ["head", "revisions", "files"], ["catalog"]);
|
|
const state = this.decodeStateRevisions(manifest.head, manifest.revisions);
|
|
const catalog = manifest.catalog === undefined ? undefined : this.decodeCatalog(manifest.catalog);
|
|
if (!manifest.files || typeof manifest.files !== "object" || Array.isArray(manifest.files)) {
|
|
throw new Error("bad manifest files");
|
|
}
|
|
const entries = Object.entries(manifest.files as Record<string, unknown>);
|
|
if (entries.some(([, contentsDigest]) => typeof contentsDigest !== "string")) {
|
|
throw new Error("bad manifest files");
|
|
}
|
|
return { ...state, ...(catalog ? { catalog } : {}), files: Object.fromEntries(entries) as Record<string, string> };
|
|
}
|
|
|
|
private 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);
|
|
if (ids.has(revision.id)) throw new Error("duplicate revision");
|
|
ids.add(revision.id);
|
|
return revision;
|
|
});
|
|
return { head, revisions };
|
|
}
|
|
|
|
private 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);
|
|
const required = ["id", "commit", "blob", "snapshotPath"];
|
|
const hasHistoricalState = Object.prototype.hasOwnProperty.call(revision, "state");
|
|
if (
|
|
keys.length !== required.length + (hasHistoricalState ? 1 : 0)
|
|
|| !required.every((key) => Object.prototype.hasOwnProperty.call(revision, key))
|
|
|| (hasHistoricalState && revision.state !== "operational")
|
|
) {
|
|
throw new Error("bad revision");
|
|
}
|
|
if (
|
|
typeof revision.id !== "string"
|
|
|| typeof revision.commit !== "string"
|
|
|| typeof revision.blob !== "string"
|
|
|| typeof revision.snapshotPath !== "string"
|
|
) {
|
|
throw new Error("bad revision");
|
|
}
|
|
const id = revision.id;
|
|
const commit = safeCommit(revision.commit);
|
|
const blob = safeBlob(revision.blob);
|
|
const snapshotPath = revision.snapshotPath;
|
|
workspacePath(id);
|
|
if (
|
|
commit !== head
|
|
|| !isAbsolute(snapshotPath)
|
|
|| snapshotPath !== this.snapshotPath(commit, id)
|
|
) {
|
|
throw new Error("bad revision");
|
|
}
|
|
// Always reconstruct a fresh public revision. The sole accepted historical state field is
|
|
// compatibility input and must never cross the registry boundary.
|
|
return { id, commit, blob, snapshotPath };
|
|
}
|
|
|
|
private decodeCatalog(value: unknown): WorkspaceCatalog {
|
|
if (typeof value !== "object" || value === null) throw new Error("bad catalog");
|
|
return parseWorkspaceCatalogYaml(JSON.stringify(value));
|
|
}
|
|
|
|
private optionalKeyObject(
|
|
value: unknown,
|
|
requiredKeys: readonly string[],
|
|
optionalKeys: 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 allowed = new Set([...requiredKeys, ...optionalKeys]);
|
|
const keys = Object.keys(record);
|
|
if (
|
|
keys.length !== requiredKeys.length + optionalKeys.length
|
|
|| !requiredKeys.every((key) => Object.prototype.hasOwnProperty.call(record, key))
|
|
|| !keys.every((key) => allowed.has(key))
|
|
) {
|
|
throw new Error("bad state");
|
|
}
|
|
return record;
|
|
}
|
|
|
|
private 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);
|
|
if (
|
|
keys.length !== expectedKeys.length
|
|
|| !expectedKeys.every((key) => Object.prototype.hasOwnProperty.call(record, key))
|
|
) {
|
|
throw new Error("bad state");
|
|
}
|
|
return record;
|
|
}
|
|
|
|
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 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,
|
|
verifyGitEvidence = true,
|
|
exactStoredShape = false,
|
|
): Promise<void> {
|
|
const directory = join(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)
|
|
|| JSON.stringify(manifest.catalog ?? null) !== JSON.stringify(state.catalog ?? null)) {
|
|
throw new Error("manifest state does not match active state");
|
|
}
|
|
await this.assertManifestFiles(directory, manifest.files, this.expectedSnapshotFiles(state, directory));
|
|
if (exactStoredShape) this.assertStoredSnapshotShape(directory, state);
|
|
if (verifyGitEvidence) await this.assertSnapshotEvidenceContexts(state);
|
|
} catch (error) {
|
|
if (error instanceof WorkspaceRegistryError) throw error;
|
|
throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot integrity check failed");
|
|
}
|
|
}
|
|
|
|
private assertStoredSnapshotShape(directory: string, state: ActiveState): void {
|
|
const expected: Record<string, "file" | "directory"> = { "snapshot.json": "file" };
|
|
for (const revision of state.revisions) {
|
|
expected[`${revision.id}.yaml`] = "file";
|
|
expected[`${revision.id}.env.example`] = "file";
|
|
expected[`${revision.id}.md`] = "file";
|
|
const workspace = parseWorkspaceYaml(readFileSync(join(directory, `${revision.id}.yaml`), "utf8"));
|
|
if (workspace.evidence?.source.type === "filesystem") expected[revision.id] = "directory";
|
|
}
|
|
this.assertExactDirectoryEntries(readdirSync(directory, { withFileTypes: true }), expected);
|
|
for (const revision of state.revisions) {
|
|
if (expected[revision.id] !== "directory") continue;
|
|
this.assertExactDirectoryEntries(
|
|
readdirSync(join(directory, revision.id), { withFileTypes: true }),
|
|
{ evidence: "directory", "evidence.manifest.json": "file" },
|
|
);
|
|
}
|
|
}
|
|
|
|
private expectedSnapshotFiles(state: ActiveState, directory: string): string[] {
|
|
return state.revisions.flatMap((revision) => {
|
|
const names = [`${revision.id}.yaml`, `${revision.id}.env.example`, `${revision.id}.md`];
|
|
const workspace = parseWorkspaceYaml(readFileSync(join(directory, `${revision.id}.yaml`), "utf8"));
|
|
if (workspace.evidence?.source.type === "filesystem") {
|
|
names.push(`${revision.id}/evidence.manifest.json`);
|
|
}
|
|
return names;
|
|
});
|
|
}
|
|
|
|
private evidenceMaterializationLimits(): Partial<{
|
|
maxEntries: number;
|
|
maxTotalBytes: number;
|
|
maxFileBytes: number;
|
|
maxPathBytes: number;
|
|
maxManifestBytes: number;
|
|
}> {
|
|
return {
|
|
...(this.config.maxEvidenceEntries === undefined ? {} : { maxEntries: this.config.maxEvidenceEntries }),
|
|
...(this.config.maxEvidenceBytes === undefined ? {} : { maxTotalBytes: this.config.maxEvidenceBytes }),
|
|
...(this.config.maxEvidenceFileBytes === undefined ? {} : { maxFileBytes: this.config.maxEvidenceFileBytes }),
|
|
...(this.config.maxEvidencePathBytes === undefined ? {} : { maxPathBytes: this.config.maxEvidencePathBytes }),
|
|
...(this.config.maxEvidenceManifestBytes === undefined ? {} : { maxManifestBytes: this.config.maxEvidenceManifestBytes }),
|
|
};
|
|
}
|
|
|
|
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;
|
|
});
|
|
}
|
|
|
|
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;
|
|
}
|
|
}
|
|
}
|