Files
ThothII/backend/src/workspaces/preprocessing-state.ts
T

403 lines
14 KiB
TypeScript

import { createHash, randomBytes } from "node:crypto";
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import {
closeSync,
constants as fsConstants,
fchmodSync,
fstatSync,
fsyncSync,
lstatSync,
mkdirSync,
openSync,
readFileSync,
renameSync,
unlinkSync,
writeFileSync,
} from "node:fs";
import { dirname, join } from "node:path";
export type PreprocessingConflictCode = "preprocessing_conflict" | "preprocessing_resume_mismatch";
export class PreprocessingStateError extends Error {
constructor(readonly code: PreprocessingConflictCode, message: string) {
super(message);
this.name = "PreprocessingStateError";
}
}
export interface SessionInventoryRow {
id: string;
status: string;
archived?: boolean;
workspaceRevision?: string | null;
}
export interface WriterLockLease {
holderPid: number;
release(): Promise<void>;
}
export interface BeginPreprocessingJobOptions {
operation: string;
workspaceRevision: string;
descriptorBlob: string;
catalogBlob: string;
configDigest: string;
bindingDigest: string;
embeddingId: string;
embeddingDimensions: number;
runId?: string;
}
export interface PreprocessingJobState {
schemaVersion: 2;
runId: string;
operation: string;
workspaceId: string;
workspaceRevision: string;
descriptorBlob: string;
catalogBlob: string;
configDigest: string;
bindingDigest: string;
embeddingId: string;
embeddingDimensions: number;
completedStages: string[];
childRuns: Record<string, string>;
status: "active" | "succeeded" | "blocked" | "failed";
candidateDigest?: string;
reviewDigest?: string;
}
export interface FkReviewRecord {
reviewedCandidatesDigest: string;
annotationsDigest: string;
workspaceRevision: string;
/** Curated Git blob id accepted at review time (P5); absent for legacy host-file reviews. */
blobId?: string;
}
function sha256(value: string | Buffer): string {
return `sha256:${createHash("sha256").update(value).digest("hex")}`;
}
function syncDirectory(directory: string): void {
if (process.platform === "win32") return;
const fd = openSync(directory, "r");
try { fsyncSync(fd); } finally { closeSync(fd); }
}
function ensureDirectory(path: string): string {
mkdirSync(path, { recursive: true, mode: 0o700 });
const entry = lstatSync(path);
if (!entry.isDirectory() || entry.isSymbolicLink()) {
throw new Error("preprocessing state directory is unavailable");
}
return path;
}
function validateRunId(runId: string): string {
if (!/^[0-9a-f]{32}$/.test(runId)) throw new Error("preprocessing run id is invalid");
return runId;
}
function writeAtomicFile(path: string, contents: string, mode: number): void {
ensureDirectory(join(path, ".."));
const directory = path.slice(0, path.lastIndexOf("/"));
ensureDirectory(directory);
const staging = `${path}.tmp-${process.pid}-${Date.now()}-${randomBytes(6).toString("hex")}`;
const fd = openSync(
staging,
fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_NOFOLLOW,
0o600,
);
let closed = false;
try {
writeFileSync(fd, contents, "utf8");
fsyncSync(fd);
fchmodSync(fd, mode);
closeSync(fd);
closed = true;
renameSync(staging, path);
syncDirectory(directory);
} catch (error) {
if (!closed) try { closeSync(fd); } catch { /* preserve original failure */ }
try { unlinkSync(staging); } catch { /* best effort */ }
throw error;
}
}
function readTrustedFile(path: string): string {
const entry = lstatSync(path);
if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("preprocessing state file is invalid");
const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW);
try {
const before = fstatSync(fd);
if (!before.isFile() || before.nlink !== 1) throw new Error("preprocessing state file is invalid");
const contents = readFileSync(fd, "utf8");
const after = fstatSync(fd);
if (
before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size || before.nlink !== after.nlink
) throw new Error("preprocessing state file changed while reading");
return contents;
} finally {
closeSync(fd);
}
}
function decodeJob(value: unknown): PreprocessingJobState {
if (!value || typeof value !== "object" || Array.isArray(value)) {
throw new Error("preprocessing job state is invalid");
}
const record = value as Record<string, unknown>;
if (
record.schemaVersion !== 2
|| typeof record.runId !== "string"
|| typeof record.operation !== "string"
|| typeof record.workspaceId !== "string"
|| typeof record.workspaceRevision !== "string"
|| typeof record.descriptorBlob !== "string"
|| typeof record.catalogBlob !== "string"
|| typeof record.configDigest !== "string"
|| typeof record.bindingDigest !== "string"
|| typeof record.embeddingId !== "string"
|| typeof record.embeddingDimensions !== "number"
|| !Array.isArray(record.completedStages)
|| typeof record.childRuns !== "object" || record.childRuns === null || Array.isArray(record.childRuns)
|| !["active", "succeeded", "blocked", "failed"].includes(String(record.status))
) {
throw new Error("preprocessing job state is invalid");
}
return record as unknown as PreprocessingJobState;
}
function decodeReview(value: unknown): FkReviewRecord {
if (!value || typeof value !== "object" || Array.isArray(value)) {
throw new Error("preprocessing review state is invalid");
}
const record = value as Record<string, unknown>;
if (
typeof record.reviewedCandidatesDigest !== "string"
|| typeof record.annotationsDigest !== "string"
|| typeof record.workspaceRevision !== "string"
|| (record.blobId !== undefined && typeof record.blobId !== "string")
) throw new Error("preprocessing review state is invalid");
return record as unknown as FkReviewRecord;
}
export class PreprocessingStateStore {
private readonly root: string;
constructor(private readonly options: { dataRoot: string; workspaceId: string }) {
if (!/^[a-z][a-z0-9-]{2,62}$/.test(options.workspaceId)) {
throw new Error("workspace id is invalid");
}
this.root = join(options.dataRoot, "sessions", options.workspaceId, "preprocessing");
}
writerLockPath(): string { return join(this.root, "writer.lock"); }
runtimeConfigDirectory(): string { return join(this.root, "runtime-config"); }
runtimeConfigManifestDirectory(): string { return join(this.root, "runtime-config-manifests"); }
jobsDirectory(): string { return join(this.root, "jobs"); }
fkCandidatesDirectory(): string { return join(this.root, "fk-candidates"); }
fkReviewsDirectory(): string { return join(this.root, "fk-reviews"); }
jobPath(runId: string): string { return join(this.jobsDirectory(), `${validateRunId(runId)}.json`); }
fkCandidatesPath(runId: string): string { return join(this.fkCandidatesDirectory(), `${validateRunId(runId)}.yaml`); }
fkReviewPath(runId: string): string { return join(this.fkReviewsDirectory(), `${validateRunId(runId)}.json`); }
private ensureLayout(): void {
ensureDirectory(this.root);
for (const directory of [
this.runtimeConfigDirectory(),
this.runtimeConfigManifestDirectory(),
this.jobsDirectory(),
this.fkCandidatesDirectory(),
this.fkReviewsDirectory(),
]) ensureDirectory(directory);
}
async acquireWriterLock(): Promise<WriterLockLease> {
this.ensureLayout();
const lockPath = this.writerLockPath();
try {
const entry = lstatSync(lockPath);
if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("invalid writer lock path");
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ENOENT") {
throw new PreprocessingStateError("preprocessing_conflict", "Workspace preprocessing lock is unavailable");
}
}
const holder = spawn("python3", ["-c", PreprocessingStateStore.HOLDER_PROGRAM, lockPath], {
stdio: ["pipe", "pipe", "pipe"],
});
await new Promise<void>((resolve, reject) => {
let output = "";
const fail = (error: Error) => {
holder.stdout.removeAllListeners("data");
reject(error);
};
holder.once("error", () => fail(new PreprocessingStateError(
"preprocessing_conflict",
"Workspace preprocessing lock is unavailable",
)));
holder.once("exit", (code) => {
fail(new PreprocessingStateError(
"preprocessing_conflict",
code === 73 ? "Workspace preprocessing is busy" : "Workspace preprocessing lock is unavailable",
));
});
holder.stdout.on("data", (chunk: Buffer) => {
output += chunk.toString("utf8");
if (output === "locked\n") {
holder.stdout.removeAllListeners("data");
resolve();
}
});
});
let released = false;
return {
holderPid: holder.pid ?? 0,
release: async () => {
if (released) return;
released = true;
if (!holder.stdin.destroyed) holder.stdin.end();
await new Promise<void>((resolve) => holder.once("exit", () => resolve()));
},
};
}
async beginJob(options: BeginPreprocessingJobOptions): Promise<PreprocessingJobState> {
this.ensureLayout();
if (!/^[0-9a-f]{40}$/.test(options.workspaceRevision)) {
throw new Error("workspace revision is invalid");
}
const runId = options.runId ?? randomBytes(16).toString("hex");
const path = this.jobPath(runId);
try {
const existing = this.readJob(runId);
if (
existing.operation !== options.operation
|| existing.workspaceRevision !== options.workspaceRevision
|| existing.descriptorBlob !== options.descriptorBlob
|| existing.catalogBlob !== options.catalogBlob
|| existing.configDigest !== options.configDigest
|| existing.bindingDigest !== options.bindingDigest
|| existing.embeddingId !== options.embeddingId
|| existing.embeddingDimensions !== options.embeddingDimensions
) {
throw new PreprocessingStateError(
"preprocessing_resume_mismatch",
"Workspace preprocessing resume no longer matches the pinned revision",
);
}
return existing;
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "ENOENT") {
if (error instanceof PreprocessingStateError) throw error;
throw error;
}
if (options.runId) {
throw new PreprocessingStateError(
"preprocessing_resume_mismatch",
"Workspace preprocessing run is unavailable",
);
}
}
const job: PreprocessingJobState = {
schemaVersion: 2,
runId,
operation: options.operation,
workspaceId: this.options.workspaceId,
workspaceRevision: options.workspaceRevision,
descriptorBlob: options.descriptorBlob,
catalogBlob: options.catalogBlob,
configDigest: options.configDigest,
bindingDigest: options.bindingDigest,
embeddingId: options.embeddingId,
embeddingDimensions: options.embeddingDimensions,
completedStages: [],
childRuns: {},
status: "active",
};
writeAtomicFile(path, `${JSON.stringify(job)}
`, 0o600);
return job;
}
readJob(runId: string): PreprocessingJobState {
return decodeJob(JSON.parse(readTrustedFile(this.jobPath(runId))));
}
writeJob(job: PreprocessingJobState): PreprocessingJobState {
this.ensureLayout();
writeAtomicFile(this.jobPath(job.runId), `${JSON.stringify(job)}
`, 0o600);
return job;
}
writeFkCandidates(runId: string, yaml: string): { path: string; digest: string } {
this.ensureLayout();
const path = this.fkCandidatesPath(runId);
writeAtomicFile(path, yaml, 0o600);
return { path, digest: sha256(yaml) };
}
readFkCandidates(runId: string): { path: string; yaml: string; digest: string } | undefined {
const path = this.fkCandidatesPath(runId);
try {
const yaml = readTrustedFile(path);
return { path, yaml, digest: sha256(yaml) };
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined;
throw error;
}
}
writeFkReview(runId: string, review: FkReviewRecord): { path: string; digest: string } {
this.ensureLayout();
const path = this.fkReviewPath(runId);
const json = `${JSON.stringify(review)}
`;
writeAtomicFile(path, json, 0o600);
return { path, digest: sha256(json) };
}
readFkReview(runId: string): FkReviewRecord | undefined {
try {
return decodeReview(JSON.parse(readTrustedFile(this.fkReviewPath(runId))));
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined;
throw error;
}
}
async assertSessionInventoryCompatible(
workspaceRevision: string,
sessions: readonly SessionInventoryRow[],
): Promise<void> {
const conflicting = sessions.find((session) => (
session.workspaceRevision
&& session.workspaceRevision !== workspaceRevision
&& session.status !== "finalized"
&& !session.archived
));
if (conflicting) {
throw new PreprocessingStateError(
"preprocessing_conflict",
"A resumable session is pinned to a different workspace revision",
);
}
}
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");
}