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

622 lines
25 KiB
TypeScript

import { execFile, spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import { lstatSync, mkdirSync } from "node:fs";
import { mkdir } from "node:fs/promises";
import { 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;
repository?: WorkspaceRepositoryIdentity;
head?: string;
ahead: number;
behind: number;
degraded: boolean;
lastError?: WorkspaceErrorCode;
}
export interface WorkspaceRepositoryIdentity {
host: string;
repository: string;
transport: "https" | "ssh" | "local";
}
export interface EvidenceTreeObject {
mode: "100644" | "100755";
oid: string;
posixPath: string;
}
export class WorkspaceRegistryError extends Error {
constructor(readonly code: WorkspaceErrorCode, message: string) {
super(message);
this.name = "WorkspaceRegistryError";
}
}
function invalidRemote(): never {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Git remote is invalid");
}
function safeRepositoryPath(raw: string): string {
let decoded: string;
try {
decoded = decodeURIComponent(raw).replace(/^\/+/, "").replace(/\/+$/, "").replace(/\.git$/, "");
} catch {
return invalidRemote();
}
if (
decoded.length === 0
|| decoded.includes("\\")
|| decoded.split("/").some((part) => part === "" || part === "." || part === "..")
|| /[\p{Cc}\s?#]/u.test(decoded)
) return invalidRemote();
return decoded;
}
/** Convert a configured remote to the only repository identity safe for API/UI responses. */
export function normalizeRepositoryIdentity(remote: string): WorkspaceRepositoryIdentity {
if (remote.length === 0 || remote.trim() !== remote || remote.includes("\0")) return invalidRemote();
if (isAbsolute(remote) || remote.startsWith("file://")) {
return { host: "local", repository: "configured-repository", transport: "local" };
}
const scp = /^git@([^:/\s]+):(.+)$/.exec(remote);
if (scp) {
return { host: scp[1].toLowerCase(), repository: safeRepositoryPath(scp[2]), transport: "ssh" };
}
let parsed: URL;
try {
parsed = new URL(remote);
} catch {
return invalidRemote();
}
if (parsed.search || parsed.hash || !parsed.hostname || parsed.port) return invalidRemote();
if (parsed.protocol === "https:") {
if (parsed.username || parsed.password) return invalidRemote();
return {
host: parsed.hostname.toLowerCase(),
repository: safeRepositoryPath(parsed.pathname),
transport: "https",
};
}
if (parsed.protocol === "ssh:") {
if (parsed.password || (parsed.username !== "" && parsed.username !== "git")) return invalidRemote();
return {
host: parsed.hostname.toLowerCase(),
repository: safeRepositoryPath(parsed.pathname),
transport: "ssh",
};
}
return invalidRemote();
}
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";
}
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;
private readonly identity?: WorkspaceRepositoryIdentity;
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");
this.identity = config.remoteUrl === undefined
? undefined
: normalizeRepositoryIdentity(config.remoteUrl);
}
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,
...(this.identity ? { repository: this.identity } : {}),
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.gitBlobBytes(blobId, 16 * 1024 * 1024, "annotations");
if (!isValidUtf8(contents)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace annotations object is not valid UTF-8");
}
return { blobId, contents };
}
/**
* Recursively enumerate a canonical `<id>/evidence` tree at an exact commit as regular Git blobs.
* Symlinks (120000), gitlinks (160000), non-regular modes, non-blob types, traversal/absolute/
* duplicate/cross-namespace paths, and NUL/newline-bearing names are refused.
*/
async evidenceTreeObjects(revision: string, id: string): Promise<EvidenceTreeObject[]> {
if (!/^[0-9a-f]{40}$/.test(revision)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence revision is invalid");
}
if (!/^[a-z][a-z0-9-]{2,62}$/.test(id)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence path is invalid");
}
const prefix = `${id}/evidence`;
const listing = await this.git(["ls-tree", "-r", "-z", "--full-tree", revision, "--", prefix]);
const entries = listing.split("\0").filter((entry) => entry.length > 0);
const seen = new Set<string>();
const objects: EvidenceTreeObject[] = [];
for (const entry of entries) {
const tab = entry.lastIndexOf("\t");
if (tab < 0) throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence object is invalid");
const name = entry.slice(tab + 1);
const meta = entry.slice(0, tab);
const match = /^([0-9]{6}) (blob|commit|tree) ([0-9a-f]{40})$/.exec(meta);
if (match === null || match[2] !== "blob" || (match[1] !== "100644" && match[1] !== "100755")) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence object is invalid");
}
if (name === prefix) {
// The Evidence root resolves to a single regular blob (or symlink/gitlink already refused above).
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence root is invalid");
}
if (!name.startsWith(`${prefix}/`)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence object escapes its namespace");
}
const rel = name.slice(prefix.length + 1);
if (rel.length === 0 || rel.includes("\0") || rel.includes("\n") || rel.includes("\r")) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence path is invalid");
}
const segments = rel.split("/");
if (segments.some((segment) => segment === "" || segment === "." || segment === "..")) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence path is invalid");
}
if (seen.has(rel)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence path is duplicated");
seen.add(rel);
objects.push({ mode: match[1] as "100644" | "100755", oid: match[3], posixPath: rel });
}
return objects;
}
/** Read one Evidence blob with a per-object byte bound. */
evidenceBlobBytes(objectId: string, maxBytes: number): Promise<Buffer> {
return this.gitBlobBytes(objectId, maxBytes, "Evidence");
}
/** Return the 40-hex tree id of a canonical Evidence root at an exact commit. */
async evidenceTreeId(revision: string, id: string): Promise<string> {
if (!/^[0-9a-f]{40}$/.test(revision) || !/^[a-z][a-z0-9-]{2,62}$/.test(id)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence revision is invalid");
}
const objectId = (await this.git(["rev-parse", `${revision}:${id}/evidence`])).trim();
const type = (await this.git(["cat-file", "-t", objectId])).trim();
if (!/^[0-9a-f]{40}$/.test(objectId) || type !== "tree") {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence root is invalid");
}
return objectId;
}
/** Return the byte size of one Git object without reading its contents. */
async gitObjectSize(objectId: string): Promise<number> {
if (!/^[0-9a-f]{40}$/.test(objectId)) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence object is invalid");
}
const raw = (await this.git(["cat-file", "-s", objectId])).trim();
const size = Number(raw);
if (!Number.isSafeInteger(size) || size < 0) {
throw new WorkspaceRegistryError("workspace_invalid", "Workspace Evidence object size is invalid");
}
return size;
}
private async gitBlobBytes(objectId: string, maxBytes: number, label: string): Promise<Buffer> {
if (!/^[0-9a-f]{40}$/.test(objectId)) {
throw new WorkspaceRegistryError("workspace_invalid", `Workspace ${label} 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 ${label} 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 ${label} object is too large`);
}
throw this.sanitizeGitError(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 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 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");
}