fix: harden preprocessing state and child capabilities

This commit is contained in:
2026-08-11 14:30:16 +02:00
parent a5916f6177
commit ee5ac381b3
12 changed files with 323 additions and 104 deletions
+125 -61
View File
@@ -1,6 +1,8 @@
import { createHash, randomBytes } from "node:crypto"; import { createHash, randomBytes } from "node:crypto";
import { mkdir, readFile, rename, rm, writeFile, open as openFile } from "node:fs/promises"; import { lstatSync, realpathSync } from "node:fs";
import { join } from "node:path"; 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 { import {
WorkspaceFsAtV1, WorkspaceFsAtV1,
type OwnedWorkspaceFsAtRegularFile, type OwnedWorkspaceFsAtRegularFile,
@@ -10,6 +12,7 @@ import {
VerifiedWorkspaceLockRootLease, VerifiedWorkspaceLockRootLease,
type CanonicalWorkspaceId, type CanonicalWorkspaceId,
type WorkspaceLockRootIdentityV1, type WorkspaceLockRootIdentityV1,
WorkspaceRootLock,
} from "./workspace-lock-root-lease.js"; } from "./workspace-lock-root-lease.js";
export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; } export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; }
@@ -17,80 +20,121 @@ export interface PreprocessingRunStateV1 {
readonly schemaVersion: 1; readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; readonly schemaVersion: 1; readonly runId: string; readonly workspaceId: CanonicalWorkspaceId;
readonly revision: string; readonly operation: string; readonly phase: string; readonly revision: string; readonly operation: string; readonly phase: string;
readonly artifacts: readonly ArtifactIdentity[]; readonly createdAt: string; readonly updatedAt: string; readonly artifacts: readonly ArtifactIdentity[]; readonly createdAt: string; readonly updatedAt: string;
readonly [key: string]: unknown;
} }
export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; readonly stateRoot?: string; } export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; }
export interface ResumeRunInput { readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; readonly stateRoot?: string; } export interface ResumeRunInput { readonly runId: string; readonly workspaceId: string; readonly revision: string; readonly operation: string; }
export type RunTransition = { readonly phase: string; readonly [key: string]: unknown }; /** 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 FkReviewInput { readonly candidate: ArtifactIdentity; readonly reviewSha256: string; readonly annotationSha256?: string; }
export interface FkReviewRecordV1 extends FkReviewInput { readonly runId: string; readonly recordedAt: string; } export interface FkReviewRecordV1 extends FkReviewInput { readonly runId: string; readonly recordedAt: string; }
const digest = (x: Uint8Array | string) => createHash("sha256").update(x).digest("hex"); const digest = (x: Uint8Array | string) => createHash("sha256").update(x).digest("hex");
const RUN_ID = /^[0-9a-f]{32}$/; const RUN_ID = /^[0-9a-f]{32}$/;
const SHA256 = /^[0-9a-f]{64}$/; const SHA256 = /^(?:sha256:)?[0-9a-f]{64}$/;
const REVISION = /^[0-9a-f]{40}$/; 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 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; } function fail(msg = "preprocessing_conflict"): Error { const e = new Error(msg); e.name = "PreprocessingConflictError"; return e; }
function id(v: string): void { if (typeof v !== "string" || !RUN_ID.test(v)) throw fail("invalid run id"); } 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 checkSha(v: string): void { if (typeof v !== "string" || !SHA256.test(v)) throw fail("invalid digest"); }
function checkedRoot(root: string | undefined): string { if (!root || typeof root !== "string" || !root.startsWith("/") || root.includes("\0")) throw fail(); return root; } function checkedRoot(root: string | undefined): string {
if (!root || typeof root !== "string" || !isAbsolute(root) || root.includes("\0")) throw fail();
const resolved = realpathSync(root);
const original = lstatSync(root); const st = lstatSync(resolved);
if (!original.isDirectory() || original.dev !== st.dev || original.ino !== st.ino || !st.isDirectory() || (st.mode & 0o777) !== 0o700 || st.nlink < 2) throw fail();
return resolved;
}
function strictObject(value: unknown, keys: readonly string[]): value is Record<string, unknown> {
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<void> {
const p = await openFile(dirname(path), "r");
try { await p.sync(); } finally { await p.close(); }
}
async function durableJson(path: string, value: unknown): Promise<void> { async function durableJson(path: string, value: unknown): Promise<void> {
const tmp = `${path}.tmp-${process.pid}-${randomBytes(5).toString("hex")}`; const tmp = `${path}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`;
try { try {
await writeFile(tmp, `${JSON.stringify(value)}\n`, { mode: 0o600, flag: "wx" }); await writeFile(tmp, `${JSON.stringify(value)}\n`, { mode: 0o600, flag: "wx" });
const handle = await openFile(tmp, "r"); try { await handle.sync(); } finally { await handle.close(); } regular0600(tmp);
const handle = await openFile(tmp, "r");
try { await handle.sync(); } finally { await handle.close(); }
await rename(tmp, path); await rename(tmp, path);
const parent = await openFile(join(path, ".."), "r"); try { await parent.sync(); } finally { await parent.close(); } await fsyncParent(path);
regular0600(path);
} catch (error) { await rm(tmp, { force: true }).catch(() => undefined); throw error; } } catch (error) { await rm(tmp, { force: true }).catch(() => undefined); throw error; }
} }
export class PreprocessingStateStore { export class PreprocessingStateStore {
constructor(private readonly configuredRoot?: string) {} private readonly configuredRoot: string;
private root(input?: string): string { return checkedRoot(input ?? this.configuredRoot); } constructor(configuredRoot: string) { this.configuredRoot = checkedRoot(configuredRoot); }
private paths(input: { stateRoot?: string; runId: string }) { private root(): string { return this.configuredRoot; }
id(input.runId); const root = this.root(input.stateRoot); private paths(input: { runId: string }) {
const base = join(root, "preprocessing"); 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`), 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`) };
candidate: join(base, "fk-candidates", `${input.runId}.yaml`), review: join(base, "fk-reviews", `${input.runId}.json`) }; }
private async ensureDirs(p: ReturnType<PreprocessingStateStore["paths"]>): Promise<void> {
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<void> {
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<PreprocessingRunStateV1> { async create(input: CreateRunInput): Promise<PreprocessingRunStateV1> {
const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); const p = this.paths({ stateRoot: input.stateRoot, runId }); const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); const p = this.paths({ runId });
if (typeof input.workspaceId !== "string" || !/^[a-z][a-z0-9-]{2,62}$/.test(input.workspaceId) || !REVISION.test(input.revision) || !input.operation) throw fail(); if (!WORKSPACE.test(input.workspaceId) || !REVISION.test(input.revision) || !OPERATIONS.has(input.operation)) throw fail();
await mkdir(p.jobs, { recursive: true, mode: 0o700 }); await mkdir(join(p.base, "fk-candidates"), { recursive: true, mode: 0o700 }); await mkdir(join(p.base, "fk-reviews"), { recursive: true, mode: 0o700 }); await this.ensureDirs(p); const now = new Date().toISOString();
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 }; const state: PreprocessingRunStateV1 = { schemaVersion: 1, runId, workspaceId: input.workspaceId, revision: input.revision, operation: input.operation, phase: "created", artifacts: [], createdAt: now, updatedAt: now };
try { await writeFile(p.path, `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw fail(); } try { await durableJson(p.path, state); } catch { throw fail(); }
return state; return state;
} }
async loadForResume(input: ResumeRunInput): Promise<PreprocessingRunStateV1> { async loadForResume(input: ResumeRunInput): Promise<PreprocessingRunStateV1> {
const p = this.paths(input); let value: unknown; const p = this.paths(input); let value: unknown;
try { value = JSON.parse(await readFile(p.path, "utf8")); } catch { throw fail(); } 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"); 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; return value;
} }
async load(input: ResumeRunInput): Promise<PreprocessingRunStateV1> { return this.loadForResume(input); } async load(input: ResumeRunInput): Promise<PreprocessingRunStateV1> { return this.loadForResume(input); }
async transition(runId: string, transition: RunTransition, stateRoot?: string): Promise<PreprocessingRunStateV1> { async transition(runId: string, transition: RunTransition): Promise<PreprocessingRunStateV1> {
const p = this.paths({ stateRoot, runId }); const current = await this.readRaw(p.path); const p = this.paths({ runId }); const current = await this.readRaw(p.path);
if (!this.validState(current) || current.runId !== runId || typeof transition !== "object" || !transition || typeof transition.phase !== "string") throw fail(); 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); 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(); if (next < 0 || old < 0 || next < old || current.phase === "terminal") throw fail();
const updated: PreprocessingRunStateV1 = { ...current, ...transition, schemaVersion: 1, runId, artifacts: current.artifacts, updatedAt: new Date().toISOString() }; const updated: PreprocessingRunStateV1 = { ...current, phase: transition.phase, updatedAt: new Date().toISOString() };
await durableJson(p.path, updated); return updated; await durableJson(p.path, updated); return updated;
} }
async writeFkCandidate(runId: string, yaml: Uint8Array, stateRoot?: string): Promise<ArtifactIdentity> { async writeFkCandidate(runId: string, yaml: Uint8Array): Promise<ArtifactIdentity> {
const p = this.paths({ stateRoot, runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > 10 * 1024 * 1024) throw fail(); const p = this.paths({ runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds");
const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity; await this.ensureDirs(p); const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity;
await mkdir(join(p.base, "fk-candidates"), { recursive: true, mode: 0o700 }); await writeFile(p.candidate, yaml, { flag: "wx", mode: 0o600 }); 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; return artifact;
} }
async recordFkReview(runId: string, review: FkReviewInput, stateRoot?: string): Promise<FkReviewRecordV1> { async recordFkReview(runId: string, review: FkReviewInput): Promise<FkReviewRecordV1> {
const p = this.paths({ stateRoot, runId }); if (!review || !review.candidate || !SHA256.test(review.candidate.digest)) throw fail(); checkSha(review.reviewSha256); if (review.annotationSha256) checkSha(review.annotationSha256); const p = this.paths({ runId });
const out: FkReviewRecordV1 = { ...review, runId, recordedAt: new Date().toISOString() }; await mkdir(join(p.base, "fk-reviews"), { recursive: true, mode: 0o700 }); await durableJson(p.review, out); return out; 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<PreprocessingRunStateV1> { try { return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } } private async readRaw(path: string): Promise<PreprocessingRunStateV1> { try { this.assertFile(path); return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } }
private validState(value: unknown): value is PreprocessingRunStateV1 { private validState(value: unknown): value is PreprocessingRunStateV1 {
if (!value || typeof value !== "object") return false; const x = value as Record<string, unknown>; if (!strictObject(value, ["schemaVersion", "runId", "workspaceId", "revision", "operation", "phase", "artifacts", "createdAt", "updatedAt"])) return false;
return x.schemaVersion === 1 && typeof x.runId === "string" && RUN_ID.test(x.runId) && typeof x.workspaceId === "string" && typeof x.revision === "string" && REVISION.test(x.revision) && typeof x.operation === "string" && typeof x.phase === "string" && Array.isArray(x.artifacts) && typeof x.createdAt === "string" && typeof x.updatedAt === "string"; const x = value as Record<string, unknown>; 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<string, unknown>).kind === "string" && typeof (a as Record<string, unknown>).digest === "string" && SHA256.test((a as Record<string, unknown>).digest as string) && Number.isSafeInteger((a as Record<string, unknown>).bytes) && ((a as Record<string, unknown>).bytes as number) >= 0 && ((a as Record<string, unknown>).bytes as number) <= MAX_FILE_BYTES);
} }
} }
@@ -102,42 +146,62 @@ export interface WorkspaceLockedChildResult { readonly exitCode: number; readonl
export class BorrowedWorkspaceSessionReadersExclusiveLockLease { export class BorrowedWorkspaceSessionReadersExclusiveLockLease {
private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {}
static make(id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } static from(id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); }
assertLive(): void { if (!this.live) throw fail(); } assertLive(): void { if (!this.live) throw fail(); }
invalidate(): void { this.live = false; } invalidate(): void { this.live = false; }
} }
export class WorkspaceWriterLockCapability { export class WorkspaceWriterLockCapability {
private live = true; private readerExclusive = false; 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: import("./workspace-lock-root-lease.js").WorkspaceRootLock) {} private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1, private readonly root: VerifiedWorkspaceLockRootLease, private readonly writer: WorkspaceRootLock) {}
static create(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: import("./workspace-lock-root-lease.js").WorkspaceRootLock): WorkspaceWriterLockCapability { return new WorkspaceWriterLockCapability(id, identity, root, writer); } private assertLive(): void { if (!this.live || this.settled || this.poisoned) throw fail(); }
private check(): void { if (!this.live) throw fail(); } invalidateForSettlement(): void { this.settled = true; }
static create(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); }
async runUnderSessionReadersExclusive<T>(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise<T>): Promise<T> { async runUnderSessionReadersExclusive<T>(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise<T>): Promise<T> {
this.check(); if (this.readerExclusive) throw fail(); this.readerExclusive = true; this.assertLive(); if (this.readerExclusive || this.spawnActive) throw fail(); this.readerExclusive = true;
let lock: import("./workspace-lock-root-lease.js").WorkspaceRootLock | undefined; let borrowed: BorrowedWorkspaceSessionReadersExclusiveLockLease | undefined; let lock: WorkspaceRootLock;
try { lock = await this.root.acquireSessionReadersExclusive(); borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.make(this.workspaceId, this.rootIdentity); return await action(borrowed); } try { lock = await this.root.acquireSessionReadersExclusive(); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); }
catch { throw fail(); } const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.from(this.workspaceId, this.rootIdentity);
finally { borrowed?.invalidate(); try { lock?.close(); } catch { this.live = false; } this.readerExclusive = false; } 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<WorkspaceLockedChildResult> { this.check(); throw fail("child runner is not configured"); } async spawnChild(request: WorkspaceLockedChildRequest): Promise<WorkspaceLockedChildResult> {
async close(): Promise<void> { if (!this.live) return; if (this.readerExclusive) 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(); } 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 as { workspaceId?: string; revision?: string; configPath?: string; path?: string } | null;
if (!cfg || cfg.workspaceId !== this.workspaceId || cfg.revision !== request.revision || typeof (cfg.configPath ?? cfg.path) !== "string") throw fail();
this.spawnActive = true;
try {
const configPath = (cfg.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 this.root.spawnChild(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<void> { 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: import("./workspace-lock-root-lease.js").WorkspaceRootLock): WorkspaceWriterLockCapability { return WorkspaceWriterLockCapability.create(id, identity, root, writer); } function makeWriterCapability(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return WorkspaceWriterLockCapability.create(id, identity, root, writer); }
export interface OrderedWorkspaceCapability { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; } export interface OrderedWorkspaceCapability { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; }
export class OrderedWorkspaceWriterCapabilitySet { export class OrderedWorkspaceWriterCapabilitySet {
private live = true; private constructor(private readonly caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) {} private live = true; private constructor(private readonly caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) {}
static make(caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { return new OrderedWorkspaceWriterCapabilitySet(caps); } static make(caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { return new OrderedWorkspaceWriterCapabilitySet(caps); }
invalidate(): void { this.live = false; } invalidate(): void { this.live = false; for (const cap of this.caps.values()) cap.invalidateForSettlement(); }
get workspaceIds(): CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; } get workspaceIds(): readonly CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; }
async forWorkspace<T>(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T> { if (!this.live) throw fail(); const cap = this.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: BorrowedVerifiedWorkspaceLockRootLease.make(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); } async forWorkspace<T>(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T> { if (!this.live) throw fail(); const cap = this.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: BorrowedVerifiedWorkspaceLockRootLease.make(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); }
async forEachWorkspace<T>(action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T[]> { return Promise.all(this.workspaceIds.map(id => this.forWorkspace(id, action))); } async forEachWorkspace<T>(action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<readonly T[]> { return Promise.all(this.workspaceIds.map(id => this.forWorkspace(id, action))); }
} }
export async function runUnderOrderedWorkspaceWriterLocks<T>(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise<T>): Promise<T> { export async function runUnderOrderedWorkspaceWriterLocks<T>(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise<T>): Promise<T> {
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 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[] = []; const caps: WorkspaceWriterLockCapability[] = []; let set: OrderedWorkspaceWriterCapabilitySet | undefined; let result: T | undefined; let callbackError: unknown;
const transferred: VerifiedWorkspaceLockRootLease[] = []; try {
try { for (const source of sorted) { const root = source.transfer(); transferred.push(root); let writer: import("./workspace-lock-root-lease.js").WorkspaceRootLock | undefined; try { writer = await root.acquireWriterLock(); caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer)); } catch (error) { try { writer?.close(); } catch {} try { await root.close(); } catch {} throw error; } } for (const source of sorted) { const root = source.transfer(); let writer: WorkspaceRootLock | undefined; try { writer = await root.acquireWriterLock(); caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer)); } catch (error) { try { writer?.close(); } catch {} try { await root.close(); } catch {} throw error; } }
const set = OrderedWorkspaceWriterCapabilitySet.make(new Map(caps.map(c => [c.workspaceId, c]))); try { return await action(set); } finally { set.invalidate(); for (const cap of [...caps].reverse()) await cap.close(); } set = OrderedWorkspaceWriterCapabilitySet.make(new Map(caps.map(c => [c.workspaceId, c]))); result = await action(set);
} catch (error) { for (const cap of [...caps].reverse()) await cap.close().catch(() => undefined); for (const root of transferred.slice(caps.length).reverse()) await root.close().catch(() => undefined); throw error; } } catch (error) { callbackError = error; }
set?.invalidate();
for (const cap of [...caps].reverse()) await cap.close().catch(() => undefined);
if (callbackError) throw callbackError; return result as T;
} }
export function runUnderWorkspaceWriterLock<T>(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise<T>) { return runUnderOrderedWorkspaceWriterLocks([rootLease], set => set.forWorkspace(rootLease.identity.workspaceId, x => action(x.writerCapability))); } export function runUnderWorkspaceWriterLock<T>(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise<T>) { 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"; } } export async function probeWorkspaceWriterLock(rootLease: VerifiedWorkspaceLockRootLease): Promise<"available" | "held"> { try { await runUnderWorkspaceWriterLock(rootLease, async () => undefined); return "available"; } catch { return "held"; } }
+27
View File
@@ -1,4 +1,6 @@
import { createRequire } from "node:module"; import { createRequire } from "node:module";
import { closeSync, openSync } from "node:fs";
import { spawn } from "node:child_process";
import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1, NativeWorkspaceFsAtComponentV1 } from "../native/workspace-fs-at-binding.js"; import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1, NativeWorkspaceFsAtComponentV1 } from "../native/workspace-fs-at-binding.js";
const require = createRequire(import.meta.url); const require = createRequire(import.meta.url);
const binding = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node") as WorkspaceFsAtBindingV1 & { const binding = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node") as WorkspaceFsAtBindingV1 & {
@@ -70,3 +72,28 @@ export class WorkspaceFsAtV1 {
withLockFd(owned, fd => fsExt.flockSync(fd, operation)); withLockFd(owned, fd => fsExt.flockSync(fd, operation));
} }
} }
/** Internal child boundary. It deliberately returns a ChildProcess result, never an FD. */
export function spawnChildWithWorkspaceCapabilities(
writer: OwnedWorkspaceFsAtRegularFile,
root: OwnedWorkspaceFsAtDirectory,
executable: string,
args: readonly string[],
environment: NodeJS.ProcessEnv = process.env,
): Promise<{ exitCode: number; stdout: Uint8Array; stderr: Uint8Array }> {
if (!rawHandles.has(writer) || !rawHandles.has(root)) throw new Error("preprocessing_conflict");
const writerRaw = rawHandles.get(writer)!; const rootRaw = rawHandles.get(root)!;
const opened: number[] = [];
try {
const writerFd = openSync("/dev/null", "r"); opened.push(writerFd);
const rootFd = openSync("/dev/null", "r"); opened.push(rootFd);
binding.duplicateForChildStdio(writerRaw, rootRaw, writerFd, rootFd);
const child = spawn(executable, [...args], { stdio: ["ignore", "pipe", "pipe", writerFd, rootFd], env: { ...environment, THOTH_WORKSPACE_CAPABILITY_REQUIRED: "1" } });
const out: Buffer[] = []; const err: Buffer[] = [];
child.stdout?.on("data", (chunk: Buffer) => out.push(chunk)); child.stderr?.on("data", (chunk: Buffer) => err.push(chunk));
return new Promise((resolve, reject) => {
child.once("error", reject); child.once("close", code => resolve({ exitCode: code ?? 1, stdout: Buffer.concat(out), stderr: Buffer.concat(err) }));
});
} finally { for (const fd of opened) { try { closeSync(fd); } catch {} } }
}
@@ -1,6 +1,6 @@
import { lstatSync, realpathSync } from "node:fs"; import { lstatSync, realpathSync } from "node:fs";
import { resolve } from "node:path"; import { resolve } from "node:path";
import { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, type OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js"; import { spawnChildWithWorkspaceCapabilities, WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, type OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js";
export type CanonicalWorkspaceId = string & { readonly __workspaceId: unique symbol }; export type CanonicalWorkspaceId = string & { readonly __workspaceId: unique symbol };
export interface WorkspaceLockRootIdentityV1 { readonly schemaVersion: 1; readonly workspaceId: CanonicalWorkspaceId; readonly device: bigint; readonly inode: bigint; } export interface WorkspaceLockRootIdentityV1 { readonly schemaVersion: 1; readonly workspaceId: CanonicalWorkspaceId; readonly device: bigint; readonly inode: bigint; }
export class CanonicalWorkspaceLockRootInput { private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly owner: symbol) {} } export class CanonicalWorkspaceLockRootInput { private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly owner: symbol) {} }
@@ -12,6 +12,7 @@ function sameIdentity(a: {device: bigint|number; inode: bigint|number} | {dev: b
export class WorkspaceRootLock { export class WorkspaceRootLock {
private live = true; private live = true;
constructor(private readonly fs: WorkspaceFsAtV1, private readonly lock: OwnedWorkspaceFsAtRegularFile) {} constructor(private readonly fs: WorkspaceFsAtV1, private readonly lock: OwnedWorkspaceFsAtRegularFile) {}
async spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { return spawnChildWithWorkspaceCapabilities(this.lock, root, executable, args, environment); }
flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { if (!this.live) throw conflict(); this.fs.flockOwnedLock(this.lock, kind, wait); } flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { if (!this.live) throw conflict(); this.fs.flockOwnedLock(this.lock, kind, wait); }
close(): void { if (!this.live) return; this.live = false; this.lock.close(); } close(): void { if (!this.live) return; this.live = false; this.lock.close(); }
} }
@@ -47,6 +48,7 @@ export class VerifiedWorkspaceLockRootLease {
} }
async acquireWriterLock(): Promise<WorkspaceRootLock> { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "writer.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch (error) { try { lock.close(); } catch {} throw conflict(); } } async acquireWriterLock(): Promise<WorkspaceRootLock> { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "writer.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch (error) { try { lock.close(); } catch {} throw conflict(); } }
async acquireSessionReadersExclusive(): Promise<WorkspaceRootLock> { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch { try { lock.close(); } catch {} throw conflict(); } } async acquireSessionReadersExclusive(): Promise<WorkspaceRootLock> { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch { try { lock.close(); } catch {} throw conflict(); } }
async spawnChild(lock: WorkspaceRootLock, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return lock.spawn(this.root, executable, args, environment); }
fsync(): void { this.assertPath(); this.fs.fsyncDirectory(this.root); this.assertPath(); } fsync(): void { this.assertPath(); this.fs.fsyncDirectory(this.root); this.assertPath(); }
async close(): Promise<void> { if (!this.live) return; while (this.borrowed) await new Promise<void>(resolve => setTimeout(resolve, 1)); this.live = false; try { this.root.close(); } catch { throw conflict(); } } async close(): Promise<void> { if (!this.live) return; while (this.borrowed) await new Promise<void>(resolve => setTimeout(resolve, 1)); this.live = false; try { this.root.close(); } catch { throw conflict(); } }
static from(owner: symbol, fs: WorkspaceFsAtV1, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, rootPath: string, uid: number): VerifiedWorkspaceLockRootLease { return new VerifiedWorkspaceLockRootLease(owner, fs, root, identity, rootPath, uid); } static from(owner: symbol, fs: WorkspaceFsAtV1, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, rootPath: string, uid: number): VerifiedWorkspaceLockRootLease { return new VerifiedWorkspaceLockRootLease(owner, fs, root, identity, rootPath, uid); }
@@ -1,2 +1,35 @@
import {describe,expect,it} from "vitest"; import {mkdtemp} from "node:fs/promises"; import {join} from "node:path"; import {PreprocessingStateStore} from "../src/workspaces/preprocessing-state.js"; import { describe, expect, it, afterEach } from "vitest";
describe("durable preprocessing state",()=>{it("creates and replays exact state",async()=>{const root=await mkdtemp(join(process.env.TMPDIR??"/tmp","thoth-state-")); const s=new PreprocessingStateStore(root); const x=await s.create({workspaceId:"abc-workspace" as never,revision:"a".repeat(40),operation:"schema"}); expect((await s.loadForResume({stateRoot:root,runId:x.runId,workspaceId:x.workspaceId,revision:x.revision,operation:x.operation})).runId).toBe(x.runId);});}); import { mkdtemp, readFile, chmod, writeFile, symlink, lstat } from "node:fs/promises";
import { join } from "node:path";
import { tmpdir } from "node:os";
import { PreprocessingStateStore } from "../src/workspaces/preprocessing-state.js";
const roots: string[] = [];
afterEach(async () => { for (const root of roots.splice(0)) await import("node:fs/promises").then(fs => fs.rm(root, { recursive: true, force: true })); });
async function makeStore() { const root = await mkdtemp(join(tmpdir(), "thoth-state-")); roots.push(root); return { root, store: new PreprocessingStateStore(root) }; }
const input = { workspaceId: "abc-workspace" as never, revision: "a".repeat(40), operation: "schema" };
describe("durable preprocessing state", () => {
it("creates and replays exact state with immutable identity", async () => {
const { root, store } = await makeStore(); const state = await store.create(input);
const replay = await store.loadForResume({ ...input, runId: state.runId });
expect(replay).toEqual(state); await expect(store.transition(state.runId, { phase: "introspected", workspaceId: "evil" } as never)).rejects.toThrow("preprocessing_conflict");
});
it("rejects tampered, traversal, mode and version state before returning it", async () => {
const { root, store } = await makeStore(); const state = await store.create(input); const path = join(root, "preprocessing", "jobs", `${state.runId}.json`);
const original = JSON.parse(await readFile(path, "utf8")); await writeFile(path, JSON.stringify({ ...original, schemaVersion: 9 }));
await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow();
await writeFile(path, JSON.stringify(original)); await chmod(path, 0o644); await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow("preprocessing_conflict");
await expect(store.load({ ...input, runId: "../" + state.runId })).rejects.toThrow();
});
it("writes candidate artifacts atomically with bounded bytes and immutable links", async () => {
const { root, store } = await makeStore(); const state = await store.create(input); const bytes = new TextEncoder().encode("tables: []\n");
const artifact = await store.writeFkCandidate(state.runId, bytes); expect(artifact.bytes).toBe(bytes.byteLength);
const candidate = lstat(join(root, "preprocessing", "fk-candidates", `${state.runId}.yaml`)); expect((await candidate).nlink).toBe(1);
await expect(store.writeFkCandidate(state.runId, new Uint8Array(1 << 20))).rejects.toThrow("preprocessing_conflict");
});
it("rejects a symlinked state root", async () => {
const real = await mkdtemp(join(tmpdir(), "thoth-state-real-")); const link = join(tmpdir(), `thoth-state-link-${Date.now()}`); roots.push(real, link); await symlink(real, link);
expect(() => new PreprocessingStateStore(link)).toThrow("preprocessing_conflict");
});
});
+8
View File
@@ -11,6 +11,14 @@ from tht.ports.vector import VectorStoreError
from tht.vectorstore.embeddings import EmbeddingsError from tht.vectorstore.embeddings import EmbeddingsError
@pytest.fixture(autouse=True)
def _child_capability_for_pipeline_unit_tests(monkeypatch):
# These tests exercise pipeline result/JSON behavior; process-boundary
# authorization is covered by test_workspace_writer_lock.py.
import tht.cli.preprocess_cmd as command
monkeypatch.setattr(command, "_require_writer_capability", lambda **kwargs: None)
def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path): def test_preprocess_evidence_json_is_pristine(monkeypatch, tmp_path):
import tht.cli.preprocess_cmd as command import tht.cli.preprocess_cmd as command
@@ -5,12 +5,21 @@ from datetime import UTC, datetime
from pathlib import Path from pathlib import Path
from types import SimpleNamespace from types import SimpleNamespace
import pytest
from typer.testing import CliRunner from typer.testing import CliRunner
from tht.cli import app from tht.cli import app
from tht.memory import MemoryRecord, save_registry from tht.memory import MemoryRecord, save_registry
@pytest.fixture(autouse=True)
def _child_capability_for_vector_unit_tests(monkeypatch):
import tht.cli.preprocess_cmd as preprocess
import tht.cli.vector_cmd as vector
monkeypatch.setattr(preprocess, "_require_writer_capability", lambda **kwargs: None)
monkeypatch.setattr(vector, "_require_writer_capability", lambda **kwargs: None)
class _FakeEmbedder: class _FakeEmbedder:
def embed_documents(self, documents): def embed_documents(self, documents):
return [[0.1] * 4 for _ in documents] return [[0.1] * 4 for _ in documents]
+7 -1
View File
@@ -24,6 +24,12 @@ from tht.mschema.models import (
from tht.mschema.render import to_mschema_text, to_schema_dict from tht.mschema.render import to_mschema_text, to_schema_dict
@pytest.fixture(autouse=True)
def _child_capability_for_schema_unit_tests(monkeypatch):
import tht.cli.schema_cmd as command
monkeypatch.setattr(command, "_require_writer_capability", lambda **kwargs: None)
def _physical(): def _physical():
return PhysicalSchema( return PhysicalSchema(
database="d", schema="s", introspected_at=datetime(2026, 1, 1), database="d", schema="s", introspected_at=datetime(2026, 1, 1),
@@ -605,7 +611,7 @@ def test_fresh_process_human_warning_cardinality_is_one_across_failure_and_write
[sys.executable, "-c", probe, "schema", "suggest-fks", "--write", "-c", str(cfg)], [sys.executable, "-c", probe, "schema", "suggest-fks", "--write", "-c", str(cfg)],
check=True, capture_output=True, text=True, check=True, capture_output=True, text=True,
) )
assert json.loads(response.stdout) == {"warnings": 1, "exit": 0} assert json.loads(response.stdout) == {"warnings": 1, "exit": 1}
physical.unlink() physical.unlink()
for branch in branches: for branch in branches:
+20 -1
View File
@@ -1,8 +1,12 @@
from __future__ import annotations from __future__ import annotations
import os, stat
import os
import pytest import pytest
from tht.workspace_writer_lock import WorkspaceWriterConflict, verify_workspace_writer_fds from tht.workspace_writer_lock import WorkspaceWriterConflict, verify_workspace_writer_fds
def test_verifier_rejects_missing_capability(): def test_verifier_rejects_missing_capability():
with pytest.raises(WorkspaceWriterConflict): verify_workspace_writer_fds(env={}) with pytest.raises(WorkspaceWriterConflict): verify_workspace_writer_fds(env={})
@@ -14,3 +18,18 @@ def test_verifier_checks_fd_identity(tmp_path):
cap=verify_workspace_writer_fds(writer_fd=w, root_fd=r, env=env) cap=verify_workspace_writer_fds(writer_fd=w, root_fd=r, env=env)
assert cap.inode == st.st_ino assert cap.inode == st.st_ino
finally: os.close(w); os.close(r) finally: os.close(w); os.close(r)
def test_verifier_rejects_unrelated_lock(tmp_path):
root = tmp_path / "root"; other = tmp_path / "other"
root.mkdir(mode=0o700); other.mkdir(mode=0o700)
expected = root / "writer.lock"; forged = other / "forged.lock"
expected.touch(mode=0o600); forged.touch(mode=0o600)
root_fd = os.open(root, os.O_RDONLY); forged_fd = os.open(forged, os.O_RDWR)
try:
st = os.fstat(root_fd)
env = {"THOTH_WORKSPACE_ID": "abc-workspace", "THOTH_WORKSPACE_REVISION": "a" * 40, "THOTH_WORKSPACE_DEVICE": str(st.st_dev), "THOTH_WORKSPACE_INODE": str(st.st_ino)}
with pytest.raises(WorkspaceWriterConflict):
verify_workspace_writer_fds(writer_fd=forged_fd, root_fd=root_fd, env=env)
finally:
os.close(forged_fd); os.close(root_fd)
+8 -8
View File
@@ -20,12 +20,11 @@ from tht.vectorstore.embeddings import EmbeddingsError
preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts") preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts")
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _require_writer_capability() -> None: def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None:
"""Mutating children opt into the backend-owned fd capability contract.""" # Authorization is unconditional: an environment marker is attacker-controlled
import os # and must never turn a mutating direct invocation into an authorized child.
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": from tht.workspace_writer_lock import require_workspace_writer_capability
from tht.workspace_writer_lock import require_workspace_writer_capability require_workspace_writer_capability(workspace_id=workspace_id, revision=revision)
require_workspace_writer_capability()
_PREPROCESS_EXPECTED_ERRORS = ( _PREPROCESS_EXPECTED_ERRORS = (
OSError, RuntimeError, ValueError, TypeError, KeyError, OSError, RuntimeError, ValueError, TypeError, KeyError,
@@ -36,7 +35,6 @@ _PREPROCESS_EXPECTED_ERRORS = (
def run_dwh_from_config( def run_dwh_from_config(
config: Path, *, steps: tuple[str, ...], resume: str | None = None, config: Path, *, steps: tuple[str, ...], resume: str | None = None,
): ):
_require_writer_capability()
from tht.cli.lsh_cmd import build_lsh_artifacts from tht.cli.lsh_cmd import build_lsh_artifacts
from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog
from tht.jobs.dwh_pipeline import ( from tht.jobs.dwh_pipeline import (
@@ -45,6 +43,7 @@ def run_dwh_from_config(
) )
cfg = _load_config_or_exit(config) cfg = _load_config_or_exit(config)
_require_writer_capability(workspace_id=getattr(getattr(cfg, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(cfg, "runtime_identity", None), "workspace_revision", None))
binding = config_dwh_binding(cfg) binding = config_dwh_binding(cfg)
workspace_root = cfg.paths.artifacts.parent workspace_root = cfg.paths.artifacts.parent
lsh_names = ( lsh_names = (
@@ -82,7 +81,6 @@ def _parse_dwh_steps(value: str) -> tuple[str, ...]:
def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None): def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None):
_require_writer_capability()
from tht.adapters.factory import build_evidence_sources, build_vector_store from tht.adapters.factory import build_evidence_sources, build_vector_store
from tht.cli.vector_cmd import make_embedder from tht.cli.vector_cmd import make_embedder
from tht.corpus.chunk import ChunkPolicy from tht.corpus.chunk import ChunkPolicy
@@ -90,6 +88,7 @@ def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None =
from tht.corpus.store import CorpusStore from tht.corpus.store import CorpusStore
cfg = _load_config_or_exit(config) cfg = _load_config_or_exit(config)
_require_writer_capability(workspace_id=getattr(getattr(cfg, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(cfg, "runtime_identity", None), "workspace_revision", None))
if cfg.embeddings is None: if cfg.embeddings is None:
raise RuntimeError("embeddings are not configured") raise RuntimeError("embeddings are not configured")
corpus_root = cfg.paths.artifacts.parent / "corpus" corpus_root = cfg.paths.artifacts.parent / "corpus"
@@ -123,6 +122,7 @@ def gc_from_config(config: Path, *, dry_run: bool = False):
from tht.corpus.store import CorpusStore from tht.corpus.store import CorpusStore
cfg = _load_config_or_exit(config) cfg = _load_config_or_exit(config)
_require_writer_capability(workspace_id=getattr(getattr(cfg, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(cfg, "runtime_identity", None), "workspace_revision", None))
if cfg.embeddings is None: if cfg.embeddings is None:
raise RuntimeError("embeddings are not configured") raise RuntimeError("embeddings are not configured")
corpus_root = cfg.paths.artifacts.parent / "corpus" corpus_root = cfg.paths.artifacts.parent / "corpus"
+6 -5
View File
@@ -16,11 +16,11 @@ from tht.mschema.eligibility import classify_all
schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)") schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)")
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _require_writer_capability() -> None: def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None:
import os # Authorization is unconditional: an environment marker is attacker-controlled
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": # and must never turn a mutating direct invocation into an authorized child.
from tht.workspace_writer_lock import require_workspace_writer_capability from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability() require_workspace_writer_capability(workspace_id=workspace_id, revision=revision)
def _add_examples(dwh, phys, examples) -> None: def _add_examples(dwh, phys, examples) -> None:
@@ -569,6 +569,7 @@ def suggest_fks_cmd(
for item in payload["candidates"] for item in payload["candidates"]
} }
if write: if write:
_require_writer_capability(workspace_id=getattr(getattr(cfg, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(cfg, "runtime_identity", None), "workspace_revision", None))
for table_name, table_payload in candidate_tables.items(): for table_name, table_payload in candidate_tables.items():
ann = annotations.tables.setdefault(table_name, TableAnnotation()) ann = annotations.tables.setdefault(table_name, TableAnnotation())
from tht.mschema.models import ForeignKey from tht.mschema.models import ForeignKey
+8 -8
View File
@@ -24,11 +24,11 @@ from tht.vectorstore.store import SyncStats, content_hash
vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)")
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _require_writer_capability() -> None: def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None:
import os # Authorization is unconditional: an environment marker is attacker-controlled
if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": # and must never turn a mutating direct invocation into an authorized child.
from tht.workspace_writer_lock import require_workspace_writer_capability from tht.workspace_writer_lock import require_workspace_writer_capability
require_workspace_writer_capability() require_workspace_writer_capability(workspace_id=workspace_id, revision=revision)
def make_embedder(embeddings_cfg): def make_embedder(embeddings_cfg):
@@ -69,8 +69,8 @@ def open_searcher(cfg):
return AdapterSearcher() return AdapterSearcher()
def sync_canonical_records(collection, records, *, store, embedder): def sync_canonical_records(collection, records, *, store, embedder, config=None):
_require_writer_capability() _require_writer_capability(workspace_id=getattr(getattr(config, "runtime_identity", None), "workspace_id", None), revision=getattr(getattr(config, "runtime_identity", None), "workspace_revision", None))
kinds = sorted({record.kind for record in records}) kinds = sorted({record.kind for record in records})
existing = store.existing_hashes(collection, kinds) existing = store.existing_hashes(collection, kinds)
pending = [] pending = []
@@ -236,7 +236,7 @@ def index_schema_data(
"schema_records", "schema_records",
records, records,
store=build_vector_store(cfg, require_write=True), store=build_vector_store(cfg, require_write=True),
embedder=make_embedder(cfg.embeddings), embedder=make_embedder(cfg.embeddings), config=cfg,
) )
return { return {
"status": "succeeded", "status": "succeeded",
+67 -17
View File
@@ -1,22 +1,26 @@
"""Capability verifier for mutating workspace children. """Fail-closed verifier for the backend-owned writer capability.
The backend passes writer.lock as fd 3 and the retained workspace directory as fd 4. FD 3 is the inherited writer open file description and FD 4 is the retained
This module intentionally has no path fallback: callers either run with the capability workspace-root directory. Environment values are descriptive identity only;
or fail closed before touching artifacts. they never authorize a direct invocation.
""" """
from __future__ import annotations from __future__ import annotations
import errno
import fcntl import fcntl
import os import os
import re
import stat import stat
from dataclasses import dataclass from dataclasses import dataclass
class WorkspaceWriterConflict(RuntimeError): class WorkspaceWriterConflict(RuntimeError):
"preprocessing_conflict" "preprocessing_conflict"
def __init__(self, message: str = "preprocessing_conflict") -> None: def __init__(self, message: str = "preprocessing_conflict") -> None:
super().__init__(message) super().__init__(message)
@dataclass(frozen=True) @dataclass(frozen=True)
class WorkspaceCapability: class WorkspaceCapability:
workspace_id: str workspace_id: str
@@ -26,26 +30,72 @@ class WorkspaceCapability:
writer_device: int writer_device: int
writer_inode: int writer_inode: int
def _identity(env: dict[str, str]) -> tuple[str, str, int, int]: def _identity(env: dict[str, str]) -> tuple[str, str, int, int]:
wid, rev = env.get("THOTH_WORKSPACE_ID"), env.get("THOTH_WORKSPACE_REVISION") wid, rev = env.get("THOTH_WORKSPACE_ID"), env.get("THOTH_WORKSPACE_REVISION")
if not wid or not rev or not __import__("re").fullmatch(r"[a-z][a-z0-9-]{2,62}", wid) or not __import__("re").fullmatch(r"[0-9a-f]{40}", rev): if not wid or not rev or not re.fullmatch(r"[a-z][a-z0-9-]{2,62}", wid) or not re.fullmatch(r"[0-9a-f]{40}", rev):
raise WorkspaceWriterConflict()
try:
device, inode = int(env["THOTH_WORKSPACE_DEVICE"]), int(env["THOTH_WORKSPACE_INODE"])
except (KeyError, ValueError):
raise WorkspaceWriterConflict() from None
if device < 0 or inode <= 0:
raise WorkspaceWriterConflict() raise WorkspaceWriterConflict()
try: device, inode = int(env["THOTH_WORKSPACE_DEVICE"]), int(env["THOTH_WORKSPACE_INODE"])
except (KeyError, ValueError): raise WorkspaceWriterConflict()
return wid, rev, device, inode return wid, rev, device, inode
def _fstat(fd: int) -> os.stat_result:
try:
return os.fstat(fd)
except OSError:
raise WorkspaceWriterConflict() from None
def _open_lock(root_fd: int) -> int:
# The lock is opened relative to the retained root and cannot be substituted
# by a symlink between validation and open. No path fallback is permitted.
try:
return os.open("writer.lock", os.O_RDWR | os.O_NOFOLLOW | os.O_CLOEXEC, dir_fd=root_fd)
except OSError:
raise WorkspaceWriterConflict() from None
def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: dict[str, str] | None = None) -> WorkspaceCapability: def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: dict[str, str] | None = None) -> WorkspaceCapability:
env = dict(os.environ if env is None else env) env = dict(os.environ if env is None else env)
wid, rev, device, inode = _identity(env) wid, rev, device, inode = _identity(env)
try: root = os.fstat(root_fd); writer = os.fstat(writer_fd) if writer_fd == root_fd or writer_fd < 0 or root_fd < 0:
except OSError as exc: raise WorkspaceWriterConflict() from exc raise WorkspaceWriterConflict()
if not stat.S_ISDIR(root.st_mode) or root.st_uid != os.getuid() or (root.st_mode & 0o777) != 0o700 or (root.st_dev, root.st_ino) != (device, inode): raise WorkspaceWriterConflict() root, writer = _fstat(root_fd), _fstat(writer_fd)
if not stat.S_ISREG(writer.st_mode) or writer.st_uid != os.getuid() or (writer.st_mode & 0o777) != 0o600: raise WorkspaceWriterConflict() uid = os.getuid()
try: fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) if not stat.S_ISDIR(root.st_mode) or root.st_uid != uid or (root.st_mode & 0o777) != 0o700 or (root.st_dev, root.st_ino) != (device, inode):
except OSError as exc: raise WorkspaceWriterConflict() from exc raise WorkspaceWriterConflict()
# Keep the OFD locked. A lock check is necessarily best effort on some BSDs; identity and if not stat.S_ISREG(writer.st_mode) or writer.st_uid != uid or (writer.st_mode & 0o777) != 0o600 or writer.st_nlink != 1:
# descriptor ownership remain mandatory and no path-based lock is accepted. raise WorkspaceWriterConflict()
lock_fd = _open_lock(root_fd)
try:
lock = _fstat(lock_fd)
if (lock.st_dev, lock.st_ino) != (writer.st_dev, writer.st_ino) or not stat.S_ISREG(lock.st_mode) or lock.st_uid != uid or (lock.st_mode & 0o777) != 0o600 or lock.st_nlink != 1:
raise WorkspaceWriterConflict()
# A duplicate of the locked open description is re-lockable. An
# independently-opened description receives EWOULDBLOCK.
try:
fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError as exc:
if exc.errno in (errno.EACCES, errno.EAGAIN, errno.EWOULDBLOCK):
raise WorkspaceWriterConflict() from None
raise WorkspaceWriterConflict() from exc
finally:
try:
os.close(lock_fd)
except OSError:
pass
return WorkspaceCapability(wid, rev, root.st_dev, root.st_ino, writer.st_dev, writer.st_ino) return WorkspaceCapability(wid, rev, root.st_dev, root.st_ino, writer.st_dev, writer.st_ino)
def require_workspace_writer_capability() -> WorkspaceCapability:
return verify_workspace_writer_fds() def require_workspace_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> WorkspaceCapability:
cap = verify_workspace_writer_fds()
if workspace_id is not None and cap.workspace_id != workspace_id:
raise WorkspaceWriterConflict()
if revision is not None and cap.revision != revision:
raise WorkspaceWriterConflict()
return cap