Files
ThothII/backend/src/workspaces/git-repository.ts
T

563 lines
22 KiB
TypeScript

import { execFile, spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import { lstatSync, mkdirSync } from "node:fs";
import { mkdir, rm, writeFile } from "node:fs/promises";
import { basename, dirname, isAbsolute, join } from "node:path";
import { promisify, TextDecoder } 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 isValidUtf8(buffer: Buffer): boolean {
try {
new TextDecoder("utf-8", { fatal: true }).decode(buffer);
return true;
} catch {
return false;
}
}
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> {
try {
for (const path of [this.root, this.snapshotsPath, this.statePath, this.locksPath, this.hooksPath]) {
await mkdir(path, { recursive: true, mode: 0o700 });
assertDirectory(path);
}
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry storage is unavailable");
}
}
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 workspaceDirectories(): Promise<string[]> {
const output = await this.git(["ls-tree", "-d", "--name-only", "HEAD"]);
const directories = output.trim() === "" ? [] : output.trim().split("\n");
for (const id of directories) {
if (id === "workspace-docs") continue;
if (!/^[a-z][a-z0-9-]{2,62}$/.test(id)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository contains an invalid path");
}
}
return directories.filter((id) => id !== "workspace-docs").sort();
}
async workspacePaths(): Promise<string[]> {
const paths: string[] = [];
for (const id of await this.workspaceDirectories()) {
const path = `${id}/workspace.yaml`;
const type = await this.gitOptional(["cat-file", "-t", `HEAD:${path}`]);
if (type === undefined) continue;
if (type.trim() !== "blob") {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace descriptor is invalid");
}
paths.push(path);
}
return paths.sort();
}
async readCatalog(revision = "HEAD"): Promise<string> {
return await this.git(["show", `${revision}:thoth-workspaces.yaml`], {}, "Workspace catalog is invalid");
}
async catalogBlob(revision = "HEAD"): Promise<string> {
return (await this.git(["rev-parse", `${revision}:thoth-workspaces.yaml`], {},
"Workspace catalog is invalid")).trim();
}
async readWorkspace(path: string, revision = "HEAD"): Promise<string> {
if (!/^(?!workspace-docs\/)[a-z][a-z0-9-]{2,62}\/workspace\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
return await this.git(["show", `${revision}:${path}`]);
}
async blob(path: string, revision = "HEAD"): Promise<string> {
if (!/^(?!workspace-docs\/)[a-z][a-z0-9-]{2,62}\/workspace\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
return (await this.git(["rev-parse", `${revision}:${path}`])).trim();
}
/** Read-only object type at an exact revision, or undefined when absent. */
async gitObjectType(revision: string, path: string): Promise<string | undefined> {
if (!/^[0-9a-f]{40}$/.test(revision)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
}
if (!/^[a-z][a-z0-9-]{2,62}\/workspace\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
return await this.gitOptional(["cat-file", "-t", `${revision}:${path}`]);
}
/** Read a generated-doc blob at an exact revision, or undefined when absent. */
async readObjectOrAbsent(revision: string, path: string): Promise<string | undefined> {
if (!/^[0-9a-f]{40}$/.test(revision)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
}
if (!/^workspace-docs\/[a-z][a-z0-9-]{2,62}\/(?:contract\.env\.example|README\.md)$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
}
const output = await this.gitOptional(["show", `${revision}:${path}`]);
return output === undefined ? undefined : output;
}
/** List committed generated-doc paths at an exact revision. */
async workspaceDocsPaths(revision: string): Promise<string[]> {
if (!/^[0-9a-f]{40}$/.test(revision)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
}
const output = await this.git(["ls-tree", "-r", "--name-only", revision, "--", "workspace-docs"]);
if (output.trim() === "") return [];
const paths = output.trim().split("\n");
for (const path of paths) {
if (!/^workspace-docs\/[a-z][a-z0-9-]{2,62}\/(?:contract\.env\.example|README\.md)$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository contains an invalid docs path");
}
}
return paths;
}
/** Assert that a canonical Evidence root is a Git tree at an exact commit. */
async assertTreeAtRevision(revision: string, repoRelativePath: string): Promise<void> {
if (!/^[0-9a-f]{40}$/.test(revision)
|| !/^[a-z][a-z0-9-]{2,62}\/evidence$/.test(repoRelativePath)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence revision is invalid");
}
const type = (await this.git(
["cat-file", "-t", `${revision}:${repoRelativePath}`],
{},
"Workspace Evidence root is invalid",
)).trim();
if (type !== "tree") {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence root is invalid");
}
}
/** Read the curated FK annotations object at an exact commit, or undefined when absent. */
async annotationsObject(revision: string, id: string): Promise<{ blobId: string; contents: Buffer } | undefined> {
if (!/^[0-9a-f]{40}$/.test(revision)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations revision is invalid");
}
if (!/^[a-z][a-z0-9-]{2,62}$/.test(id) || id === "workspace-docs") {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations path is invalid");
}
const path = `${id}/schema/annotations.yaml`;
// ls-tree -z reports the exact object at the path (or its children when the path is a tree).
const listing = await this.git(["ls-tree", "-z", "--full-tree", revision, "--", path]);
const entries = listing.split("\0").filter((entry) => entry.length > 0);
if (entries.length === 0) return undefined;
const exact = entries.find((entry) => entry.slice(entry.lastIndexOf("\t") + 1) === path);
if (exact === undefined) {
// The path resolves to a tree (its children are listed) or another non-blob object.
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is invalid");
}
const match = /^([0-9]{6})\s+(blob|tree|commit)\s+([0-9a-f]{40})\t/.exec(exact);
// Only regular Git blobs are accepted: symlinks (120000) and gitlinks (160000) are refused.
if (match === null || match[2] !== "blob" || (match[1] !== "100644" && match[1] !== "100755")) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is invalid");
}
const blobId = match[3];
const contents = await this.gitBlobBuffer(blobId, 16 * 1024 * 1024);
if (!isValidUtf8(contents)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is not valid UTF-8");
}
return { blobId, contents };
}
private async gitBlobBuffer(objectId: string, maxBytes: number): Promise<Buffer> {
if (!/^[0-9a-f]{40}$/.test(objectId)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is invalid");
}
try {
const { stdout } = await execFileAsync(
"git",
["-c", `core.hooksPath=${this.hooksPath}`, "cat-file", "blob", objectId],
{
cwd: this.repoPath,
env: { ...process.env, GIT_TERMINAL_PROMPT: "0" },
encoding: "buffer",
maxBuffer: maxBytes + 1024 * 1024,
},
);
if (stdout.length > maxBytes) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is too large");
}
return stdout;
} catch (error) {
if (error instanceof WorkspaceRegistryError) throw error;
const detail = error instanceof Error ? error.message : "";
if (/maxBuffer|stdout maxBuffer/i.test(detail)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is too large");
}
throw this.sanitizeGitError(error);
}
}
/** Write only a validated API-owned 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 });
}
/** Create a descriptor only when no filesystem entry exists at its exact path. */
async createRegistryFile(path: string, source: string): Promise<void> {
this.assertRegistryArtifactPath(path);
if (!/^(?!workspace-docs\/)[a-z][a-z0-9-]{2,62}\/workspace\.yaml$/.test(path)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace descriptor path is invalid");
}
const target = join(this.repoPath, path);
await mkdir(dirname(target), { recursive: true, mode: 0o700 });
try {
await writeFile(target, source, { encoding: "utf8", mode: 0o600, flag: "wx" });
} catch {
throw new WorkspaceRegistryError("workspace_curator_owned", "Workspace descriptor is curator-owned");
}
}
async removeRegistryFile(path: string): Promise<void> {
this.assertRegistryArtifactPath(path);
await rm(join(this.repoPath, path), { force: true });
}
private pendingPublicationPaths: string[] = [];
/** 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");
}
this.pendingPublicationPaths = [...paths];
try {
await this.git(["add", "--", ...paths]);
await this.git(["commit", "-m", message], this.publicationIdentity());
await this.git(["push", "origin", `HEAD:${this.config.branch}`]);
return await this.status();
} catch (error) {
// A failed commit leaves staged/working changes; a failed push leaves an ahead commit.
// Restore the last fetched remote revision so the next refresh or explicit retry starts clean.
await this.restoreFailedPublication();
throw error;
}
}
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 isRegistryArtifactPath(path: string): boolean {
return /^(?!workspace-docs\/)[a-z][a-z0-9-]{2,62}\/workspace\.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");
}
if (this.config.remoteUrl) {
await this.git(["remote", "set-url", "origin", "--", this.config.remoteUrl]);
}
await this.git(["fetch", "--no-tags", "origin", this.config.branch]);
const remoteHead = (await this.git(["rev-parse", "FETCH_HEAD"])).trim();
const localHead = (await this.git(["rev-parse", "HEAD"])).trim();
if (localHead !== remoteHead) {
const commonAncestor = (await this.git(["merge-base", "HEAD", "FETCH_HEAD"])).trim();
if (commonAncestor !== localHead) {
throw new WorkspaceRegistryError("git_non_fast_forward", "Workspace checkout diverged from remote");
}
await this.git(["merge", "--ff-only", "FETCH_HEAD"]);
}
if ((await this.git(["rev-parse", "HEAD"])).trim() !== remoteHead) {
throw new WorkspaceRegistryError("git_non_fast_forward", "Workspace checkout does not match remote");
}
}
private publicationIdentity(): NodeJS.ProcessEnv {
return {
GIT_AUTHOR_NAME: this.config.gitAuthorName,
GIT_AUTHOR_EMAIL: this.config.gitAuthorEmail,
GIT_COMMITTER_NAME: this.config.gitAuthorName,
GIT_COMMITTER_EMAIL: this.config.gitAuthorEmail,
};
}
private async restoreFailedPublication(): Promise<void> {
try {
await this.git(["reset", "--hard", `refs/remotes/origin/${this.config.branch}`]);
// Remove only the exact untracked files this publication created, never curated content.
const untracked = this.pendingPublicationPaths.filter((path) => {
try {
lstatSync(join(this.repoPath, path));
return true;
} catch {
return false;
}
});
if (untracked.length > 0) {
await this.git(["clean", "-fd", "--", ...untracked]);
}
this.pendingPublicationPaths = [];
} catch {
// Keep the original sanitized publish failure. A future refresh will surface any recovery
// problem without leaking the Git failure details through the API.
}
}
private async git(
args: string[],
env: NodeJS.ProcessEnv = {},
invalidObjectMessage?: 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", ...env } },
);
return stdout;
} catch (error) {
const stderr = typeof error === "object" && error !== null && "stderr" in error
&& typeof error.stderr === "string" ? error.stderr : "";
if (invalidObjectMessage
&& /^fatal: path '[^']+' does not exist in '[0-9a-f]{40}'\s*$/u.test(stderr)) {
throw new WorkspaceRegistryError("workspace_invalid", invalidObjectMessage);
}
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;
let holder: ChildProcessWithoutNullStreams | undefined;
const lockPath = join(this.locksPath, "repository.lock");
try {
try {
mkdirSync(this.locksPath, { recursive: true, mode: 0o700 });
assertDirectory(this.locksPath);
} catch (error) {
throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable");
}
try {
holder = await this.acquire(lockPath);
} catch (error) {
throw this.lockError(error);
}
return await operation();
} finally {
try {
if (holder !== undefined) await this.release(holder);
} finally {
releaseQueue();
}
}
}
private async acquire(lockPath: string): Promise<ChildProcessWithoutNullStreams> {
try {
const entry = lstatSync(lockPath);
if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("invalid lock path");
} catch (error) {
if (!(typeof error === "object" && error !== null && "code" in error && error.code === "ENOENT")) {
throw error;
}
}
const holder = spawn("python3", ["-c", WorkspaceRepositoryLock.HOLDER_PROGRAM, lockPath], {
stdio: ["pipe", "pipe", "pipe"],
});
await new Promise<void>((resolve, reject) => {
let output = "";
const fail = (error: WorkspaceRegistryError) => {
holder.stdout.removeAllListeners("data");
reject(error);
};
holder.once("error", () => fail(new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable")));
holder.once("exit", (code) => {
fail(new WorkspaceRegistryError(
code === 73 ? "workspace_stale" : "git_unavailable",
code === 73 ? "Workspace registry is busy" : "Workspace registry lock is unavailable",
));
});
holder.stdout.on("data", (chunk: Buffer) => {
output += chunk.toString("utf8");
if (output === "locked\n") {
holder.stdout.removeAllListeners("data");
resolve();
}
});
});
return holder;
}
private async release(holder: ChildProcessWithoutNullStreams): Promise<void> {
if (!holder.stdin.destroyed) holder.stdin.end();
await new Promise<void>((resolve) => holder.once("exit", () => resolve()));
}
private lockError(error: unknown): WorkspaceRegistryError {
if (error instanceof WorkspaceRegistryError) return error;
if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") {
return new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy");
}
return new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable");
}
private static readonly HOLDER_PROGRAM = [
"import fcntl, os, sys",
"fd = os.open(sys.argv[1], os.O_RDWR | os.O_CREAT | getattr(os, 'O_NOFOLLOW', 0), 0o600)",
"try:",
" fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)",
"except BlockingIOError:",
" sys.exit(73)",
"sys.stdout.write('locked\\n')",
"sys.stdout.flush()",
"sys.stdin.buffer.read()",
].join("\n");
}