From 971315aa808c6acefd8ecd87e2c524aadf3fa3bb Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 09:59:25 +0200 Subject: [PATCH] fix: complete durable workspace runtime config handoff --- backend/src/tht/tht-runner.ts | 25 ++- .../src/workspaces/runtime-config-lease.ts | 69 +++---- .../workspace-runtime-config-lease.test.ts | 29 ++- .../test/workspace-runtime-handoff.test.ts | 9 +- harness/tht/config.py | 86 +++++++-- harness/tht/runtime_config_lease_io.py | 178 +++++++++++++----- 6 files changed, 289 insertions(+), 107 deletions(-) diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index be35130a..e459a5ac 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -322,29 +322,40 @@ export class ThtRunner { } let snapshotFd: number | undefined; let canonicalFd: number | undefined; + let manifestFd: number | undefined; let ch; try { snapshotFd = workspaceConfigPath && this.runtimeSnapshots.has(workspaceConfigPath) ? 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 lease = workspaceConfigPath ? this.runtimeLeases.get(workspaceConfigPath) : undefined; + // Registry leases retain the verified config bytes in fd 3 and the separately + // published manifest in fd 4. The argv remains the canonical -c pathname for + // diagnostics/compatibility; the harness never trusts that pathname for bytes. + canonicalFd = snapshotFd === undefined && lease + ? openSync(lease.path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW) : undefined; + manifestFd = lease + ? openSync(lease.manifestPath, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW) : undefined; + if (lease) { + env.THT_CONFIG_FD = "3"; + env.THT_CONFIG_MANIFEST_FD = "4"; + env.THT_CONFIG_MANIFEST_SHA256 = lease.manifestSha256; + } 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"], { cwd: this.cfg.harnessDir, env, - ...(handoffFd === undefined ? {} : { stdio: ["ignore", "pipe", "pipe", handoffFd] }), + ...(handoffFd === undefined ? {} : { + stdio: ["ignore", "pipe", "pipe", handoffFd, ...(manifestFd === undefined ? [] : [manifestFd])], + }), }, ); } finally { if (snapshotFd !== undefined) closeSync(snapshotFd); if (canonicalFd !== undefined) closeSync(canonicalFd); + if (manifestFd !== undefined) closeSync(manifestFd); } let stdout = ""; let stderr = ""; diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index e017493a..7d480305 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -3,8 +3,7 @@ import { spawnSync } from "node:child_process"; import { isIP } from "node:net"; import { domainToASCII } from "node:url"; import { - closeSync, constants as fsConstants, fstatSync, lstatSync, - existsSync, mkdirSync, openSync, readFileSync, + existsSync, lstatSync, readFileSync, } from "node:fs"; import { dirname, isAbsolute, join, relative, resolve } from "node:path"; import { parseAllDocuments } from "yaml"; @@ -20,6 +19,8 @@ import { parseWorkspaceYaml, validateOperationalWorkspace, type WorkspaceDescrip export interface RuntimeConfigLease { path: string; manifestPath: string; + /** SHA-256 of the exact durable manifest bytes handed to the child. */ + manifestSha256: string; workspaceId: string; workspaceRevision: string; release(): void; @@ -53,12 +54,13 @@ interface SnapshotIdentity { interface PublishedIdentity { path: string; manifestPath: string; + /** SHA-256 of the exact durable manifest bytes handed to the child. */ + manifestSha256: string; workspaceId: string; workspaceRevision: string; digest: string; content: string; manifest: string; - refs: number; } @@ -119,11 +121,12 @@ export class WorkspaceRuntimeConfigLeaseFactory { private readonly env: NodeJS.ProcessEnv; private readonly secretRoots: readonly string[]; private readonly installation: RuntimeInstallationOverlay; - 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"); - if (!existsSync(input.runtimeSnapshotRoot)) mkdirSync(input.runtimeSnapshotRoot, { recursive: true, mode: 0o700 }); + // Registry snapshots are produced by WorkspaceRegistry. Never recursively + // create this security boundary from a pathname (an ancestor could be swapped). + if (!existsSync(input.runtimeSnapshotRoot)) throw new Error("runtime snapshot root is unavailable"); this.assertDirectory(input.runtimeSnapshotRoot, "runtime snapshot root"); this.env = { ...(input.env ?? process.env) }; this.secretRoots = [...(input.secretRoots ?? [])]; @@ -149,7 +152,7 @@ export class WorkspaceRuntimeConfigLeaseFactory { const renderedDigest = digest(rendered); const base = { workspace_id: snapshot.workspaceId, workspace_revision: snapshot.workspaceRevision, - descriptor_git_blob: snapshot.descriptorBlob ?? "unknown", + descriptor_git_blob: snapshot.descriptorBlob!, descriptor_sha256: snapshot.digest, config_sha256: renderedDigest, config_dwh_binding: this.computeBinding(rendered), @@ -158,9 +161,8 @@ export class WorkspaceRuntimeConfigLeaseFactory { const identity: PublishedIdentity = { path: result.path, manifestPath: result.manifestPath, workspaceId: snapshot.workspaceId, workspaceRevision: snapshot.workspaceRevision, digest: renderedDigest, content: rendered, - manifest: result.manifest, refs: 1, + manifest: result.manifest, manifestSha256: result.manifest_sha256, }; - this.published.set(identity.path, identity); return this.lease(identity); } @@ -194,14 +196,14 @@ export class WorkspaceRuntimeConfigLeaseFactory { 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} { + private publishSecure(workspaceId: string, revision: string, content: string, manifestBase: Record): {path:string; manifestPath:string; manifest:string; manifest_sha256: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, + return { path: identity.path, manifestPath: identity.manifestPath, manifestSha256: identity.manifestSha256, workspaceId: identity.workspaceId, workspaceRevision: identity.workspaceRevision, release: () => { if (released) return; released = true; /* Durable revision-owned state: release only drops our local handle/ref. */ }, }; } @@ -217,31 +219,30 @@ export class WorkspaceRuntimeConfigLeaseFactory { private readSnapshot(path: string): SnapshotIdentity { if (!isAbsolute(path)) throw new Error("workspace snapshot path must be absolute"); - const rel = relative(this.input.runtimeSnapshotRoot, path); + const root = resolve(this.input.runtimeSnapshotRoot); + const rel = relative(root, 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 || (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"); - 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); } + // The helper is the canonical registry capability boundary. It opens the exact + // production snapshot.json and descriptor component-by-component, and MUST prove + // the Git commit/blob identity; a pathname-shaped file is never sufficient. + const repositoryRoot = join(dirname(root), "repo"); + const verified = this.helper("verified-snapshot", { + snapshots_root: root, repository_root: repositoryRoot, + workspace_revision: match[1], workspace_id: match[2], + }); + if (!verified || typeof verified.source !== "string" + || verified.sha256 !== digest(verified.source) || verified.snapshot_path !== path) { + throw new Error("workspace snapshot integrity check failed"); + } + const workspace = validateOperationalWorkspace(parseWorkspaceYaml(verified.source)); + if (workspace.workspace.id !== match[2] || typeof verified.descriptor_git_blob !== "string") { + throw new Error("workspace snapshot integrity check failed"); + } + return { + workspace, workspaceId: match[2], workspaceRevision: match[1], + revisionContentRoot: join(root, match[1]), digest: verified.sha256, + descriptorBlob: verified.descriptor_git_blob, + }; } } diff --git a/backend/test/workspace-runtime-config-lease.test.ts b/backend/test/workspace-runtime-config-lease.test.ts index a93365c5..b414f992 100644 --- a/backend/test/workspace-runtime-config-lease.test.ts +++ b/backend/test/workspace-runtime-config-lease.test.ts @@ -1,10 +1,11 @@ import { test, expect } from "vitest"; import { chmodSync, existsSync, lstatSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; -import { dirname, join } from "node:path"; +import { join } from "node:path"; +import { execFileSync } from "node:child_process"; +import { createHash } from "node:crypto"; import { WorkspaceRuntimeConfigLeaseFactory } from "../src/workspaces/runtime-config-lease.js"; -const commit = "a".repeat(40); const workspace = "abc"; const descriptor = `workspace: schema_version: 3 @@ -33,11 +34,27 @@ llm_policy: function fixture() { const root = mkdtempSync(join(tmpdir(), "runtime-config-lease-")); const snapshots = join(root, "snapshots"); - const snapshotPath = join(snapshots, commit, `${workspace}.yaml`); + const repo = join(root, "repo"); + mkdirSync(join(repo, "workspaces"), { recursive: true }); + execFileSync("git", ["init", "--initial-branch=main"], { cwd: repo }); + execFileSync("git", ["config", "user.name", "Fixture"], { cwd: repo }); + execFileSync("git", ["config", "user.email", "fixture@example.invalid"], { cwd: repo }); + writeFileSync(join(repo, "workspaces", `${workspace}.yaml`), descriptor); + execFileSync("git", ["add", "."], { cwd: repo }); + execFileSync("git", ["commit", "-m", "fixture"], { cwd: repo }); + const actualCommit = execFileSync("git", ["rev-parse", "HEAD"], { cwd: repo, encoding: "utf8" }).trim(); + const blob = execFileSync("git", ["rev-parse", `HEAD:workspaces/${workspace}.yaml`], { cwd: repo, encoding: "utf8" }).trim(); + const snapshotsDir = join(snapshots, actualCommit); + const snapshotPath = join(snapshotsDir, `${workspace}.yaml`); const dataRoot = join(root, "data"); const harness = join(root, "harness"); - mkdirSync(join(snapshots, commit), { recursive: true, mode: 0o700 }); + mkdirSync(snapshotsDir, { recursive: true, mode: 0o700 }); chmodSync(snapshots, 0o700); + writeFileSync(join(snapshotsDir, "snapshot.json"), JSON.stringify({ + head: actualCommit, + revisions: [{ id: workspace, commit: actualCommit, blob, snapshotPath }], + files: { [`${workspace}.yaml`]: createHash("sha256").update(descriptor).digest("hex") }, + }), { mode: 0o400 }); mkdirSync(harness); writeFileSync(snapshotPath, descriptor, { mode: 0o400 }); const secret = join(root, "password"); @@ -66,9 +83,9 @@ test("session and maintenance share deterministic bytes and path", () => { expect(session.path).toBe(maintenance.path); expect(readFileSync(session.path, "utf8")).toBe(readFileSync(maintenance.path, "utf8")); expect(lstatSync(session.path).mode & 0o777).toBe(0o400); - expect(existsSync(join(dirname(session.path), `${commit}.manifest.json`))).toBe(true); + expect(existsSync(maintenance.manifestPath)).toBe(true); session.release(); maintenance.release(); - expect(existsSync(session.path)).toBe(false); + expect(existsSync(session.path)).toBe(true); } finally { rmSync(f.root, { recursive: true, force: true }); } }); diff --git a/backend/test/workspace-runtime-handoff.test.ts b/backend/test/workspace-runtime-handoff.test.ts index e5f98c3a..b407c651 100644 --- a/backend/test/workspace-runtime-handoff.test.ts +++ b/backend/test/workspace-runtime-handoff.test.ts @@ -101,7 +101,8 @@ async function fixture(workspaceSource = filesystemWorkspace) { writeFileSync(path, contents, { mode: 0o600 }); chmodSync(path, 0o600); } - mkdirSync(dataRoot); + mkdirSync(dataRoot, { mode: 0o700 }); + chmodSync(dataRoot, 0o700); const registryConfig: WorkspaceRegistryConfig = { root: registryRoot, remoteUrl: remote, @@ -138,7 +139,7 @@ function runnerFor(f: Awaited>): ThtRunner { harnessDir, configPath: "config/tht.yaml", dataRoot: f.dataRoot, - runtimeSnapshotRoot: join(f.registryConfig.root, "snapshots", "runtime"), + runtimeSnapshotRoot: join(f.registryConfig.root, "snapshots"), secretRoots: f.registryConfig.secretRoots, } as any); } @@ -162,7 +163,7 @@ test("real schema-v3 registry revision loads through ThtRunner and the harness c expect(existsSync(join( f.dataRoot, "sessions", "psd-clinical", "sessions", created.id, "session_manifest.yaml", ))).toBe(true); - expect(readdirSync(join(f.registryConfig.root, "snapshots", "runtime"))).toEqual([]); + expect(readdirSync(join(f.registryConfig.root, "snapshots"))).toContain(f.revision.commit); }); test("separate runtime leases hand off byte-identical revision Evidence configs accepted by tht", async () => { @@ -215,7 +216,7 @@ test("separate runtime leases hand off byte-identical revision Evidence configs expect(existsSync(first.path)).toBe(true); expect(existsSync(second.path)).toBe(true); second.release(); - expect(existsSync(second.path)).toBe(false); + expect(existsSync(second.path)).toBe(true); } finally { first.release(); second.release(); diff --git a/harness/tht/config.py b/harness/tht/config.py index 8a5fcff5..1973fab3 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -1,3 +1,4 @@ +import hashlib import json import os import re @@ -572,25 +573,80 @@ def _validate_raw_config_shape(raw: dict[str, Any], path: Path) -> None: ) +def _read_runtime_fd(fd: int, label: str, expected_mode: int = 0o400) -> tuple[bytes, os.stat_result]: + try: + info = os.fstat(fd) + if (not stat.S_ISREG(info.st_mode) or info.st_nlink != 1 + or stat.S_IMODE(info.st_mode) != expected_mode + or info.st_uid != os.getuid()): + raise OSError("unsafe runtime descriptor") + os.lseek(fd, 0, os.SEEK_SET) + chunks: list[bytes] = [] + total = 0 + while chunk := os.read(fd, 1024 * 1024): + total += len(chunk) + if total > 16 * 1024 * 1024: + raise OSError("runtime descriptor too large") + chunks.append(chunk) + return b"".join(chunks), info + except OSError as exc: + raise ConfigError(f"File runtime {label} non attendibile") from exc + + +def _strict_runtime_manifest(raw: object) -> dict[str, object]: + required = { + "version", "workspace_id", "workspace_revision", "descriptor_git_blob", + "descriptor_sha256", "config_sha256", "config_dwh_binding", "config_dev", + "config_ino", "config_size", "config_mode", "config_uid", "config_nlink", + } + if not isinstance(raw, dict) or set(raw) != required or raw.get("version") != 1: + raise ConfigError("Manifest runtime non valido") + if not isinstance(raw.get("config_dwh_binding"), dict): + raise ConfigError("Manifest runtime non valido") + binding = raw["config_dwh_binding"] + if set(binding) != {"workspace_id", "config_fingerprint", "input_fingerprint"} or any(not isinstance(v, str) for v in binding.values()): + raise ConfigError("Manifest runtime non valido") + for key in ("descriptor_sha256", "config_sha256"): + if not isinstance(raw[key], str) or not re.fullmatch(r"[0-9a-f]{64}", raw[key]): + raise ConfigError("Manifest runtime non valido") + for key in ("config_dev", "config_ino", "config_size", "config_uid", "config_nlink"): + if not isinstance(raw[key], str) or not raw[key].isdigit(): + raise ConfigError("Manifest runtime non valido") + if raw["config_mode"] != "400": + raise ConfigError("Manifest runtime non valido") + return raw + + def load_config(path: Path) -> Config: # 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_manifest: dict[str, object] | None = None runtime_fd = os.environ.get("THT_CONFIG_FD") - if runtime_fd is not None: + manifest_fd = os.environ.get("THT_CONFIG_MANIFEST_FD") + expected_manifest = os.environ.get("THT_CONFIG_MANIFEST_SHA256") + if runtime_fd is not None or manifest_fd is not None or expected_manifest is not None: + if runtime_fd is None or manifest_fd is None or expected_manifest is None or not re.fullmatch(r"[0-9a-f]{64}", expected_manifest): + raise ConfigError("Handoff runtime incompleto") 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: + config_bytes, config_info = _read_runtime_fd(int(runtime_fd), "config") + manifest_bytes, manifest_info = _read_runtime_fd(int(manifest_fd), "manifest", 0o600) + if hashlib.sha256(manifest_bytes).hexdigest() != expected_manifest: + raise ConfigError("Manifest runtime modificato") + runtime_manifest = _strict_runtime_manifest(json.loads(manifest_bytes.decode("utf-8"))) + if (runtime_manifest["config_sha256"] != hashlib.sha256(config_bytes).hexdigest() + or int(runtime_manifest["config_dev"]) != config_info.st_dev + or int(runtime_manifest["config_ino"]) != config_info.st_ino + or int(runtime_manifest["config_size"]) != config_info.st_size + or int(runtime_manifest["config_uid"]) != config_info.st_uid + or int(runtime_manifest["config_nlink"]) != config_info.st_nlink + or config_info.st_dev == manifest_info.st_dev and config_info.st_ino == manifest_info.st_ino): + raise ConfigError("Identità config runtime non valida") + source_text = config_bytes.decode("utf-8") + except (OSError, UnicodeError, ValueError, json.JSONDecodeError) as exc: + if isinstance(exc, ConfigError): + raise raise ConfigError("File di configurazione runtime non attendibile") from exc else: if not path.exists(): @@ -679,6 +735,12 @@ def load_config(path: Path) -> Config: ) _validate_active_embeddings_config(cfg.embeddings, path) _validate_active_vector_config(cfg.vectors, path) + if runtime_manifest is not None: + from tht.jobs.dwh_pipeline import config_dwh_binding + if runtime_manifest["workspace_id"] != cfg._workspace_id or runtime_manifest["workspace_revision"] != cfg._workspace_revision: + raise ConfigError("Identità workspace runtime non valida") + if config_dwh_binding(cfg) != runtime_manifest["config_dwh_binding"]: + raise ConfigError("Binding DWH runtime modificato") return cfg diff --git a/harness/tht/runtime_config_lease_io.py b/harness/tht/runtime_config_lease_io.py index ae1806dc..424622ea 100644 --- a/harness/tht/runtime_config_lease_io.py +++ b/harness/tht/runtime_config_lease_io.py @@ -26,8 +26,16 @@ def safe_rev(v: str) -> bool: def open_dir(parent: int | None, name: str, create: bool = False) -> int: - flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | os.O_NOFOLLOW + # Darwin rejects O_NOFOLLOW|openat for directories (ELOOP); lstat the + # component before opening and verify the resulting descriptor below. Linux + # uses the stronger flag where available. + flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) + if sys.platform != "darwin": + flags |= os.O_NOFOLLOW try: + entry = os.stat(name, dir_fd=parent, follow_symlinks=False) + if stat.S_ISLNK(entry.st_mode): + fail("runtime config directory is not trusted") return os.open(name, flags, dir_fd=parent) except FileNotFoundError: if not create: @@ -48,25 +56,39 @@ def checked_dir(fd: int, expected_mode: int = 0o700) -> None: def walk(root: str, comps: list[str], create: bool = True) -> int: + """Open an absolute path component-by-component without following symlinks. + + In particular, never use os.makedirs/root pathname resolution here: an attacker + replacing an ancestor between those calls must not redirect publication. + """ 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) + # macOS exposes temporary directories through the conventional /var and + # /tmp symlinks. Resolve only these OS-owned aliases; workspace-owned + # ancestors remain component checked and are never realpath-followed. + if root == "/var" or root == "/tmp" or root.startswith(("/var/", "/tmp/")): + root = "/private" + root + parts = [part for part in Path(root).parts if part not in ("", "/")] + if any(part in (".", "..") or "/" in part for part in parts + comps): + fail("unsafe path component") + fd = os.open("/", os.O_RDONLY | getattr(os, "O_DIRECTORY", 0)) try: - checked_dir(fd) - except: + all_components = [*parts, *comps] + for index, component in enumerate(all_components): + nxt = open_dir(fd, component, create) + # Ancestors such as /var/folders are installation-owned and commonly + # 0755; the trusted runtime root and every workspace child are private. + info = os.fstat(nxt) + if (not stat.S_ISDIR(info.st_mode) or info.st_nlink < 1 + or (index >= len(parts) - 1 and (info.st_uid != os.getuid() or stat.S_IMODE(info.st_mode) != 0o700))): + os.close(nxt) + fail("runtime config directory is not trusted") + os.close(fd) + fd = nxt + return fd + except BaseException: 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: @@ -100,6 +122,48 @@ def write_all(fd: int, data: bytes) -> None: pos += n +def read_all(fd: int, limit: int = 16 * 1024 * 1024) -> bytes: + os.lseek(fd, 0, os.SEEK_SET) + chunks: list[bytes] = [] + total = 0 + while True: + chunk = os.read(fd, min(1024 * 1024, limit - total)) + if not chunk: + return b"".join(chunks) + chunks.append(chunk) + total += len(chunk) + if total > limit: + fail("runtime config file is too large") + + +def strict_manifest(value: object) -> dict: + if not isinstance(value, dict): + fail("runtime config manifest is invalid") + required = { + "version", "workspace_id", "workspace_revision", "descriptor_git_blob", + "descriptor_sha256", "config_sha256", "config_dwh_binding", "config_dev", + "config_ino", "config_size", "config_mode", "config_uid", "config_nlink", + } + if set(value) != required or value.get("version") != 1: + fail("runtime config manifest is invalid") + if not safe_id(value.get("workspace_id")) or not safe_rev(value.get("workspace_revision")): + fail("runtime config manifest is invalid") + if not isinstance(value.get("descriptor_git_blob"), str) or not safe_rev(value["descriptor_git_blob"]): + fail("runtime config manifest is invalid") + for key in ("descriptor_sha256", "config_sha256"): + if not isinstance(value[key], str) or not __import__("re").fullmatch(r"[0-9a-f]{64}", value[key]): + fail("runtime config manifest is invalid") + binding_value = value.get("config_dwh_binding") + if not isinstance(binding_value, dict) or set(binding_value) != {"workspace_id", "config_fingerprint", "input_fingerprint"} or any(not isinstance(x, str) for x in binding_value.values()): + fail("runtime config manifest is invalid") + for key in ("config_dev", "config_ino", "config_size", "config_uid", "config_nlink"): + if not isinstance(value[key], str) or not value[key].isdigit(): + fail("runtime config manifest is invalid") + if value["config_mode"] != "400": + fail("runtime config manifest is invalid") + return value + + def publish(inp: dict) -> dict: root = inp.get("data_root") wid = inp.get("workspace_id") @@ -160,7 +224,7 @@ def publish(inp: dict) -> dict: if got: fd, s = got os.lseek(fd, 0, os.SEEK_SET) - old = os.read(fd, len(content) + 1) + old = read_all(fd) os.close(fd) if old != content: fail("same-revision runtime configuration changed") @@ -179,6 +243,10 @@ def publish(inp: dict) -> dict: ) except FileExistsError: pass + # Keep metadata ordering explicit even on filesystems where a + # hardlink publication does not retain fchmod as expected. + os.fchmod(fd, 0o400) + os.fsync(fd) finally: os.close(fd) try: @@ -191,7 +259,7 @@ def publish(inp: dict) -> dict: fd, s = got try: os.lseek(fd, 0, os.SEEK_SET) - if os.read(fd, len(content) + 1) != content: + if read_all(fd) != content: fail("same-revision runtime configuration changed") finally: os.close(fd) @@ -201,7 +269,7 @@ def publish(inp: dict) -> dict: assert got fd, s = got os.close(fd) - manifest = dict(base) + manifest = {"version": 1, **dict(base)} manifest.update( { "config_sha256": hashlib.sha256(content).hexdigest(), @@ -213,13 +281,16 @@ def publish(inp: dict) -> dict: "config_nlink": str(s.st_nlink), } ) + strict_manifest(manifest) 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) + existing = read_all(mfd) os.close(mfd) + try: strict_manifest(json.loads(existing.decode())) + except (ValueError, TypeError, UnicodeError, RuntimeError): fail("runtime config manifest is invalid") if existing != mb: fail("same-revision runtime configuration changed") else: @@ -245,6 +316,7 @@ def publish(inp: dict) -> dict: "path": f"{root}/sessions/{wid}/preprocessing/runtime-config/{name}", "manifestPath": f"{root}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}", "manifest": mb.decode(), + "manifest_sha256": hashlib.sha256(mb).hexdigest(), "dev": s.st_dev, "ino": s.st_ino, } @@ -295,35 +367,53 @@ def verified_snapshot(inp: dict) -> dict: 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) + try: + manifest = json.loads(payload.decode()) + except (UnicodeDecodeError, json.JSONDecodeError): + fail("workspace snapshot integrity check failed") + if not isinstance(manifest, dict) or set(manifest) != {"head", "revisions", "files"}: + fail("workspace snapshot integrity check failed") + records = manifest.get("revisions") + files = manifest.get("files") + record = next((r for r in records if isinstance(r, dict) and r.get("id") == wid), None) if isinstance(records, list) else None + expected_path = f"{root}/{rev}/{wid}.yaml" 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() + manifest.get("head") != rev or not isinstance(records, list) or not record + or set(record) != {"id", "commit", "blob", "snapshotPath"} + or record.get("commit") != rev or record.get("snapshotPath") != expected_path + or not isinstance(record.get("blob"), str) or not safe_rev(record.get("blob")) + or not isinstance(files, dict) + or 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") + 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") + try: + import re + from collections import Counter + normalize = lambda value: re.findall(r"[A-Za-z0-9_.:/@+-]+", value) + git_tokens = Counter(normalize(git_source.decode("utf-8"))) + snapshot_tokens = Counter(normalize(source.decode("utf-8"))) + # The registry canonicalizer may add schema defaults/reorder mappings. + # Every token from the exact Git descriptor must nevertheless survive; + # replacements (including non-rendered workspace.name) are rejected. + equivalent = all(snapshot_tokens[k] >= count for k, count in git_tokens.items()) + except UnicodeDecodeError: + equivalent = False + if blob != record.get("blob") or not equivalent: + fail("workspace Git descriptor identity mismatch") return { "workspace_id": wid, "workspace_revision": rev,