From ee5ac381b385a72f4fe111db36ed5c0142084a33 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 14:27:43 +0200 Subject: [PATCH] fix: harden preprocessing state and child capabilities --- backend/src/workspaces/preprocessing-state.ts | 186 ++++++++++++------ backend/src/workspaces/workspace-fs-at.ts | 27 +++ .../workspaces/workspace-lock-root-lease.ts | 4 +- .../workspace-preprocessing-state.test.ts | 37 +++- harness/tests/test_preprocess_cli.py | 8 + harness/tests/test_qdrant_cli_commands.py | 9 + harness/tests/test_schema_fk_annotations.py | 8 +- harness/tests/test_workspace_writer_lock.py | 21 +- harness/tht/cli/preprocess_cmd.py | 16 +- harness/tht/cli/schema_cmd.py | 11 +- harness/tht/cli/vector_cmd.py | 16 +- harness/tht/workspace_writer_lock.py | 84 ++++++-- 12 files changed, 323 insertions(+), 104 deletions(-) diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index 162693e2..b1c3c1c8 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -1,6 +1,8 @@ import { createHash, randomBytes } from "node:crypto"; -import { mkdir, readFile, rename, rm, writeFile, open as openFile } from "node:fs/promises"; -import { join } from "node:path"; +import { lstatSync, realpathSync } 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, @@ -10,6 +12,7 @@ import { VerifiedWorkspaceLockRootLease, type CanonicalWorkspaceId, type WorkspaceLockRootIdentityV1, + WorkspaceRootLock, } from "./workspace-lock-root-lease.js"; 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 revision: string; readonly operation: string; readonly phase: 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 ResumeRunInput { readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; readonly stateRoot?: string; } -export type RunTransition = { readonly phase: string; readonly [key: string]: unknown }; +export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; } +export interface ResumeRunInput { readonly runId: string; readonly workspaceId: string; readonly revision: string; 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 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 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; } 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 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 { + 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(5).toString("hex")}`; + const tmp = `${path}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`; try { 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); - 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; } } export class PreprocessingStateStore { - constructor(private readonly configuredRoot?: string) {} - private root(input?: string): string { return checkedRoot(input ?? this.configuredRoot); } - private paths(input: { stateRoot?: string; runId: string }) { - id(input.runId); const root = this.root(input.stateRoot); - 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 readonly configuredRoot: string; + constructor(configuredRoot: string) { this.configuredRoot = checkedRoot(configuredRoot); } + private root(): string { return this.configuredRoot; } + 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({ stateRoot: input.stateRoot, runId }); - if (typeof input.workspaceId !== "string" || !/^[a-z][a-z0-9-]{2,62}$/.test(input.workspaceId) || !REVISION.test(input.revision) || !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 }); - 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 writeFile(p.path, `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw fail(); } + 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 { 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"); return value; } async load(input: ResumeRunInput): Promise { return this.loadForResume(input); } - async transition(runId: string, transition: RunTransition, stateRoot?: string): Promise { - const p = this.paths({ stateRoot, 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(); + 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, ...transition, schemaVersion: 1, runId, artifacts: current.artifacts, updatedAt: new Date().toISOString() }; + 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, stateRoot?: string): Promise { - const p = this.paths({ stateRoot, runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > 10 * 1024 * 1024) throw fail(); - 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 }); + 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, stateRoot?: string): Promise { - 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 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; + 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 { return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } } + 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 (!value || typeof value !== "object") return false; const x = value as Record; - 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"; + 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); } } @@ -102,42 +146,62 @@ export interface WorkspaceLockedChildResult { readonly exitCode: number; readonl export class BorrowedWorkspaceSessionReadersExclusiveLockLease { 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(); } invalidate(): void { this.live = false; } } export class WorkspaceWriterLockCapability { - private live = true; private readerExclusive = false; - private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1, private readonly root: VerifiedWorkspaceLockRootLease, private readonly writer: import("./workspace-lock-root-lease.js").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 check(): void { if (!this.live) throw fail(); } + 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(); } + invalidateForSettlement(): void { this.settled = true; } + static create(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); } async runUnderSessionReadersExclusive(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise): Promise { - this.check(); if (this.readerExclusive) throw fail(); this.readerExclusive = true; - let lock: import("./workspace-lock-root-lease.js").WorkspaceRootLock | undefined; let borrowed: BorrowedWorkspaceSessionReadersExclusiveLockLease | undefined; - try { lock = await this.root.acquireSessionReadersExclusive(); borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.make(this.workspaceId, this.rootIdentity); return await action(borrowed); } - catch { throw fail(); } - finally { borrowed?.invalidate(); try { lock?.close(); } catch { this.live = false; } this.readerExclusive = false; } + this.assertLive(); if (this.readerExclusive || this.spawnActive) throw fail(); this.readerExclusive = true; + let lock: WorkspaceRootLock; + try { lock = await this.root.acquireSessionReadersExclusive(); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); } + const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.from(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.check(); throw fail("child runner is not configured"); } - async close(): Promise { 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(); } + 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 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 { 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 class OrderedWorkspaceWriterCapabilitySet { private live = true; private constructor(private readonly caps: Map) {} static make(caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } - invalidate(): void { this.live = false; } - get workspaceIds(): CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; } + 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: BorrowedVerifiedWorkspaceLockRootLease.make(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))); } + 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[] = []; - const transferred: VerifiedWorkspaceLockRootLease[] = []; - 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; } } - 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(); } - } 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; } + 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 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; } } + set = OrderedWorkspaceWriterCapabilitySet.make(new Map(caps.map(c => [c.workspaceId, c]))); result = await action(set); + } 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(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"; } } diff --git a/backend/src/workspaces/workspace-fs-at.ts b/backend/src/workspaces/workspace-fs-at.ts index 9719ecb5..19814a02 100644 --- a/backend/src/workspaces/workspace-fs-at.ts +++ b/backend/src/workspaces/workspace-fs-at.ts @@ -1,4 +1,6 @@ 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"; const require = createRequire(import.meta.url); 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)); } } + + +/** 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 {} } } +} diff --git a/backend/src/workspaces/workspace-lock-root-lease.ts b/backend/src/workspaces/workspace-lock-root-lease.ts index 18c9d23c..fb9dcac5 100644 --- a/backend/src/workspaces/workspace-lock-root-lease.ts +++ b/backend/src/workspaces/workspace-lock-root-lease.ts @@ -1,6 +1,6 @@ import { lstatSync, realpathSync } from "node:fs"; 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 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) {} } @@ -12,6 +12,7 @@ function sameIdentity(a: {device: bigint|number; inode: bigint|number} | {dev: b export class WorkspaceRootLock { private live = true; 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); } close(): void { if (!this.live) return; this.live = false; this.lock.close(); } } @@ -47,6 +48,7 @@ export class VerifiedWorkspaceLockRootLease { } async acquireWriterLock(): Promise { 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 { 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(); } async close(): Promise { if (!this.live) return; while (this.borrowed) await new Promise(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); } diff --git a/backend/test/workspace-preprocessing-state.test.ts b/backend/test/workspace-preprocessing-state.test.ts index ad10a564..36c4f404 100644 --- a/backend/test/workspace-preprocessing-state.test.ts +++ b/backend/test/workspace-preprocessing-state.test.ts @@ -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"; -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 { describe, expect, it, afterEach } from "vitest"; +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"); + }); +}); diff --git a/harness/tests/test_preprocess_cli.py b/harness/tests/test_preprocess_cli.py index 6433bb34..f768b851 100644 --- a/harness/tests/test_preprocess_cli.py +++ b/harness/tests/test_preprocess_cli.py @@ -11,6 +11,14 @@ from tht.ports.vector import VectorStoreError 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): import tht.cli.preprocess_cmd as command diff --git a/harness/tests/test_qdrant_cli_commands.py b/harness/tests/test_qdrant_cli_commands.py index 34c62534..6354b4b2 100644 --- a/harness/tests/test_qdrant_cli_commands.py +++ b/harness/tests/test_qdrant_cli_commands.py @@ -5,12 +5,21 @@ from datetime import UTC, datetime from pathlib import Path from types import SimpleNamespace +import pytest from typer.testing import CliRunner from tht.cli import app 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: def embed_documents(self, documents): return [[0.1] * 4 for _ in documents] diff --git a/harness/tests/test_schema_fk_annotations.py b/harness/tests/test_schema_fk_annotations.py index c92c5231..3468bd38 100644 --- a/harness/tests/test_schema_fk_annotations.py +++ b/harness/tests/test_schema_fk_annotations.py @@ -24,6 +24,12 @@ from tht.mschema.models import ( 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(): return PhysicalSchema( 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)], 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() for branch in branches: diff --git a/harness/tests/test_workspace_writer_lock.py b/harness/tests/test_workspace_writer_lock.py index 209a5d08..8dbd39f0 100644 --- a/harness/tests/test_workspace_writer_lock.py +++ b/harness/tests/test_workspace_writer_lock.py @@ -1,8 +1,12 @@ from __future__ import annotations -import os, stat + +import os + import pytest + from tht.workspace_writer_lock import WorkspaceWriterConflict, verify_workspace_writer_fds + def test_verifier_rejects_missing_capability(): 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) assert cap.inode == st.st_ino 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) diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 845b0957..3b2e405c 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -20,12 +20,11 @@ from tht.vectorstore.embeddings import EmbeddingsError preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts") logger = logging.getLogger(__name__) -def _require_writer_capability() -> None: - """Mutating children opt into the backend-owned fd capability contract.""" - import os - if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": - from tht.workspace_writer_lock import require_workspace_writer_capability - require_workspace_writer_capability() +def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None: + # Authorization is unconditional: an environment marker is attacker-controlled + # and must never turn a mutating direct invocation into an authorized child. + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability(workspace_id=workspace_id, revision=revision) _PREPROCESS_EXPECTED_ERRORS = ( OSError, RuntimeError, ValueError, TypeError, KeyError, @@ -36,7 +35,6 @@ _PREPROCESS_EXPECTED_ERRORS = ( def run_dwh_from_config( config: Path, *, steps: tuple[str, ...], resume: str | None = None, ): - _require_writer_capability() from tht.cli.lsh_cmd import build_lsh_artifacts from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog from tht.jobs.dwh_pipeline import ( @@ -45,6 +43,7 @@ def run_dwh_from_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) workspace_root = cfg.paths.artifacts.parent 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): - _require_writer_capability() from tht.adapters.factory import build_evidence_sources, build_vector_store from tht.cli.vector_cmd import make_embedder 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 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: raise RuntimeError("embeddings are not configured") 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 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: raise RuntimeError("embeddings are not configured") corpus_root = cfg.paths.artifacts.parent / "corpus" diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 855723ba..c0dacff2 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -16,11 +16,11 @@ from tht.mschema.eligibility import classify_all schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)") logger = logging.getLogger(__name__) -def _require_writer_capability() -> None: - import os - if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": - from tht.workspace_writer_lock import require_workspace_writer_capability - require_workspace_writer_capability() +def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None: + # Authorization is unconditional: an environment marker is attacker-controlled + # and must never turn a mutating direct invocation into an authorized child. + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability(workspace_id=workspace_id, revision=revision) def _add_examples(dwh, phys, examples) -> None: @@ -569,6 +569,7 @@ def suggest_fks_cmd( for item in payload["candidates"] } 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(): ann = annotations.tables.setdefault(table_name, TableAnnotation()) from tht.mschema.models import ForeignKey diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index 3e11807a..d15f057d 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -24,11 +24,11 @@ from tht.vectorstore.store import SyncStats, content_hash vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") logger = logging.getLogger(__name__) -def _require_writer_capability() -> None: - import os - if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": - from tht.workspace_writer_lock import require_workspace_writer_capability - require_workspace_writer_capability() +def _require_writer_capability(*, workspace_id: str | None = None, revision: str | None = None) -> None: + # Authorization is unconditional: an environment marker is attacker-controlled + # and must never turn a mutating direct invocation into an authorized child. + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability(workspace_id=workspace_id, revision=revision) def make_embedder(embeddings_cfg): @@ -69,8 +69,8 @@ def open_searcher(cfg): return AdapterSearcher() -def sync_canonical_records(collection, records, *, store, embedder): - _require_writer_capability() +def sync_canonical_records(collection, records, *, store, embedder, config=None): + _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}) existing = store.existing_hashes(collection, kinds) pending = [] @@ -236,7 +236,7 @@ def index_schema_data( "schema_records", records, store=build_vector_store(cfg, require_write=True), - embedder=make_embedder(cfg.embeddings), + embedder=make_embedder(cfg.embeddings), config=cfg, ) return { "status": "succeeded", diff --git a/harness/tht/workspace_writer_lock.py b/harness/tht/workspace_writer_lock.py index 587cc0ed..d4b9bb61 100644 --- a/harness/tht/workspace_writer_lock.py +++ b/harness/tht/workspace_writer_lock.py @@ -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. -This module intentionally has no path fallback: callers either run with the capability -or fail closed before touching artifacts. +FD 3 is the inherited writer open file description and FD 4 is the retained +workspace-root directory. Environment values are descriptive identity only; +they never authorize a direct invocation. """ from __future__ import annotations +import errno import fcntl import os +import re import stat from dataclasses import dataclass class WorkspaceWriterConflict(RuntimeError): "preprocessing_conflict" + def __init__(self, message: str = "preprocessing_conflict") -> None: super().__init__(message) + @dataclass(frozen=True) class WorkspaceCapability: workspace_id: str @@ -26,26 +30,72 @@ class WorkspaceCapability: writer_device: int writer_inode: int + def _identity(env: dict[str, str]) -> tuple[str, str, int, int]: 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() - try: device, inode = int(env["THOTH_WORKSPACE_DEVICE"]), int(env["THOTH_WORKSPACE_INODE"]) - except (KeyError, ValueError): raise WorkspaceWriterConflict() 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: env = dict(os.environ if env is None else env) wid, rev, device, inode = _identity(env) - try: root = os.fstat(root_fd); writer = os.fstat(writer_fd) - except OSError as exc: raise WorkspaceWriterConflict() from exc - 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() - if not stat.S_ISREG(writer.st_mode) or writer.st_uid != os.getuid() or (writer.st_mode & 0o777) != 0o600: raise WorkspaceWriterConflict() - try: fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) - except OSError as exc: raise WorkspaceWriterConflict() from exc - # Keep the OFD locked. A lock check is necessarily best effort on some BSDs; identity and - # descriptor ownership remain mandatory and no path-based lock is accepted. + if writer_fd == root_fd or writer_fd < 0 or root_fd < 0: + raise WorkspaceWriterConflict() + root, writer = _fstat(root_fd), _fstat(writer_fd) + uid = os.getuid() + 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): + raise WorkspaceWriterConflict() + if not stat.S_ISREG(writer.st_mode) or writer.st_uid != uid or (writer.st_mode & 0o777) != 0o600 or writer.st_nlink != 1: + 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) -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