feat: expose workspace registry API

This commit is contained in:
2026-08-04 01:01:57 +02:00
parent 49fa7030a5
commit f95a18ab0d
9 changed files with 829 additions and 50 deletions
+5 -7
View File
@@ -17,6 +17,7 @@ import { loadSettings, saveSettings, type Settings } from "./settings/settings-s
import { ReadinessManager } from "./runtime/readiness-manager.js";
import { WorkspaceRegistry } from "./workspaces/registry.js";
import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js";
import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js";
export interface BuildAppDeps {
thtRunner?: ThtRunner;
@@ -27,6 +28,7 @@ export interface BuildAppDeps {
readiness?: ReadinessManager;
hub?: SseHub;
workspaceRegistry?: WorkspaceRegistry;
workspaceDiagnoser?: WorkspaceDiagnoser;
}
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
@@ -50,14 +52,9 @@ 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();
// Routes are introduced in Task 6; construction here keeps production and injected-test
// dependencies on the same registry lifecycle without performing Git I/O at startup.
const workspaceRegistry = deps?.workspaceRegistry ?? new WorkspaceRegistry(config.workspaceRegistry);
void workspaceRegistry;
// Task 6 consumes this dependency from the registry route. Construct it from the effective
// application configuration here so production diagnostics never silently use test defaults.
const workspaceDiagnoser = createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs);
void workspaceDiagnoser;
const workspaceDiagnoser = deps?.workspaceDiagnoser
?? createProductionWorkspaceDiagnoser(config.workspaceDiagnosticTimeoutMs);
const readiness = deps?.readiness ?? new ReadinessManager(
tht as ThtRunner,
Math.round(config.ollamaEnsureTimeoutMs / 1000),
@@ -113,6 +110,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
});
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
workspaceRoutes(app, { registry: workspaceRegistry, config: config.workspaceRegistry, diagnose: workspaceDiagnoser });
settingsRoutes(app, { cfg: config, listModels, getSettings, saveSettings: saveUserSettings });
return app;
-4
View File
@@ -26,10 +26,6 @@ export function metaRoutes(
app: FastifyInstance,
deps: { harnessDir: string; listModels?: ListModelsFn },
): void {
app.get("/workspaces", async () => {
return listWorkspaces(deps.harnessDir);
});
app.get("/models", async () => {
const fn = deps.listModels ?? (async () => []);
try {
+362
View File
@@ -0,0 +1,362 @@
import { createHash } from "node:crypto";
import { Buffer } from "node:buffer";
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
import multipart from "@fastify/multipart";
import yauzl from "yauzl";
import yazl from "yazl";
import { z } from "zod";
import type { WorkspaceRegistryConfig } from "../workspaces/types.js";
import { WorkspaceRegistryError } from "../workspaces/git-repository.js";
import {
WorkspaceConflictError,
type PublishWorkspaceRequest,
type WorkspaceRegistry,
} from "../workspaces/registry.js";
import { resolveRuntimeBindings } from "../workspaces/bindings.js";
import { buildInstallationContract, renderWorkspaceDocs } from "../workspaces/contracts.js";
import {
isCanonicalWorkspace,
parseWorkspaceYaml,
serializeWorkspaceYaml,
validateCanonicalWorkspace,
type CanonicalWorkspace,
type WorkspaceDescriptor,
} from "../workspaces/schema.js";
import type { RuntimeBindings } from "../workspaces/runtime-renderer.js";
import type { WorkspaceDiagnostics } from "../workspaces/diagnostics.js";
export type WorkspaceDiagnoser = (
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
options: { writeProbe: boolean },
) => Promise<WorkspaceDiagnostics>;
interface WorkspaceRoutesDeps {
registry: WorkspaceRegistry;
config: WorkspaceRegistryConfig;
diagnose: WorkspaceDiagnoser;
}
const workspaceId = z.string().regex(/^[a-z][a-z0-9-]{2,62}$/);
const commit = z.string().regex(/^[0-9a-f]{40}$/);
const workspacePayload = z.object({ workspace: z.unknown() }).strict();
const publishPayload = z.discriminatedUnion("action", [
z.object({ action: z.literal("create"), workspace: z.unknown(), baseCommit: commit }).strict(),
z.object({ action: z.literal("update"), workspace: z.unknown(), baseCommit: commit, baseBlob: commit }).strict(),
z.object({ action: z.literal("delete"), id: workspaceId, baseCommit: commit, baseBlob: commit }).strict(),
]);
const bundleManifest = z.object({
schema_version: z.literal(1),
workspace_id: workspaceId,
files: z.object({
"workspace.yaml": z.string().regex(/^[0-9a-f]{64}$/),
"contract.env.example": z.string().regex(/^[0-9a-f]{64}$/),
"README.md": z.string().regex(/^[0-9a-f]{64}$/),
}).strict(),
}).strict();
const BUNDLE_FILES = ["manifest.json", "workspace.yaml", "contract.env.example", "README.md"] as const;
type BundleFile = (typeof BUNDLE_FILES)[number];
const SAFE_MESSAGES = {
workspace_invalid: "Workspace request or bundle is invalid.",
binding_missing: "Installation binding is missing or invalid.",
workspace_not_activatable: "Workspace cannot be activated on this installation.",
workspace_stale: "Workspace revision is stale.",
workspace_conflict: "Workspace changed in the registry.",
git_unavailable: "Workspace Git service is unavailable.",
git_auth_failed: "Workspace Git authentication failed.",
git_non_fast_forward: "Workspace Git branch has changed.",
git_push_rejected: "Workspace Git publication was rejected.",
connector_unavailable: "Workspace connector is unavailable.",
semantic_index_incompatible: "Semantic index is incompatible with this workspace.",
} as const;
function sha256(value: string | Buffer): string {
return createHash("sha256").update(value).digest("hex");
}
function invalidBundle(): WorkspaceRegistryError {
return new WorkspaceRegistryError("workspace_invalid", "Workspace bundle is invalid");
}
function isBundleFile(value: string): value is BundleFile {
return (BUNDLE_FILES as readonly string[]).includes(value);
}
function unsafeArchiveEntry(entry: yauzl.Entry): boolean {
const name = entry.fileName;
const unixType = (entry.externalFileAttributes >>> 16) & 0o170000;
return name.length === 0
|| name.startsWith("/")
|| name.startsWith("\\")
|| name.includes("\\")
|| name.split("/").includes("..")
|| name.endsWith("/")
|| unixType === 0o120000
|| !isBundleFile(name);
}
async function readZipBundle(source: Buffer, config: WorkspaceRegistryConfig): Promise<Record<BundleFile, Buffer>> {
if (source.length === 0 || source.length > config.maxImportBytes) throw invalidBundle();
return await new Promise<Record<BundleFile, Buffer>>((resolve, reject) => {
yauzl.fromBuffer(source, {
lazyEntries: true,
strictFileNames: true,
validateEntrySizes: true,
decodeStrings: true,
}, (error, archive) => {
if (error || !archive) return reject(invalidBundle());
const files = new Map<BundleFile, Buffer>();
let entries = 0;
let settled = false;
const fail = () => {
if (settled) return;
settled = true;
archive.close();
reject(invalidBundle());
};
archive.on("error", fail);
archive.on("entry", (entry) => {
entries += 1;
if (entries > config.maxImportEntries || unsafeArchiveEntry(entry) || files.has(entry.fileName as BundleFile)) {
fail();
return;
}
if (entry.uncompressedSize > config.maxImportBytes) {
fail();
return;
}
archive.openReadStream(entry, (streamError, stream) => {
if (streamError || !stream) return fail();
const chunks: Buffer[] = [];
let size = 0;
stream.on("data", (chunk: Buffer) => {
size += chunk.length;
if (size > config.maxImportBytes) return fail();
chunks.push(chunk);
});
stream.on("error", fail);
stream.on("end", () => {
if (settled || size !== entry.uncompressedSize) return fail();
files.set(entry.fileName as BundleFile, Buffer.concat(chunks));
archive.readEntry();
});
});
});
archive.on("end", () => {
if (settled) return;
settled = true;
if (entries !== BUNDLE_FILES.length || BUNDLE_FILES.some((name) => !files.has(name))) return reject(invalidBundle());
resolve(Object.fromEntries(files) as Record<BundleFile, Buffer>);
});
archive.readEntry();
});
});
}
function utf8(buffer: Buffer): string {
const text = buffer.toString("utf8");
if (!Buffer.from(text, "utf8").equals(buffer) || text.includes("\0")) throw invalidBundle();
return text;
}
async function importDraft(source: Buffer, config: WorkspaceRegistryConfig): Promise<CanonicalWorkspace> {
const files = await readZipBundle(source, config);
let manifest: z.infer<typeof bundleManifest>;
try {
manifest = bundleManifest.parse(JSON.parse(utf8(files["manifest.json"])));
} catch {
throw invalidBundle();
}
for (const name of ["workspace.yaml", "contract.env.example", "README.md"] as const) {
if (sha256(files[name]) !== manifest.files[name]) throw invalidBundle();
}
try {
const descriptor = parseWorkspaceYaml(utf8(files["workspace.yaml"]));
const workspace = validateCanonicalWorkspace(descriptor);
const docs = renderWorkspaceDocs(workspace);
if (
workspace.workspace.id !== manifest.workspace_id
|| serializeWorkspaceYaml(workspace) !== utf8(files["workspace.yaml"])
|| docs.envExample !== utf8(files["contract.env.example"])
|| docs.markdown !== utf8(files["README.md"])
) throw invalidBundle();
return workspace;
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
throw invalidBundle();
}
}
async function exportBundle(workspace: CanonicalWorkspace): Promise<Buffer> {
const yaml = serializeWorkspaceYaml(workspace);
const docs = renderWorkspaceDocs(workspace);
const files: Record<BundleFile, string> = {
"manifest.json": JSON.stringify({
schema_version: 1,
workspace_id: workspace.workspace.id,
files: {
"workspace.yaml": sha256(yaml),
"contract.env.example": sha256(docs.envExample),
"README.md": sha256(docs.markdown),
},
}),
"workspace.yaml": yaml,
"contract.env.example": docs.envExample,
"README.md": docs.markdown,
};
const archive = new yazl.ZipFile();
const chunks: Buffer[] = [];
archive.outputStream.on("data", (chunk: Buffer) => chunks.push(chunk));
for (const name of BUNDLE_FILES) archive.addBuffer(Buffer.from(files[name]), name);
archive.end();
await new Promise<void>((resolve, reject) => {
archive.outputStream.once("end", resolve);
archive.outputStream.once("error", reject);
});
return Buffer.concat(chunks);
}
function workspaceErrorCode(error: unknown): keyof typeof SAFE_MESSAGES {
return error instanceof WorkspaceRegistryError ? error.code : "workspace_invalid";
}
function workspaceErrorStatus(code: keyof typeof SAFE_MESSAGES): number {
if (code === "workspace_conflict" || code === "workspace_stale" || code === "git_non_fast_forward") return 409;
if (code === "git_unavailable" || code === "git_auth_failed" || code === "git_push_rejected") return 503;
return 400;
}
function errorReply(reply: FastifyReply, error: unknown) {
const code = workspaceErrorCode(error);
const body: Record<string, unknown> = { code, message: SAFE_MESSAGES[code] };
if (error instanceof WorkspaceConflictError) {
body.fields = error.fields;
if (error.base) body.base = error.base;
if (error.local) body.local = error.local;
if (error.remote) body.remote = error.remote;
} else if (code === "workspace_conflict" && error && typeof error === "object") {
const conflict = error as Partial<WorkspaceConflictError>;
if (Array.isArray(conflict.fields) && conflict.fields.every((field) => typeof field === "string")) body.fields = conflict.fields;
for (const key of ["base", "local", "remote"] as const) {
if (conflict[key] && isCanonicalWorkspace(conflict[key] as WorkspaceDescriptor)) body[key] = conflict[key];
}
}
return reply.code(workspaceErrorStatus(code)).send(body);
}
function publishRequest(value: unknown): PublishWorkspaceRequest {
const parsed = publishPayload.parse(value);
if (parsed.action === "delete") return parsed;
return { ...parsed, workspace: validateCanonicalWorkspace(parsed.workspace) };
}
export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps): void {
app.register(multipart, {
limits: { fileSize: deps.config.maxImportBytes, files: 1, fields: 0, parts: 1 },
throwFileSizeLimit: true,
});
app.get("/workspace-registry/status", async (_request, reply) => {
try {
return await deps.registry.bootstrap();
} catch (error) {
return errorReply(reply, error);
}
});
app.post("/workspace-registry/pull", async (_request, reply) => {
try {
return await deps.registry.pull();
} catch (error) {
return errorReply(reply, error);
}
});
app.get("/workspaces", async (_request, reply) => {
try {
const revisions = await deps.registry.list();
return await Promise.all(revisions.map(async (revision) => {
const { workspace } = await deps.registry.read(revision.id);
return {
id: revision.id,
// Retain the metadata endpoint's selector fields while adding registry summary data.
name: revision.id,
file: `${revision.id}.yaml`,
displayName: workspace.workspace.name,
description: workspace.workspace.description,
language: workspace.workspace.language,
revision,
};
}));
} catch (error) {
return errorReply(reply, error);
}
});
app.get("/workspaces/:id", async (request, reply) => {
try {
const { id } = z.object({ id: workspaceId }).parse(request.params);
return await deps.registry.read(id);
} catch (error) {
return errorReply(reply, error);
}
});
app.post("/workspaces/validate", async (request, reply) => {
try {
const { workspace } = workspacePayload.parse(request.body);
const canonical = validateCanonicalWorkspace(workspace);
return { workspace: canonical, contract: buildInstallationContract(canonical) };
} catch (error) {
return errorReply(reply, error);
}
});
app.post("/workspaces/:id/test", async (request, reply) => {
try {
const { id } = z.object({ id: workspaceId }).parse(request.params);
const { workspace } = await deps.registry.read(id);
const bindings = resolveRuntimeBindings(workspace, process.env, deps.config.secretRoots);
return await deps.diagnose(workspace, bindings, { writeProbe: false });
} catch (error) {
return errorReply(reply, error);
}
});
app.post("/workspaces/publish", async (request, reply) => {
try {
const result = await deps.registry.publish(publishRequest(request.body));
return result ? { revision: result } : reply.code(204).send();
} catch (error) {
return errorReply(reply, error);
}
});
app.get("/workspaces/:id/export", async (request, reply) => {
try {
const { id } = z.object({ id: workspaceId }).parse(request.params);
const { workspace } = await deps.registry.read(id);
const canonical = validateCanonicalWorkspace(workspace);
const bundle = await exportBundle(canonical);
return reply
.type("application/zip")
.header("content-disposition", `attachment; filename=\"${id}.zip\"`)
.send(bundle);
} catch (error) {
return errorReply(reply, error);
}
});
app.post("/workspaces/import", async (request: FastifyRequest, reply) => {
try {
const file = await request.file();
if (!file || file.fieldname !== "bundle" || file.mimetype !== "application/zip") throw invalidBundle();
const draft = await importDraft(await file.toBuffer(), deps.config);
return { draft: { workspace: draft, contract: buildInstallationContract(draft) } };
} catch (error) {
return errorReply(reply, error);
}
});
}
+37 -2
View File
@@ -1,7 +1,7 @@
import { execFile, spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import { lstatSync, mkdirSync } from "node:fs";
import { mkdir } from "node:fs/promises";
import { basename, isAbsolute, join } from "node:path";
import { mkdir, rm, writeFile } from "node:fs/promises";
import { basename, dirname, isAbsolute, join } from "node:path";
import { promisify } from "node:util";
import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js";
@@ -154,6 +154,30 @@ export class GitWorkspaceRepository {
return (await this.git(["rev-parse", `HEAD:${path}`])).trim();
}
/** Write only a validated registry artifact below the checked-out repository. */
async writeRegistryFile(path: string, source: string): Promise<void> {
this.assertRegistryArtifactPath(path);
const target = join(this.repoPath, path);
await mkdir(dirname(target), { recursive: true, mode: 0o700 });
await writeFile(target, source, { encoding: "utf8", mode: 0o600 });
}
async removeRegistryFile(path: string): Promise<void> {
this.assertRegistryArtifactPath(path);
await rm(join(this.repoPath, path), { force: true });
}
/** Commit and push a fixed set of validated artifact paths without exposing Git output. */
async commitAndPush(paths: readonly string[], message: string): Promise<GitStatus> {
if (paths.length === 0 || paths.some((path) => !this.isRegistryArtifactPath(path))) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
await this.git(["add", "--", ...paths]);
await this.git(["commit", "-m", message]);
await this.git(["push", "origin", `HEAD:${this.config.branch}`]);
return await this.status();
}
private async clone(): Promise<void> {
try {
await execFileAsync("git", [
@@ -166,6 +190,17 @@ export class GitWorkspaceRepository {
}
}
private isRegistryArtifactPath(path: string): boolean {
return /^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)
|| /^workspace-docs\/[a-z][a-z0-9-]{2,62}\/(?:contract\.env\.example|README\.md)$/.test(path);
}
private assertRegistryArtifactPath(path: string): void {
if (!this.isRegistryArtifactPath(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
}
private async refresh(): Promise<void> {
if ((await this.git(["status", "--porcelain"])).trim() !== "") {
throw new WorkspaceRegistryError("workspace_stale", "Workspace checkout has local changes");
+106 -3
View File
@@ -33,6 +33,18 @@ export type PublishWorkspaceRequest =
| { 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 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[];
@@ -139,9 +151,100 @@ export class WorkspaceRegistry {
}
}
/** Publication is deliberately deferred until Task 6 adds validated route-level concurrency controls. */
async publish(_request: PublishWorkspaceRequest): Promise<WorkspaceRevision> {
throw new WorkspaceRegistryError("workspace_stale", "Workspace publication is unavailable");
/**
* 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, existing, local);
}
if (request.action === "create" && existing) throw await this.conflictFor(request, existing, local);
if (request.action !== "create" && !existing) throw await this.conflictFor(request, 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,
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), 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 || !remote) return ["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> {