feat: P2 operator, preprocessing state/service, and runtime config lease
This commit is contained in:
@@ -0,0 +1,383 @@
|
||||
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;
|
||||
runId?: string;
|
||||
}
|
||||
|
||||
export interface PreprocessingJobState {
|
||||
schemaVersion: 1;
|
||||
runId: string;
|
||||
operation: string;
|
||||
workspaceId: string;
|
||||
workspaceRevision: string;
|
||||
descriptorBlob: string;
|
||||
catalogBlob: string;
|
||||
configDigest: string;
|
||||
bindingDigest: string;
|
||||
completedStages: string[];
|
||||
childRuns: Record<string, string>;
|
||||
status: "active" | "succeeded" | "blocked" | "failed";
|
||||
candidateDigest?: string;
|
||||
reviewDigest?: string;
|
||||
}
|
||||
|
||||
export interface FkReviewRecord {
|
||||
reviewedCandidatesDigest: string;
|
||||
annotationsDigest: string;
|
||||
workspaceRevision: 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 !== 1
|
||||
|| 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"
|
||||
|| !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"
|
||||
) 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
|
||||
) {
|
||||
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;
|
||||
if (options.runId) throw error;
|
||||
}
|
||||
}
|
||||
const job: PreprocessingJobState = {
|
||||
schemaVersion: 1,
|
||||
runId,
|
||||
operation: options.operation,
|
||||
workspaceId: this.options.workspaceId,
|
||||
workspaceRevision: options.workspaceRevision,
|
||||
descriptorBlob: options.descriptorBlob,
|
||||
catalogBlob: options.catalogBlob,
|
||||
configDigest: options.configDigest,
|
||||
bindingDigest: options.bindingDigest,
|
||||
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");
|
||||
}
|
||||
Reference in New Issue
Block a user