feat: manage workspace Git checkout and snapshots

This commit is contained in:
2026-08-03 22:21:10 +02:00
parent 5d7ebc5b01
commit 2087fbb0c9
5 changed files with 707 additions and 0 deletions
+6
View File
@@ -15,6 +15,7 @@ import { settingsRoutes, effectiveSettings } from "./routes/settings.js";
import { createPiModelLister } from "./pi/list-models.js";
import { loadSettings, saveSettings, type Settings } from "./settings/settings-store.js";
import { ReadinessManager } from "./runtime/readiness-manager.js";
import { WorkspaceRegistry } from "./workspaces/registry.js";
export interface BuildAppDeps {
thtRunner?: ThtRunner;
@@ -24,6 +25,7 @@ export interface BuildAppDeps {
getSettings?: (principal?: PrincipalContext) => Settings | Promise<Settings>;
readiness?: ReadinessManager;
hub?: SseHub;
workspaceRegistry?: WorkspaceRegistry;
}
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
@@ -47,6 +49,10 @@ 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;
const readiness = deps?.readiness ?? new ReadinessManager(
tht as ThtRunner,
Math.round(config.ollamaEnsureTimeoutMs / 1000),
+230
View File
@@ -0,0 +1,230 @@
import { execFile } from "node:child_process";
import { constants, lstatSync, mkdirSync, openSync, closeSync, unlinkSync } from "node:fs";
import { access, lstat, mkdir } from "node:fs/promises";
import { basename, isAbsolute, join } from "node:path";
import { promisify } from "node:util";
import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js";
const execFileAsync = promisify(execFile);
export interface GitStatus {
branch: string;
head?: string;
ahead: number;
behind: number;
degraded: boolean;
lastError?: WorkspaceErrorCode;
}
export class WorkspaceRegistryError extends Error {
constructor(readonly code: WorkspaceErrorCode, message: string) {
super(message);
this.name = "WorkspaceRegistryError";
}
}
function isMissing(path: string): boolean {
try {
lstatSync(path);
return false;
} catch {
return true;
}
}
function assertDirectory(path: string): void {
const entry = lstatSync(path);
if (!entry.isDirectory() || entry.isSymbolicLink()) {
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry path is unavailable");
}
}
function gitErrorCode(error: unknown): WorkspaceErrorCode {
const detail = [
error instanceof Error ? error.message : "",
typeof error === "object" && error !== null && "stderr" in error
? String((error as { stderr?: unknown }).stderr ?? "")
: "",
].join("\n").toLowerCase();
if (/authentication failed|could not read username|permission denied|publickey/.test(detail)) {
return "git_auth_failed";
}
if (/non-fast-forward|not possible to fast-forward|fast-forward/.test(detail)) {
return "git_non_fast_forward";
}
if (/remote rejected|pre-receive hook declined|push.*rejected/.test(detail)) {
return "git_push_rejected";
}
return "git_unavailable";
}
/**
* A persistent checkout that executes Git only through fixed argument vectors. Git's stdout and
* stderr are intentionally never exposed: they can contain remote URLs or credential hints.
*/
export class GitWorkspaceRepository {
readonly root: string;
readonly repoPath: string;
readonly snapshotsPath: string;
readonly statePath: string;
readonly locksPath: string;
private readonly hooksPath: string;
constructor(private readonly config: WorkspaceRegistryConfig) {
if (!isAbsolute(config.root)) {
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry root is unavailable");
}
this.root = config.root;
this.repoPath = join(this.root, "repo");
this.snapshotsPath = join(this.root, "snapshots");
this.statePath = join(this.root, "state");
this.locksPath = join(this.root, "locks");
this.hooksPath = join(this.locksPath, "empty-hooks");
}
async ensureLayout(): Promise<void> {
for (const path of [this.root, this.snapshotsPath, this.statePath, this.locksPath, this.hooksPath]) {
await mkdir(path, { recursive: true, mode: 0o700 });
assertDirectory(path);
}
}
async bootstrap(): Promise<GitStatus> {
await this.ensureLayout();
if (isMissing(this.repoPath)) {
if (!this.config.remoteUrl) {
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry remote is unavailable");
}
await this.clone();
} else {
assertDirectory(this.repoPath);
await this.refresh();
}
return await this.status();
}
async pull(): Promise<GitStatus> {
await this.ensureLayout();
if (isMissing(this.repoPath)) return await this.bootstrap();
assertDirectory(this.repoPath);
await this.refresh();
return await this.status();
}
async status(): Promise<GitStatus> {
const head = (await this.git(["rev-parse", "HEAD"])).trim();
const tracking = await this.gitOptional(["rev-list", "--left-right", "--count", "HEAD...@{upstream}"]);
const [ahead = "0", behind = "0"] = tracking ? tracking.trim().split(/\s+/) : [];
return {
branch: this.config.branch,
head,
ahead: Number(ahead),
behind: Number(behind),
degraded: false,
};
}
async workspacePaths(): Promise<string[]> {
const output = await this.git(["ls-tree", "-r", "--name-only", "HEAD", "--", "workspaces"]);
const paths = output.trim() === "" ? [] : output.trim().split("\n");
for (const path of paths) {
if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository contains an invalid path");
}
}
return paths;
}
async readWorkspace(path: string): Promise<string> {
if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
return await this.git(["show", `HEAD:${path}`]);
}
async blob(path: string): Promise<string> {
if (!/^workspaces\/[a-z][a-z0-9-]{2,62}\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
return (await this.git(["rev-parse", `HEAD:${path}`])).trim();
}
private async clone(): Promise<void> {
try {
await execFileAsync("git", [
"-c", `core.hooksPath=${this.hooksPath}`,
"clone", "--branch", this.config.branch, "--single-branch", "--", this.config.remoteUrl!, this.repoPath,
], { cwd: this.root, env: { ...process.env, GIT_TERMINAL_PROMPT: "0" } });
assertDirectory(this.repoPath);
} catch (error) {
throw this.sanitizeGitError(error);
}
}
private async refresh(): Promise<void> {
if (this.config.remoteUrl) {
await this.git(["remote", "set-url", "origin", "--", this.config.remoteUrl]);
}
await this.git(["fetch", "--no-tags", "origin", this.config.branch]);
await this.git(["merge", "--ff-only", "FETCH_HEAD"]);
}
private async git(args: string[]): Promise<string> {
try {
const { stdout } = await execFileAsync(
"git",
["-c", `core.hooksPath=${this.hooksPath}`, ...args],
{ cwd: this.repoPath, env: { ...process.env, GIT_TERMINAL_PROMPT: "0" } },
);
return stdout;
} catch (error) {
throw this.sanitizeGitError(error);
}
}
private async gitOptional(args: string[]): Promise<string | undefined> {
try {
return await this.git(args);
} catch (error) {
if (error instanceof WorkspaceRegistryError && error.code === "git_unavailable") return undefined;
throw error;
}
}
private sanitizeGitError(error: unknown): WorkspaceRegistryError {
return new WorkspaceRegistryError(gitErrorCode(error), "Workspace Git operation failed");
}
}
export class WorkspaceRepositoryLock {
private queue = Promise.resolve();
constructor(private readonly locksPath: string) {}
async run<T>(operation: () => Promise<T>): Promise<T> {
const previous = this.queue;
let releaseQueue!: () => void;
this.queue = new Promise<void>((resolve) => { releaseQueue = resolve; });
await previous;
mkdirSync(this.locksPath, { recursive: true, mode: 0o700 });
assertDirectory(this.locksPath);
let descriptor: number | undefined;
const lockPath = join(this.locksPath, "repository.lock");
try {
descriptor = openSync(lockPath, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600);
return await operation();
} catch (error) {
if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") {
throw new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy");
}
throw error;
} finally {
if (descriptor !== undefined) closeSync(descriptor);
if (descriptor !== undefined) {
try { unlinkSync(lockPath); } catch { /* stale lock cleanup is retried by the operator */ }
}
releaseQueue();
}
}
}
+230
View File
@@ -0,0 +1,230 @@
import { 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 { parseWorkspaceYaml, serializeWorkspaceYaml, type CanonicalWorkspace } 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 type PublishWorkspaceRequest =
| { action: "create"; workspace: CanonicalWorkspace; baseCommit: string }
| { action: "update"; workspace: CanonicalWorkspace; baseCommit: string; baseBlob: string }
| { action: "delete"; id: string; baseCommit: string; baseBlob: string };
interface ActiveState {
head: string;
revisions: WorkspaceRevision[];
}
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 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: CanonicalWorkspace; 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);
}
}
/** 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");
}
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: CanonicalWorkspace; blob: string }> = [];
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");
}
buildInstallationContract(workspace);
renderWorkspaceDocs(workspace);
snapshots.push({ id, source: serializeWorkspaceYaml(workspace), workspace, blob: await this.repository.blob(path) });
}
} catch (error) {
throw workspaceError(error);
}
const snapshotDirectory = join(this.repository.snapshotsPath, safeHead);
if (!this.pathExists(snapshotDirectory)) {
const staging = join(this.repository.snapshotsPath, `.staging-${randomUUID()}`);
await mkdir(staging, { mode: 0o700 });
try {
const revisions: WorkspaceRevision[] = [];
for (const snapshot of snapshots) {
const path = join(staging, `${snapshot.id}.yaml`);
const docs = renderWorkspaceDocs(snapshot.workspace);
await writeFile(path, snapshot.source, { encoding: "utf8", mode: 0o400 });
await writeFile(join(staging, `${snapshot.id}.env.example`), docs.envExample, { encoding: "utf8", mode: 0o400 });
await writeFile(join(staging, `${snapshot.id}.md`), docs.markdown, { encoding: "utf8", mode: 0o400 });
revisions.push({ id: snapshot.id, commit: safeHead, blob: snapshot.blob, snapshotPath: this.snapshotPath(safeHead, snapshot.id) });
}
await writeFile(join(staging, "snapshot.json"), JSON.stringify({ head: safeHead, revisions }), {
encoding: "utf8", mode: 0o400,
});
await rename(staging, snapshotDirectory);
} catch (error) {
await rm(staging, { recursive: true, force: true });
throw error;
}
}
const revisions = snapshots.map((snapshot) => ({
id: snapshot.id,
commit: safeHead,
blob: snapshot.blob,
snapshotPath: this.snapshotPath(safeHead, snapshot.id),
}));
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 state = JSON.parse(await readFile(file, "utf8")) as ActiveState;
safeCommit(state.head);
if (!Array.isArray(state.revisions) || state.revisions.length === 0) throw new Error("bad state");
for (const revision of state.revisions) {
safeCommit(revision.commit);
workspacePath(revision.id);
if (!isAbsolute(revision.snapshotPath) || revision.snapshotPath !== this.snapshotPath(revision.commit, revision.id)) {
throw new Error("bad snapshot path");
}
}
return state;
} catch {
return undefined;
}
}
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 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;
}
}
}