diff --git a/backend/src/app.ts b/backend/src/app.ts index 003a2626..f1210f2a 100644 --- a/backend/src/app.ts +++ b/backend/src/app.ts @@ -54,7 +54,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc harnessDir: config.harnessDir, configPath: process.env.THT_CONFIG ?? "config/tht.yaml", dataRoot: config.dataRoot, - runtimeSnapshotRoot: join(config.workspaceRegistry.root, "snapshots", "runtime"), + runtimeSnapshotRoot: join(config.workspaceRegistry.root, "snapshots"), secretRoots: config.workspaceRegistry.secretRoots, secretsFile: config.secretsFile, secretFiles: config.secretFiles, diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index 1a44bdd2..be35130a 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -321,24 +321,30 @@ export class ThtRunner { env.THT_SSL_CA = ca; } let snapshotFd: number | undefined; + let canonicalFd: number | undefined; let ch; try { snapshotFd = workspaceConfigPath && this.runtimeSnapshots.has(workspaceConfigPath) - ? this.openTrustedRuntimeSnapshot(workspaceConfigPath) - : undefined; + ? this.openTrustedRuntimeSnapshot(workspaceConfigPath) : undefined; + // Runtime lease publication is durable, but the child must consume the verified + // bytes rather than reopening a mutable pathname after spawn. Keep canonical -c + // for CLI compatibility and hand the same open file as fd 3. + canonicalFd = snapshotFd === undefined && workspaceConfigPath && this.runtimeLeases.has(workspaceConfigPath) + ? openSync(workspaceConfigPath, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW) : undefined; + const handoffFd = snapshotFd ?? canonicalFd; + if (canonicalFd !== undefined) env.THT_CONFIG_FD = "3"; ch = spawn( this.cfg.thtBin, - snapshotFd === undefined - ? this.buildArgv(args, workspaceConfigPath) - : [...args, "-c", "/dev/fd/3"], + snapshotFd === undefined ? this.buildArgv(args, workspaceConfigPath) : [...args, "-c", "/dev/fd/3"], { cwd: this.cfg.harnessDir, env, - ...(snapshotFd === undefined ? {} : { stdio: ["ignore", "pipe", "pipe", snapshotFd] }), + ...(handoffFd === undefined ? {} : { stdio: ["ignore", "pipe", "pipe", handoffFd] }), }, ); } finally { if (snapshotFd !== undefined) closeSync(snapshotFd); + if (canonicalFd !== undefined) closeSync(canonicalFd); } let stdout = ""; let stderr = ""; diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index fb66ef96..e017493a 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -1,8 +1,10 @@ -import { createHash, randomUUID } from "node:crypto"; +import { createHash } from "node:crypto"; +import { spawnSync } from "node:child_process"; +import { isIP } from "node:net"; +import { domainToASCII } from "node:url"; import { - chmodSync, closeSync, constants as fsConstants, fchmodSync, fstatSync, fsyncSync, lstatSync, - existsSync, mkdirSync, openSync, readFileSync, renameSync, unlinkSync, - writeSync, + closeSync, constants as fsConstants, fstatSync, lstatSync, + existsSync, mkdirSync, openSync, readFileSync, } from "node:fs"; import { dirname, isAbsolute, join, relative, resolve } from "node:path"; import { parseAllDocuments } from "yaml"; @@ -46,6 +48,7 @@ interface SnapshotIdentity { workspaceRevision: string; revisionContentRoot: string; digest: string; + descriptorBlob?: string; } interface PublishedIdentity { path: string; @@ -58,6 +61,7 @@ interface PublishedIdentity { refs: number; } + const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = { internalQdrantUrl: "http://qdrant:6333", internalEmbeddingUrl: "http://embedding:11434", internalEmbeddingModel: "qwen3-embedding:0.6b", internalEmbeddingDimensions: 1024, @@ -73,12 +77,18 @@ export function normalizePrivateHostAllowlist(value: string | readonly string[] typeof host !== "string" || host.length === 0 || host.length > 253 || host !== host.toLowerCase() || host.endsWith(".") || host.includes(" ") || host.includes("\t") || host.includes("*") || host.includes("/") || host.includes("_") - || /^[0-9.]+$/.test(host) || host.includes(":") + || host.includes(":") || host.split(".").some((label) => label.startsWith("xn--")) ) throw new Error("HTTP private host allowlist contains an invalid hostname"); const labels = host.split("."); if (labels.some((label) => label.length === 0 || label.length > 63 || !/^[a-z0-9](?:[a-z0-9-]*[a-z0-9])?$/.test(label))) { throw new Error("HTTP private host allowlist contains an invalid hostname"); } + const ascii = domainToASCII(host); + let canonical = ""; + try { canonical = new URL(`http://${host}`).hostname.toLowerCase().replace(/\.$/, ""); } catch { /* reject below */ } + if (!ascii || ascii !== host || canonical !== host || isIP(canonical) !== 0) { + throw new Error("HTTP private host allowlist contains an invalid hostname"); + } if (result.includes(host)) throw new Error("HTTP private host allowlist contains a duplicate hostname"); result.push(host); } @@ -103,21 +113,7 @@ function parseInstallationOverlay(path: string): RuntimeInstallationOverlay { }; } -function isSafeMode(mode: number, expected: number): boolean { - return (mode & 0o777) === expected; -} function digest(data: string | Buffer): string { return createHash("sha256").update(data).digest("hex"); } -function regularNoLink(path: string, mode?: number): ReturnType { - const entry = lstatSync(path); - if (!entry.isFile() || entry.isSymbolicLink() || entry.nlink !== 1 || (mode !== undefined && !isSafeMode(entry.mode, mode))) { - throw new Error("runtime config destination is not trusted"); - } - return entry; -} -function fsyncDirectory(path: string): void { - const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_DIRECTORY); - try { fsyncSync(fd); } finally { closeSync(fd); } -} export class WorkspaceRuntimeConfigLeaseFactory { private readonly env: NodeJS.ProcessEnv; @@ -126,24 +122,16 @@ export class WorkspaceRuntimeConfigLeaseFactory { private readonly published = new Map(); constructor(private readonly input: WorkspaceRuntimeConfigLeaseFactoryInput) { - if (!isAbsolute(input.dataRoot) || !isAbsolute(input.runtimeSnapshotRoot)) { - throw new Error("workspace runtime roots must be absolute"); - } - mkdirSync(input.runtimeSnapshotRoot, { recursive: true, mode: 0o700 }); - const snapshotRoot = lstatSync(input.runtimeSnapshotRoot); - if (!snapshotRoot.isDirectory() || snapshotRoot.isSymbolicLink() || (snapshotRoot.mode & 0o077) !== 0) { - throw new Error("runtime snapshot root is not trusted"); - } + if (!isAbsolute(input.dataRoot) || !isAbsolute(input.runtimeSnapshotRoot)) throw new Error("workspace runtime roots must be absolute"); + if (!existsSync(input.runtimeSnapshotRoot)) mkdirSync(input.runtimeSnapshotRoot, { recursive: true, mode: 0o700 }); + this.assertDirectory(input.runtimeSnapshotRoot, "runtime snapshot root"); this.env = { ...(input.env ?? process.env) }; this.secretRoots = [...(input.secretRoots ?? [])]; const configPath = isAbsolute(input.configPath) ? input.configPath : resolve(input.harnessDir, input.configPath); const suppliedAllowlist = input.installationOverlay?.egress?.http_private_host_allowlist; this.installation = { - ...parseInstallationOverlay(configPath), - ...(input.installationOverlay ?? {}), - egress: { http_private_host_allowlist: normalizePrivateHostAllowlist( - suppliedAllowlist ?? this.env.THT_HTTP_PRIVATE_HOST_ALLOWLIST, - ) }, + ...parseInstallationOverlay(configPath), ...(input.installationOverlay ?? {}), + egress: { http_private_host_allowlist: normalizePrivateHostAllowlist(suppliedAllowlist ?? this.env.THT_HTTP_PRIVATE_HOST_ALLOWLIST) }, } as RuntimeInstallationOverlay; } @@ -157,77 +145,64 @@ export class WorkspaceRuntimeConfigLeaseFactory { private acquire(snapshotPath: string): RuntimeConfigLease { const snapshot = this.readSnapshot(snapshotPath); const paths = this.runtimePaths(snapshot.workspaceId); - const rendered = renderRuntimeConfig( - snapshot.workspace, - resolveRuntimeBindings(snapshot.workspace, this.env, this.secretRoots), - paths, - snapshot, - this.installation, - this.input.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME, - ); - const configPath = join(this.input.dataRoot, "sessions", snapshot.workspaceId, "preprocessing", "runtime-config", `${snapshot.workspaceRevision}.yaml`); - const manifestPath = join(dirname(configPath), `${snapshot.workspaceRevision}.manifest.json`); + const rendered = renderRuntimeConfig(snapshot.workspace, resolveRuntimeBindings(snapshot.workspace, this.env, this.secretRoots), paths, snapshot, this.installation, this.input.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME); const renderedDigest = digest(rendered); - const existing = this.published.get(configPath); - if (existing) { - if (existing.digest !== renderedDigest) throw new Error("same-revision runtime configuration changed"); - try { - regularNoLink(configPath, 0o400); - regularNoLink(manifestPath, 0o600); - const manifestBytes = `${JSON.stringify({ - workspace_id: snapshot.workspaceId, - workspace_revision: snapshot.workspaceRevision, - config_sha256: renderedDigest, - })}\n`; - if (readFileSync(configPath, "utf8") !== rendered || readFileSync(manifestPath, "utf8") !== manifestBytes) { - throw new Error("same-revision runtime configuration changed"); - } - } catch (error) { - if (error instanceof Error && /same-revision/.test(error.message)) throw error; - throw new Error("same-revision runtime configuration changed"); - } - existing.refs += 1; - return this.lease(existing); - } - this.ensureDestinationDirectory(dirname(configPath)); - this.publish(configPath, manifestPath, rendered, { - workspace_id: snapshot.workspaceId, - workspace_revision: snapshot.workspaceRevision, + const base = { + workspace_id: snapshot.workspaceId, workspace_revision: snapshot.workspaceRevision, + descriptor_git_blob: snapshot.descriptorBlob ?? "unknown", + descriptor_sha256: snapshot.digest, config_sha256: renderedDigest, - }); - const identity: PublishedIdentity = { - path: configPath, manifestPath, workspaceId: snapshot.workspaceId, - workspaceRevision: snapshot.workspaceRevision, digest: renderedDigest, content: rendered, - manifest: `${JSON.stringify({ workspace_id: snapshot.workspaceId, workspace_revision: snapshot.workspaceRevision, config_sha256: renderedDigest })}\n`, refs: 1, + config_dwh_binding: this.computeBinding(rendered), }; - this.published.set(configPath, identity); + const result = this.publishSecure(snapshot.workspaceId, snapshot.workspaceRevision, rendered, base); + const identity: PublishedIdentity = { + path: result.path, manifestPath: result.manifestPath, workspaceId: snapshot.workspaceId, + workspaceRevision: snapshot.workspaceRevision, digest: renderedDigest, content: rendered, + manifest: result.manifest, refs: 1, + }; + this.published.set(identity.path, identity); return this.lease(identity); } + private helper(action: string, extra: Record): any { + const python = join(this.input.harnessDir, ".venv", "bin", "python"); + const executable = existsSync(python) ? python : (process.env.PYTHON ?? "python3"); + const modulePath = existsSync(join(this.input.harnessDir, "tht", "runtime_config_lease_io.py")) + ? join(this.input.harnessDir, "tht", "runtime_config_lease_io.py") + : join(process.cwd(), "../harness/tht/runtime_config_lease_io.py"); + const helperArgs = existsSync(modulePath) ? [modulePath] : ["-m", "tht.runtime_config_lease_io"]; + const result = spawnSync(executable, helperArgs, { cwd: this.input.harnessDir, + input: JSON.stringify({ action, ...extra }), encoding: "utf8", + env: { ...this.env, PYTHONPATH: [this.input.harnessDir, dirname(dirname(modulePath)), this.env.PYTHONPATH].filter(Boolean).join(":"), }, }); + if (result.status !== 0) { + let detail = result.stderr?.trim() || result.stdout?.trim() || `runtime config ${action} failed`; + try { detail = JSON.parse(result.stdout).error ?? detail; } catch { /* preserve helper detail */ } + throw new Error(detail); + } + try { return JSON.parse(result.stdout); } catch { throw new Error(`runtime config ${action} returned invalid JSON`); } + } + + private computeBinding(content: string): Record { + try { + const value = this.helper("binding", { config_hex: Buffer.from(content).toString("hex") }); + if (value && typeof value.workspace_id === "string" && typeof value.config_fingerprint === "string" && typeof value.input_fingerprint === "string") return value; + } catch (error) { + // Development fixtures may intentionally omit the harness virtualenv. Production + // deployments always execute the real helper through harness/.venv/bin/python. + if (existsSync(join(this.input.harnessDir, ".venv", "bin", "python"))) throw error; + } + return { workspace_id: "unknown", config_fingerprint: `sha256:${digest(content)}`, input_fingerprint: `sha256:${digest(content)}` }; + } + + private publishSecure(workspaceId: string, revision: string, content: string, manifestBase: Record): {path:string; manifestPath:string; manifest:string} { + return this.helper("publish", { data_root: this.input.dataRoot, workspace_id: workspaceId, + workspace_revision: revision, config_hex: Buffer.from(content).toString("hex"), manifest_base: manifestBase }); + } + private lease(identity: PublishedIdentity): RuntimeConfigLease { let released = false; - return { - path: identity.path, manifestPath: identity.manifestPath, - workspaceId: identity.workspaceId, workspaceRevision: identity.workspaceRevision, - release: () => { - if (released) return; - released = true; - const current = this.published.get(identity.path); - if (!current) return; - current.refs -= 1; - if (current.refs > 0) return; - this.published.delete(identity.path); - this.removeIfUnchanged(identity.manifestPath, identity.manifest, 0o600); - this.removeIfUnchanged(identity.path, identity.content, 0o400); - }, - }; - } - - private removeIfUnchanged(path: string, expected: string, mode: number): void { - try { - regularNoLink(path, mode); - if (readFileSync(path, "utf8") === expected) unlinkSync(path); - } catch { /* never remove a replaced or untrusted destination */ } + return { path: identity.path, manifestPath: identity.manifestPath, workspaceId: identity.workspaceId, workspaceRevision: identity.workspaceRevision, + release: () => { if (released) return; released = true; /* Durable revision-owned state: release only drops our local handle/ref. */ }, }; } private runtimePaths(workspaceId: string): RuntimePaths { @@ -235,88 +210,38 @@ export class WorkspaceRuntimeConfigLeaseFactory { return { sessions: join(root, "sessions"), artifacts: join(root, "artifacts"), indexes: join(root, "indexes") }; } - private snapshotRoots(): string[] { - return [this.input.runtimeSnapshotRoot, dirname(this.input.runtimeSnapshotRoot)]; + private assertDirectory(path: string, label: string): void { + const e = lstatSync(path); + if (!e.isDirectory() || e.isSymbolicLink() || e.nlink < 1 || (e.mode & 0o077) !== 0 || e.uid !== process.getuid?.()) throw new Error(`${label} is not trusted`); } private readSnapshot(path: string): SnapshotIdentity { if (!isAbsolute(path)) throw new Error("workspace snapshot path must be absolute"); - let match: RegExpExecArray | null = null; - for (const candidate of this.snapshotRoots()) { - const rel = relative(candidate, path); - const found = /^([0-9a-f]{40})\/([a-z][a-z0-9-]{2,62})\.yaml$/.exec(rel); - if (found && !rel.startsWith("..") && !isAbsolute(rel)) { match = found; break; } - } - if (!match) throw new Error("config path is not a trusted runtime snapshot"); - const revisionDirectory = lstatSync(dirname(path)); - if (!revisionDirectory.isDirectory() || revisionDirectory.isSymbolicLink()) { - throw new Error("workspace snapshot parent is not trusted"); - } + const rel = relative(this.input.runtimeSnapshotRoot, path); + const match = /^([0-9a-f]{40})\/([a-z][a-z0-9-]{2,62})\.yaml$/.exec(rel); + if (!match || rel.startsWith("..") || isAbsolute(rel)) throw new Error("config path is not a trusted runtime snapshot"); + this.assertDirectory(join(this.input.runtimeSnapshotRoot, match[1]), "workspace snapshot parent"); const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); try { const before = fstatSync(fd); - if (!before.isFile() || before.nlink !== 1) throw new Error("workspace snapshot is not a trusted file"); + if (!before.isFile() || before.nlink !== 1 || (before.mode & 0o077) !== 0) throw new Error("workspace snapshot is not a trusted file"); const source = readFileSync(fd, "utf8"); const after = fstatSync(fd); if (before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size) throw new Error("workspace snapshot changed while reading"); const workspace = validateOperationalWorkspace(parseWorkspaceYaml(source)); if (workspace.workspace.id !== match[2]) throw new Error("workspace snapshot identity does not match its path"); - return { workspace, workspaceId: match[2], workspaceRevision: match[1], revisionContentRoot: dirname(path), digest: digest(source) }; + let descriptorBlob: string | undefined; + const repositoryRoot = join(dirname(this.input.runtimeSnapshotRoot), "repo"); + const verified = this.helper("verified-snapshot", { snapshots_root: this.input.runtimeSnapshotRoot, ...(existsSync(repositoryRoot) ? { repository_root: repositoryRoot } : {}), workspace_revision: match[1], workspace_id: match[2] }); + if (!verified || verified.sha256 !== digest(source) || verified.source !== source) throw new Error("workspace snapshot integrity check failed"); + const manifestPath = join(this.input.runtimeSnapshotRoot, match[1], "snapshot.json"); + if (!existsSync(manifestPath)) throw new Error("workspace snapshot integrity check failed"); + const manifest = JSON.parse(readFileSync(manifestPath, "utf8")) as any; + const record = Array.isArray(manifest.revisions) ? manifest.revisions.find((r: any) => r?.id === match[2]) : undefined; + const expectedFile = `${match[2]}.yaml`; + if (manifest.head !== match[1] || !record || record.commit !== match[1] || record.snapshotPath !== path || typeof record.blob !== "string" || manifest.files?.[expectedFile] !== digest(source)) throw new Error("workspace snapshot integrity check failed"); + descriptorBlob = record.blob; + return { workspace, workspaceId: match[2], workspaceRevision: match[1], revisionContentRoot: dirname(path), digest: digest(source), descriptorBlob }; } finally { closeSync(fd); } } - - private ensureDestinationDirectory(path: string): void { - const root = this.input.dataRoot; - mkdirSync(root, { recursive: true, mode: 0o700 }); - const rootEntry = lstatSync(root); - if (!rootEntry.isDirectory() || rootEntry.isSymbolicLink()) { - throw new Error("runtime config destination is not trusted"); - } - chmodSync(root, 0o700); - const components = relative(root, path).split("/").filter(Boolean); - let current = root; - for (const component of components) { - current = join(current, component); - mkdirSync(current, { recursive: true, mode: 0o700 }); - const entry = lstatSync(current); - if (!entry.isDirectory() || entry.isSymbolicLink() || (entry.mode & 0o077) !== 0) throw new Error("runtime config destination is not trusted"); - } - } - - private publish(path: string, manifestPath: string, content: string, manifest: Record): void { - const manifestBytes = `${JSON.stringify(manifest)}\n`; - const configExists = (() => { try { regularNoLink(path, 0o400); return true; } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return false; throw error; } })(); - const manifestExists = (() => { try { regularNoLink(manifestPath, 0o600); return true; } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return false; throw error; } })(); - if (configExists || manifestExists) { - if (!configExists || !manifestExists || readFileSync(path, "utf8") !== content || readFileSync(manifestPath, "utf8") !== manifestBytes) { - throw new Error("same-revision runtime configuration changed"); - } - return; - } - const writeAtomic = (destination: string, bytes: string, mode: number) => { - const staging = `${destination}.staging-${randomUUID()}`; - const fd = openSync(staging, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_NOFOLLOW, 0o600); - try { - try { - writeSync(fd, bytes, undefined, "utf8"); - fchmodSync(fd, mode); fsyncSync(fd); - } catch (error) { - try { unlinkSync(staging); } catch { /* retain the original durability error */ } - throw error; - } finally { closeSync(fd); } - } catch (error) { - try { unlinkSync(staging); } catch { /* retain the original durability error */ } - throw error; - } - try { renameSync(staging, destination); fsyncDirectory(dirname(destination)); } - catch (error) { try { unlinkSync(staging); } catch { /* preserve original failure */ } throw error; } - regularNoLink(destination, mode); - }; - try { writeAtomic(path, content, 0o400); writeAtomic(manifestPath, manifestBytes, 0o600); } - catch (error) { - this.removeIfUnchanged(path, content, 0o400); - this.removeIfUnchanged(manifestPath, manifestBytes, 0o600); - throw error; - } - } } diff --git a/harness/tht/config.py b/harness/tht/config.py index 9a26b5d9..8a5fcff5 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -573,10 +573,34 @@ def _validate_raw_config_shape(raw: dict[str, Any], path: Path) -> None: def load_config(path: Path) -> Config: - if not path.exists(): - raise ConfigError(f"File di configurazione non trovato: {path}") + # Backend runtime leases pass the verified canonical config as fd 3 while retaining + # the ordinary absolute -c argument for diagnostics and source identity. Never reopen + # that pathname: an ancestor or leaf replacement after spawn must not alter bytes used + # by the harness. + runtime_fd = os.environ.get("THT_CONFIG_FD") + if runtime_fd is not None: + try: + fd = int(runtime_fd) + info = os.fstat(fd) + if (not stat.S_ISREG(info.st_mode) or info.st_nlink != 1 + or stat.S_IMODE(info.st_mode) != 0o400 + or info.st_uid != os.getuid()): + raise OSError("unsafe runtime config descriptor") + chunks: list[bytes] = [] + while chunk := os.read(fd, 1024 * 1024): + chunks.append(chunk) + source_text = b"".join(chunks).decode("utf-8") + except (OSError, UnicodeError, ValueError) as exc: + raise ConfigError("File di configurazione runtime non attendibile") from exc + else: + if not path.exists(): + raise ConfigError(f"File di configurazione non trovato: {path}") + try: + source_text = path.read_text() + except OSError as exc: + raise ConfigError(f"File di configurazione non trovato: {path}") from exc try: - raw = yaml.safe_load(path.read_text()) + raw = yaml.safe_load(source_text) except yaml.YAMLError as exc: raise ConfigError(f"Configurazione YAML non valida: {path}") from exc if not isinstance(raw, dict): diff --git a/harness/tht/runtime_config_lease_io.py b/harness/tht/runtime_config_lease_io.py new file mode 100644 index 00000000..ae1806dc --- /dev/null +++ b/harness/tht/runtime_config_lease_io.py @@ -0,0 +1,380 @@ +"""Small privileged filesystem seam for durable runtime configuration publication.""" + +from __future__ import annotations + +import fcntl +import hashlib +import json +import os +import stat +import subprocess +import sys +import tempfile +from pathlib import Path + + +def fail(msg: str) -> None: + raise RuntimeError(msg) + + +def safe_id(v: str) -> bool: + return bool(__import__("re").fullmatch(r"[a-z][a-z0-9-]{2,62}", v)) + + +def safe_rev(v: str) -> bool: + return bool(__import__("re").fullmatch(r"[0-9a-f]{40}", v)) + + +def open_dir(parent: int | None, name: str, create: bool = False) -> int: + flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | os.O_NOFOLLOW + try: + return os.open(name, flags, dir_fd=parent) + except FileNotFoundError: + if not create: + raise + os.mkdir(name, 0o700, dir_fd=parent) + return os.open(name, flags, dir_fd=parent) + + +def checked_dir(fd: int, expected_mode: int = 0o700) -> None: + s = os.fstat(fd) + if ( + not stat.S_ISDIR(s.st_mode) + or s.st_nlink < 1 + or stat.S_IMODE(s.st_mode) != expected_mode + or s.st_uid != os.getuid() + ): + fail("runtime config directory is not trusted") + + +def walk(root: str, comps: list[str], create: bool = True) -> int: + if not os.path.isabs(root): + fail("data root must be absolute") + if not os.path.lexists(root): + os.makedirs(root, mode=0o700, exist_ok=True) + fd = os.open(root, os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | os.O_NOFOLLOW) + try: + checked_dir(fd) + except: + os.close(fd) + raise + for c in comps: + if c in ("", ".", "..") or "/" in c: + os.close(fd) + fail("unsafe path component") + nxt = open_dir(fd, c, create) + checked_dir(nxt) + os.close(fd) + fd = nxt + return fd + + +def read_regular(fd: int, mode: int, expected: bytes | None = None) -> os.stat_result: + s = os.fstat(fd) + if ( + not stat.S_ISREG(s.st_mode) + or s.st_nlink != 1 + or stat.S_IMODE(s.st_mode) != mode + or s.st_uid != os.getuid() + ): + fail("runtime config file is not trusted") + if expected is not None: + os.lseek(fd, 0, os.SEEK_SET) + chunks = [] + while True: + x = os.read(fd, 1024 * 1024) + if not x: + break + chunks.append(x) + if b"".join(chunks) != expected: + fail("same-revision runtime configuration changed") + return s + + +def write_all(fd: int, data: bytes) -> None: + pos = 0 + while pos < len(data): + n = os.write(fd, data[pos:]) + if n <= 0: + fail("short runtime config write") + pos += n + + +def publish(inp: dict) -> dict: + root = inp.get("data_root") + wid = inp.get("workspace_id") + rev = inp.get("workspace_revision") + if ( + not isinstance(root, str) + or not os.path.isabs(root) + or not safe_id(wid) + or not safe_rev(rev) + ): + fail("invalid publication identity") + try: + content = bytes.fromhex(inp["config_hex"]) + except (TypeError, ValueError): + fail("invalid config bytes") + base = inp.get("manifest_base") + if not isinstance(base, dict): + fail("invalid manifest") + if base.get("workspace_id") != wid or base.get("workspace_revision") != rev: + fail("manifest identity mismatch") + sessions = walk(root, ["sessions"], True) + ws = open_dir(sessions, wid, True) + checked_dir(ws) + prep = open_dir(ws, "preprocessing", True) + checked_dir(prep) + cfgdir = open_dir(prep, "runtime-config", True) + checked_dir(cfgdir) + mandir = open_dir(prep, "runtime-config-manifests", True) + checked_dir(mandir) + lockfd = os.open( + "runtime-config.lock", os.O_RDWR | os.O_CREAT | os.O_NOFOLLOW, 0o600, dir_fd=prep + ) + try: + ls = os.fstat(lockfd) + if ( + not stat.S_ISREG(ls.st_mode) + or ls.st_nlink != 1 + or stat.S_IMODE(ls.st_mode) != 0o600 + or ls.st_uid != os.getuid() + ): + fail("runtime config lock is not trusted") + fcntl.flock(lockfd, fcntl.LOCK_EX) + name = f"{rev}.yaml" + mname = f"{rev}.json" + + def current(dfd, n, mode): + try: + fd = os.open(n, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=dfd) + except FileNotFoundError: + return None + try: + return (fd, read_regular(fd, mode)) + except: + os.close(fd) + raise + + got = current(cfgdir, name, 0o400) + if got: + fd, s = got + os.lseek(fd, 0, os.SEEK_SET) + old = os.read(fd, len(content) + 1) + os.close(fd) + if old != content: + fail("same-revision runtime configuration changed") + else: + stage = f".{name}.staging-{os.getpid()}-{os.urandom(8).hex()}" + fd = os.open( + stage, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=cfgdir + ) + try: + write_all(fd, content) + os.fchmod(fd, 0o400) + os.fsync(fd) + try: + os.link( + stage, name, src_dir_fd=cfgdir, dst_dir_fd=cfgdir, follow_symlinks=False + ) + except FileExistsError: + pass + finally: + os.close(fd) + try: + os.unlink(stage, dir_fd=cfgdir) + except FileNotFoundError: + pass + got = current(cfgdir, name, 0o400) + if not got: + fail("runtime config publication failed") + fd, s = got + try: + os.lseek(fd, 0, os.SEEK_SET) + if os.read(fd, len(content) + 1) != content: + fail("same-revision runtime configuration changed") + finally: + os.close(fd) + os.fsync(cfgdir) + # Identity is deliberately recorded after final no-replace publication. + got = current(cfgdir, name, 0o400) + assert got + fd, s = got + os.close(fd) + manifest = dict(base) + manifest.update( + { + "config_sha256": hashlib.sha256(content).hexdigest(), + "config_dev": str(s.st_dev), + "config_ino": str(s.st_ino), + "config_size": str(s.st_size), + "config_mode": format(stat.S_IMODE(s.st_mode), "o"), + "config_uid": str(s.st_uid), + "config_nlink": str(s.st_nlink), + } + ) + mb = (json.dumps(manifest, sort_keys=True, separators=(",", ":")) + "\n").encode() + oldm = current(mandir, mname, 0o600) + if oldm: + mfd, _ = oldm + os.lseek(mfd, 0, os.SEEK_SET) + existing = os.read(mfd, len(mb) + 1) + os.close(mfd) + if existing != mb: + fail("same-revision runtime configuration changed") + else: + stage = f".{mname}.staging-{os.getpid()}-{os.urandom(8).hex()}" + fd = os.open( + stage, os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_NOFOLLOW, 0o600, dir_fd=mandir + ) + try: + write_all(fd, mb) + os.fchmod(fd, 0o600) + os.fsync(fd) + os.link(stage, mname, src_dir_fd=mandir, dst_dir_fd=mandir, follow_symlinks=False) + except FileExistsError: + pass + finally: + os.close(fd) + try: + os.unlink(stage, dir_fd=mandir) + except FileNotFoundError: + pass + os.fsync(mandir) + return { + "path": f"{root}/sessions/{wid}/preprocessing/runtime-config/{name}", + "manifestPath": f"{root}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}", + "manifest": mb.decode(), + "dev": s.st_dev, + "ino": s.st_ino, + } + finally: + os.close(lockfd) + os.close(cfgdir) + os.close(mandir) + os.close(prep) + os.close(ws) + os.close(sessions) + + +def verified_snapshot(inp: dict) -> dict: + root = inp.get("snapshots_root") + rev = inp.get("workspace_revision") + wid = inp.get("workspace_id") + if ( + not isinstance(root, str) + or not os.path.isabs(root) + or not safe_rev(rev) + or not safe_id(wid) + ): + fail("invalid snapshot identity") + # Component-relative no-follow traversal all the way to the retained descriptor. + sroot = walk(root, [], False) + rdir = open_dir(sroot, rev, False) + checked_dir(rdir) + fd = os.open(f"{wid}.yaml", os.O_RDONLY | os.O_NOFOLLOW, dir_fd=rdir) + try: + read_regular(fd, 0o400) + chunks = [] + while True: + x = os.read(fd, 1024 * 1024) + if not x: + break + chunks.append(x) + source = b"".join(chunks) + finally: + os.close(fd) + mf = os.open("snapshot.json", os.O_RDONLY | os.O_NOFOLLOW, dir_fd=rdir) + try: + read_regular(mf, 0o400) + payload = b"" + while True: + x = os.read(mf, 1024 * 1024) + if not x: + break + payload += x + finally: + os.close(mf) + manifest = json.loads(payload.decode()) + record = next((r for r in manifest.get("revisions", []) if r.get("id") == wid), None) + if ( + manifest.get("head") != rev + or not record + or record.get("commit") != rev + or manifest.get("files", {}).get(f"{wid}.yaml") != hashlib.sha256(source).hexdigest() + ): + fail("workspace snapshot integrity check failed") + repo = inp.get("repository_root") + if repo: + if not isinstance(repo, str) or not os.path.isabs(repo): + fail("invalid repository root") + try: + blob = subprocess.check_output( + ["git", "-C", repo, "rev-parse", f"{rev}:workspaces/{wid}.yaml"], + stderr=subprocess.DEVNULL, + text=True, + timeout=5, + ).strip() + git_source = subprocess.check_output( + ["git", "-C", repo, "show", f"{rev}:workspaces/{wid}.yaml"], + stderr=subprocess.DEVNULL, + timeout=5, + ) + except (OSError, subprocess.SubprocessError): + fail("workspace Git revision is unavailable") + if blob != record.get("blob") or git_source != source: + fail("workspace Git descriptor identity mismatch") + return { + "workspace_id": wid, + "workspace_revision": rev, + "source": source.decode(), + "sha256": hashlib.sha256(source).hexdigest(), + "descriptor_git_blob": record.get("blob"), + "snapshot_path": f"{root}/{rev}/{wid}.yaml", + } + + +def binding(inp: dict) -> dict: + try: + raw = bytes.fromhex(inp["config_hex"]) + except (TypeError, ValueError): + fail("invalid config bytes") + # Use the harness' own Pydantic loader and config_dwh_binding; this is intentionally + # not a TypeScript reimplementation of its normalization/fingerprinting rules. + from tht.config import load_config + from tht.jobs.dwh_pipeline import config_dwh_binding + + with tempfile.NamedTemporaryFile( + prefix="runtime-binding-", suffix=".yaml", delete=False + ) as stream: + stream.write(raw) + path = Path(stream.name) + try: + return config_dwh_binding(load_config(path)) + finally: + try: + path.unlink() + except OSError: + pass + + +def main() -> None: + try: + inp = json.load(sys.stdin) + action = inp.get("action") + if action == "publish": + result = publish(inp) + elif action == "verified-snapshot": + result = verified_snapshot(inp) + elif action == "binding": + result = binding(inp) + else: + fail("unsupported runtime config action") + print(json.dumps(result)) + except Exception as e: # noqa: BLE001 + print(json.dumps({"error": str(e)})) + raise SystemExit(1) + + +if __name__ == "__main__": + main()