import { createHash, randomBytes } from "node:crypto"; import { lstatSync } from "node:fs"; import { mkdir, open as openFile, readFile, readdir, rename, rm, writeFile } from "node:fs/promises"; import { join, dirname, isAbsolute, relative } from "node:path"; import { spawn } from "node:child_process"; import { WorkspaceFsAtV1, type OwnedWorkspaceFsAtRegularFile, type OwnedWorkspaceFsAtDirectory, } from "./workspace-fs-at.js"; import type { RuntimeConfigLease } from "./runtime-config-lease.js"; import { BorrowedVerifiedWorkspaceLockRootLease, VerifiedWorkspaceLockRootLease, type CanonicalWorkspaceId, type WorkspaceLockRootIdentityV1, Revision40, workspaceRootInternals, makeBorrowedRootLease, } from "./workspace-lock-root-lease.js"; export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; } export interface PreprocessingRunStateV1 { readonly schemaVersion: 1; readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly operation: string; readonly phase: string; readonly artifacts: readonly ArtifactIdentity[]; readonly createdAt: string; readonly updatedAt: string; } export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly operation: string; } export interface ResumeRunInput { readonly runId: string; readonly workspaceId: string; readonly revision: Revision40; readonly operation: string; } /** A transition is deliberately closed: identity and artifacts are never caller-writable. */ export interface RunTransition { readonly phase: string; } export interface FkReviewInput { readonly candidate: ArtifactIdentity; readonly reviewSha256: string; readonly annotationSha256?: string; } export interface FkReviewRecordV1 extends FkReviewInput { readonly runId: string; readonly recordedAt: string; } const digest = (x: Uint8Array | string) => createHash("sha256").update(x).digest("hex"); const INTERNAL_STATE = Symbol("preprocessing-state-internal"); const RUN_ID = /^[0-9a-f]{32}$/; const SHA256 = /^(?:sha256:)?[0-9a-f]{64}$/; const REVISION = /^[0-9a-f]{40}$/; const WORKSPACE = /^[a-z][a-z0-9-]{2,62}$/; const OPERATIONS = new Set(["dwh", "schema", "evidence"]); const PHASES = ["created", "introspected", "fk_suggested", "fk_reviewed", "schema_indexed", "evidence_indexed", "lsh_built", "terminal"] as const; const MAX_FILE_BYTES = 1 << 20; const MAX_AGGREGATE_BYTES = 64 << 20; const MAX_ENTRIES = 4096; function fail(msg = "preprocessing_conflict"): Error { const e = new Error(msg); e.name = "PreprocessingConflictError"; return e; } type WorkspaceRootLock = { assertPath(): void; close(): void; flock(kind: "shared" | "exclusive", wait: "blocking" | "nonblocking"): void; spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv): Promise; }; function id(v: string): void { if (typeof v !== "string" || !RUN_ID.test(v)) throw fail("invalid run id"); } function checkSha(v: string): void { if (typeof v !== "string" || !SHA256.test(v)) throw fail("invalid digest"); } function strictObject(value: unknown, keys: readonly string[]): value is Record { if (!value || typeof value !== "object" || Array.isArray(value)) return false; const got = Object.keys(value as object).sort(); return got.length === keys.length && got.every((key, i) => key === [...keys].sort()[i]); } function regular0600(path: string): void { const st = lstatSync(path); if (!st.isFile() || (st.mode & 0o777) !== 0o600 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink !== 1) throw fail(); } async function fsyncParent(path: string): Promise { const p = await openFile(dirname(path), "r"); try { await p.sync(); } finally { await p.close(); } } async function durableJson(path: string, value: unknown): Promise { const tmp = `${path}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`; try { await writeFile(tmp, `${JSON.stringify(value)}\n`, { mode: 0o600, flag: "wx" }); regular0600(tmp); const handle = await openFile(tmp, "r"); try { await handle.sync(); } finally { await handle.close(); } await rename(tmp, path); await fsyncParent(path); regular0600(path); } catch (error) { await rm(tmp, { force: true }).catch(() => undefined); throw error; } } export class PreprocessingStateStore { private readonly rootLease: { assertLive(): void; anchoredPath(): string }; constructor(rootLease: VerifiedWorkspaceLockRootLease | string) { if (typeof rootLease === "string") { let st: ReturnType; try { st = lstatSync(rootLease); } catch { throw fail(); } if (!st.isDirectory() || st.isSymbolicLink()) throw fail(); this.rootLease = { assertLive: () => { const current = lstatSync(rootLease); if (!current.isDirectory() || current.isSymbolicLink()) throw fail(); }, anchoredPath: () => rootLease }; } else this.rootLease = { assertLive: () => workspaceRootInternals.assertLive(rootLease), anchoredPath: () => { throw fail(); } }; } private root(): string { this.rootLease.assertLive(); return this.rootLease.anchoredPath(); } private paths(input: { runId: string }) { id(input.runId); const root = this.root(); const base = join(root, "preprocessing"); return { base, jobs: join(base, "jobs"), path: join(base, "jobs", `${input.runId}.json`), candidate: join(base, "fk-candidates", `${input.runId}.yaml`), review: join(base, "fk-reviews", `${input.runId}.json`) }; } private async ensureDirs(p: ReturnType): Promise { for (const path of [p.base, p.jobs, dirname(p.candidate), dirname(p.review)]) { try { await mkdir(path, { mode: 0o700 }); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw fail(); } const st = lstatSync(path); if (!st.isDirectory() || (st.mode & 0o777) !== 0o700 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink < 2) throw fail(); } } private assertFile(path: string, max = MAX_FILE_BYTES): void { regular0600(path); if (lstatSync(path).size > max) throw fail("preprocessing_bounds"); } private async assertBounds(jobs: string): Promise { let entries: string[]; try { entries = await readdir(jobs); } catch { throw fail(); } if (entries.length > MAX_ENTRIES) throw fail("preprocessing_bounds"); let total = 0; for (const name of entries) { if (!RUN_ID.test(name.replace(/\.json$/, "")) || !name.endsWith(".json")) throw fail("preprocessing_bounds"); const path = join(jobs, name); const st = lstatSync(path); if (!st.isFile() || st.nlink !== 1 || (st.mode & 0o777) !== 0o600) throw fail("preprocessing_bounds"); if (st.size > MAX_FILE_BYTES || (total += st.size) > MAX_AGGREGATE_BYTES) throw fail("preprocessing_bounds"); } } async create(input: CreateRunInput): Promise { const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); const p = this.paths({ runId }); if (!WORKSPACE.test(input.workspaceId) || !REVISION.test(input.revision) || !OPERATIONS.has(input.operation)) throw fail(); await this.ensureDirs(p); const now = new Date().toISOString(); const state: PreprocessingRunStateV1 = { schemaVersion: 1, runId, workspaceId: input.workspaceId, revision: input.revision, operation: input.operation, phase: "created", artifacts: [], createdAt: now, updatedAt: now }; try { await durableJson(p.path, state); } catch { throw fail(); } return state; } async loadForResume(input: ResumeRunInput): Promise { const p = this.paths(input); let value: unknown; try { await this.assertBounds(p.jobs); this.assertFile(p.path); value = JSON.parse(await readFile(p.path, "utf8")); } catch { throw fail(); } if (!this.validState(value) || value.runId !== input.runId || value.workspaceId !== input.workspaceId || value.revision !== input.revision || value.operation !== input.operation) throw fail("preprocessing_resume_mismatch"); return value; } async load(input: ResumeRunInput): Promise { return this.loadForResume(input); } async transition(runId: string, transition: RunTransition): Promise { const p = this.paths({ runId }); const current = await this.readRaw(p.path); if (!this.validState(current) || current.runId !== runId || !strictObject(transition, ["phase"]) || typeof transition.phase !== "string") throw fail(); const old = PHASES.indexOf(current.phase as never); const next = PHASES.indexOf(transition.phase as never); if (next < 0 || old < 0 || next < old || current.phase === "terminal") throw fail(); const updated: PreprocessingRunStateV1 = { ...current, phase: transition.phase, updatedAt: new Date().toISOString() }; await durableJson(p.path, updated); return updated; } async writeFkCandidate(runId: string, yaml: Uint8Array): Promise { const p = this.paths({ runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds"); await this.ensureDirs(p); const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity; try { await writeFile(p.candidate, yaml, { flag: "wx", mode: 0o600 }); this.assertFile(p.candidate); const h = await openFile(p.candidate, "r"); try { await h.sync(); } finally { await h.close(); } await fsyncParent(p.candidate); } catch { throw fail(); } return artifact; } async recordFkReview(runId: string, review: FkReviewInput): Promise { const p = this.paths({ runId }); if (!strictObject(review, ["candidate", "reviewSha256", ...(review && "annotationSha256" in review ? ["annotationSha256"] : [])]) || !review?.candidate || !strictObject(review.candidate, ["kind", "digest", "bytes"]) || review.candidate.kind !== "fk-candidate" || !Number.isSafeInteger(review.candidate.bytes) || review.candidate.bytes < 0 || review.candidate.bytes > MAX_FILE_BYTES) throw fail(); checkSha(review.reviewSha256); if (review.annotationSha256 !== undefined) checkSha(review.annotationSha256); const out: FkReviewRecordV1 = { ...review, runId, recordedAt: new Date().toISOString() }; await this.ensureDirs(p); await durableJson(p.review, out); return out; } private async readRaw(path: string): Promise { try { this.assertFile(path); return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } } private validState(value: unknown): value is PreprocessingRunStateV1 { if (!strictObject(value, ["schemaVersion", "runId", "workspaceId", "revision", "operation", "phase", "artifacts", "createdAt", "updatedAt"])) return false; const x = value as Record; if (x.schemaVersion !== 1 || typeof x.runId !== "string" || !RUN_ID.test(x.runId as string) || typeof x.workspaceId !== "string" || !WORKSPACE.test(x.workspaceId as string) || typeof x.revision !== "string" || !REVISION.test(x.revision as string) || typeof x.operation !== "string" || !OPERATIONS.has(x.operation as string) || typeof x.phase !== "string" || !PHASES.includes(x.phase as never) || !Array.isArray(x.artifacts) || typeof x.createdAt !== "string" || typeof x.updatedAt !== "string") return false; if (x.artifacts.length > MAX_ENTRIES) return false; return (x.artifacts as unknown[]).every(a => strictObject(a, ["kind", "digest", "bytes"]) && typeof (a as Record).kind === "string" && typeof (a as Record).digest === "string" && SHA256.test((a as Record).digest as string) && Number.isSafeInteger((a as Record).bytes) && ((a as Record).bytes as number) >= 0 && ((a as Record).bytes as number) <= MAX_FILE_BYTES); } } export interface DwhLockedChildRequest { readonly kind: "dwh_preprocess"; readonly stage: "introspect" | "lsh"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; } export interface SchemaLockedChildRequest { readonly kind: "schema_preprocess"; readonly stage: "fk_suggest" | "fk_check" | "schema_index"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; readonly reviewedArtifact: ArtifactIdentity | null; } export interface EvidenceLockedChildRequest { readonly kind: "evidence_preprocess"; readonly stage: "http_publish"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; } export type WorkspaceLockedChildRequest = DwhLockedChildRequest | SchemaLockedChildRequest | EvidenceLockedChildRequest; export interface WorkspaceLockedChildResult { readonly exitCode: number; readonly stdout: Uint8Array; readonly stderr: Uint8Array; } export class BorrowedWorkspaceSessionReadersExclusiveLockLease { private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} static [INTERNAL_STATE](id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } assertLive(): void { if (!this.live) throw fail(); } invalidate(): void { this.live = false; } } export class WorkspaceWriterLockCapability { private live = true; private settled = false; private readerExclusive = false; private spawnActive = false; private poisoned = false; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1, private readonly root: VerifiedWorkspaceLockRootLease, private readonly writer: WorkspaceRootLock) {} private assertLive(): void { if (!this.live || this.settled || this.poisoned) throw fail(); } assertWriterPath(): void { this.assertLive(); workspaceRootInternals.assertLive(this.root); this.writer.assertPath(); } invalidateForSettlement(): void { this.settled = true; } static [INTERNAL_STATE](id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); } async runUnderSessionReadersExclusive(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise): Promise { this.assertLive(); if (this.readerExclusive || this.spawnActive) throw fail(); this.readerExclusive = true; let lock: WorkspaceRootLock; try { lock = await workspaceRootInternals.acquireReadersExclusive(this.root); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); } const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease[INTERNAL_STATE](this.workspaceId, this.rootIdentity); try { return await action(borrowed); } catch (error) { throw error; } finally { borrowed.invalidate(); let cleanupError: unknown; try { lock.close(); } catch (error) { cleanupError = error; } this.readerExclusive = false; if (cleanupError) { this.poisoned = true; this.live = false; throw fail(); } } } async spawnChild(request: WorkspaceLockedChildRequest): Promise { this.assertLive(); if (!this.readerExclusive || this.spawnActive || !request || request.workspaceId !== this.workspaceId) throw fail(); if (!RUN_ID.test(request.childRunId) || request.rootIdentity.device !== this.rootIdentity.device || request.rootIdentity.inode !== this.rootIdentity.inode || request.rootIdentity.workspaceId !== this.workspaceId) throw fail(); const cfg = request.runtimeConfig; if (cfg.workspaceId !== this.workspaceId || cfg.workspaceRevision !== request.revision || typeof cfg.path !== "string") throw fail(); this.spawnActive = true; try { const configPath = cfg.path; let argv: string[]; switch (request.kind) { case "dwh_preprocess": argv = ["-m", "tht.cli", "preprocess", "dwh", "--steps", request.stage, "--json", "-c", configPath]; break; case "schema_preprocess": argv = ["-m", "tht.cli", "schema", request.stage === "fk_suggest" ? "suggest-fks" : request.stage === "fk_check" ? "check" : "index", "--json", "-c", configPath]; break; case "evidence_preprocess": argv = ["-m", "tht.cli", "preprocess", "evidence", "--json", "-c", configPath]; break; default: throw fail(); } return await workspaceRootInternals.spawn(this.root, this.writer, process.env.THT_PYTHON ?? "python3", argv, { ...process.env, THOTH_WORKSPACE_ID: this.workspaceId, THOTH_WORKSPACE_REVISION: request.revision, THOTH_WORKSPACE_DEVICE: String(this.rootIdentity.device), THOTH_WORKSPACE_INODE: String(this.rootIdentity.inode) }); } finally { this.spawnActive = false; } } async close(): Promise { if (!this.live) return; if (this.readerExclusive || this.spawnActive) throw fail(); this.live = false; let error: unknown; try { this.writer.close(); } catch (e) { error = e; } try { await this.root.close(); } catch (e) { error ??= e; } if (error) throw fail(); } } function makeWriterCapability(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return WorkspaceWriterLockCapability[INTERNAL_STATE](id, identity, root, writer); } export interface BorrowedOrderedWorkspaceWriterLeaseV1 { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; } export interface OrderedWorkspaceCapability { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; } export class OrderedWorkspaceWriterCapabilitySet { private live = true; private constructor(private readonly caps: Map) {} static [INTERNAL_STATE](caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } invalidate(): void { this.live = false; for (const cap of this.caps.values()) cap.invalidateForSettlement(); } get workspaceIds(): readonly CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; } async forWorkspace(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise): Promise { if (!this.live) throw fail(); const cap = this.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: makeBorrowedRootLease(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); } async forEachWorkspace(action: (lease: OrderedWorkspaceCapability) => Promise): Promise { return Promise.all(this.workspaceIds.map(id => this.forWorkspace(id, action))); } } export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise): Promise { const sorted = [...rootLeases].sort((a, b) => a.identity.workspaceId.localeCompare(b.identity.workspaceId)); if (new Set(sorted.map(x => x.identity.workspaceId)).size !== sorted.length) throw fail(); const caps: WorkspaceWriterLockCapability[] = []; let set: OrderedWorkspaceWriterCapabilitySet | undefined; let result: T | undefined; let callbackError: unknown; try { for (const source of sorted) { const root = source.transfer(); let writer: WorkspaceRootLock | undefined; try { writer = await workspaceRootInternals.acquireWriter(root); caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer)); } catch (error) { try { writer?.close(); } catch {} try { await root.close(); } catch {} callbackError = error; break; } } if (!callbackError) { set = OrderedWorkspaceWriterCapabilitySet[INTERNAL_STATE](new Map(caps.map(c => [c.workspaceId, c]))); for (const capability of caps) capability.assertWriterPath(); try { result = await action(set); for (const capability of caps) capability.assertWriterPath(); } catch (error) { callbackError = error; } } } catch (error) { callbackError ??= error; } set?.invalidate(); for (const cap of [...caps].reverse()) { try { await cap.close(); } catch (error) { callbackError ??= error; } } if (callbackError) throw callbackError; return result as T; } export function runUnderWorkspaceWriterLock(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise) { return runUnderOrderedWorkspaceWriterLocks([rootLease], set => set.forWorkspace(rootLease.identity.workspaceId, x => action(x.writerCapability))); } export async function probeWorkspaceWriterLock(rootLease: VerifiedWorkspaceLockRootLease): Promise<"available" | "held"> { try { await runUnderWorkspaceWriterLock(rootLease, async () => undefined); return "available"; } catch { return "held"; } }