From 11cc8628cf341221202c62786da09e88756a5577 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 13:03:37 +0200 Subject: [PATCH] fix: harden runtime config helper protocol and lifecycle --- .../src/workspaces/runtime-config-lease.ts | 189 +++++++++++++++--- .../workspace-runtime-config-lease.test.ts | 22 +- harness/tests/test_runtime_config_lease_io.py | 101 ++++++++++ harness/tht/runtime_config_lease_io.py | 12 +- 4 files changed, 290 insertions(+), 34 deletions(-) diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index d067a06c..ea48f8cc 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -54,6 +54,10 @@ interface SnapshotIdentity { descriptorIno: string; } interface PublishedResponse { + protocol_version: 1; + kind: "publication"; + workspace_id: string; + workspace_revision: string; path: string; manifestPath: string; manifest: string; @@ -129,6 +133,120 @@ function parseInstallationOverlay(path: string): RuntimeInstallationOverlay { function digest(data: string | Buffer): string { return createHash("sha256").update(data).digest("hex"); } + +export interface BoundedHelperOptions { + cwd: string; + env: NodeJS.ProcessEnv; + action: string; + payload: string; + timeoutMs?: number; + outputLimit?: number; + /** Package-internal seam used by lifecycle tests; production uses node spawn. */ + spawnProcess?: typeof spawn; +} + +/** Run one JSON helper with bounded output and deterministic child cleanup. */ +export async function runBoundedHelper( + executable: string, + args: readonly string[], + options: BoundedHelperOptions, +): Promise { + const timeoutMs = options.timeoutMs ?? 10_000; + const outputLimit = options.outputLimit ?? 16 * 1024 * 1024; + const spawnProcess = options.spawnProcess ?? spawn; + const child = spawnProcess(executable, [...args], { + cwd: options.cwd, + env: options.env, + detached: true, + stdio: ["pipe", "pipe", "pipe"], + }); + let stdout = ""; + let stderr = ""; + let closed = false; + let finishing = false; + let timer: ReturnType | undefined; + let closeResolve: () => void = () => undefined; + const closePromise = new Promise((resolveClose) => { closeResolve = resolveClose; }); + + const killGroup = () => { + try { + if (child.pid !== undefined && child.pid !== null) process.kill(-child.pid, "SIGKILL"); + } catch { + try { child.kill("SIGKILL"); } catch { /* process already gone */ } + } + }; + const destroyStreams = () => { + try { child.stdin?.destroy(); } catch { /* already closed */ } + try { child.stdout?.destroy(); } catch { /* already closed */ } + try { child.stderr?.destroy(); } catch { /* already closed */ } + }; + const awaitClose = async () => { + if (closed) return; + // A descendant can keep stdio open even after the direct child exits. Streams + // are destroyed before this bounded reap wait so the backend cannot hang. + await Promise.race([closePromise, new Promise((resolveWait) => setTimeout(resolveWait, 1_000))]); + }; + + return new Promise((resolveResult, rejectResult) => { + const finish = async (error?: Error, value?: unknown, terminate = false) => { + if (finishing) return; + finishing = true; + if (timer !== undefined) clearTimeout(timer); + if (terminate) killGroup(); + destroyStreams(); + await awaitClose(); + destroyStreams(); + if (error) rejectResult(error); else resolveResult(value); + }; + const append = (target: "stdout" | "stderr", data: Buffer) => { + const next = target === "stdout" ? stdout + data.toString() : stderr + data.toString(); + if (Buffer.byteLength(next) > outputLimit) { + void finish(new Error(`runtime config ${options.action} output exceeded limit`), undefined, true); + return; + } + if (target === "stdout") stdout = next; else stderr = next; + }; + + // Every stream gets an error listener before any data is written. In + // particular, EPIPE from end() must become the same bounded failure path. + child.stdin?.on("error", (error) => void finish(error instanceof Error ? error : new Error(String(error)), undefined, true)); + child.stdout?.on("error", (error) => void finish(error instanceof Error ? error : new Error(String(error)), undefined, true)); + child.stderr?.on("error", (error) => void finish(error instanceof Error ? error : new Error(String(error)), undefined, true)); + child.stdout?.on("data", (data: Buffer) => append("stdout", data)); + child.stderr?.on("data", (data: Buffer) => append("stderr", data)); + child.on("error", (error) => void finish(error instanceof Error ? error : new Error(String(error)), undefined, true)); + child.on("close", (code) => { + closed = true; + closeResolve(); + if (finishing) return; + if (code !== 0) { + let detail = stderr.trim() || stdout.trim() || `runtime config ${options.action} failed`; + try { detail = (JSON.parse(stdout) as { error?: string }).error ?? detail; } catch { /* preserve detail */ } + void finish(new Error(detail)); + return; + } + try { void finish(undefined, JSON.parse(stdout)); } + catch { void finish(new Error(`runtime config ${options.action} returned invalid JSON`)); } + }); + timer = setTimeout(() => { + void finish(new Error(`runtime config ${options.action} timed out`), undefined, true); + }, timeoutMs); + try { + // The listener above is intentionally installed before end(), since a + // helper may close its input immediately and emit EPIPE synchronously. + child.stdin?.end(options.payload); + } catch (error) { + void finish(error instanceof Error ? error : new Error(String(error)), undefined, true); + } + }); +} + +function sameBinding(left: unknown, right: unknown): boolean { + if (!left || typeof left !== "object" || Array.isArray(left) || !right || typeof right !== "object" || Array.isArray(right)) return false; + const a = left as Record; const b = right as Record; + return a.workspace_id === b.workspace_id && a.config_fingerprint === b.config_fingerprint && a.input_fingerprint === b.input_fingerprint; +} + export class WorkspaceRuntimeConfigLeaseFactory { private readonly env: NodeJS.ProcessEnv; private readonly secretRoots: readonly string[]; @@ -186,40 +304,32 @@ export class WorkspaceRuntimeConfigLeaseFactory { const executable = existsSync(python) ? python : existsSync(projectPython) ? projectPython : (process.env.PYTHON ?? "python3"); const helperArgs = existsSync(modulePath) ? [modulePath] : ["-m", "tht.runtime_config_lease_io"]; const env = { ...this.env }; - // Never pass ambient capability variables to binding/snapshot/publication. + // Capability variables are never inherited. Fault seams are only available in + // tests, and are copied explicitly rather than forwarding ambient state. for (const key of Object.keys(env)) { - if ((key.startsWith("THT_RUNTIME_CONFIG_") && !["THT_RUNTIME_CONFIG_FSYNC_FAIL", "THT_RUNTIME_CONFIG_RENAME_FAIL"].includes(key)) || key.startsWith("THT_CONFIG_")) delete env[key]; + if (key.startsWith("THT_RUNTIME_CONFIG_") || key.startsWith("THT_CONFIG_")) delete env[key]; + } + if (env.NODE_ENV === "test") { + for (const key of ["THT_RUNTIME_CONFIG_FSYNC_FAIL", "THT_RUNTIME_CONFIG_RENAME_FAIL"]) { + const value = this.env[key]; + if (value !== undefined) env[key] = value; + } } env.PYTHONPATH = [this.input.harnessDir, dirname(dirname(modulePath)), env.PYTHONPATH].filter(Boolean).join(":"); const payload = JSON.stringify({ protocol_version: 1, action, ...extra }); - const timeoutMs = 10_000; - return await new Promise((resolveResult, reject) => { - const child = spawn(executable, helperArgs, { cwd: this.input.harnessDir, env, detached: true, stdio: ["pipe", "pipe", "pipe"] }); - let stdout = ""; let stderr = ""; let settled = false; - const finish = (error?: Error, value?: unknown) => { if (settled) return; settled = true; clearTimeout(timer); error ? reject(error) : resolveResult(value); }; - const kill = () => { try { process.kill(-child.pid!, "SIGKILL"); } catch { try { child.kill("SIGKILL"); } catch { /* gone */ } } }; - const timer = setTimeout(() => { kill(); finish(new Error(`runtime config ${action} timed out`)); }, timeoutMs); - const append = (target: "stdout" | "stderr", data: Buffer) => { - const next = target === "stdout" ? stdout + data.toString() : stderr + data.toString(); - if (next.length > 16 * 1024 * 1024) { kill(); finish(new Error(`runtime config ${action} output exceeded limit`)); return; } - if (target === "stdout") stdout = next; else stderr = next; - }; - child.stdout.on("data", (d: Buffer) => append("stdout", d)); child.stderr.on("data", (d: Buffer) => append("stderr", d)); - child.on("error", (error) => finish(error)); - child.on("close", (code) => { - if (code !== 0) { let detail = stderr.trim() || stdout.trim() || `runtime config ${action} failed`; try { detail = (JSON.parse(stdout) as {error?: string}).error ?? detail; } catch { /* preserve detail */ } finish(new Error(detail)); return; } - try { finish(undefined, JSON.parse(stdout)); } catch { finish(new Error(`runtime config ${action} returned invalid JSON`)); } - }); - child.stdin.end(payload); + return runBoundedHelper(executable, helperArgs, { + cwd: this.input.harnessDir, env, action, payload, }); } private async computeBinding(content: string): Promise> { const value = await this.helper("binding", { config_hex: Buffer.from(content).toString("hex") }); - if (!value || typeof value !== "object" || Array.isArray(value) || Object.keys(value).sort().join(",") !== "config_fingerprint,input_fingerprint,workspace_id") throw new Error("runtime config binding helper returned malformed output"); + if (!value || typeof value !== "object" || Array.isArray(value) || Object.keys(value).sort().join(",") !== "config_fingerprint,input_fingerprint,kind,protocol_version,workspace_id") throw new Error("runtime config binding helper returned malformed output"); const record = value as Record; - if (Object.values(record).some((v) => typeof v !== "string")) throw new Error("runtime config binding helper returned malformed output"); - return record as Record; + if (record.protocol_version !== 1 || record.kind !== "binding" + || Object.values(record).some((v) => typeof v !== "string" && typeof v !== "number")) throw new Error("runtime config binding helper returned malformed output"); + if (typeof record.workspace_id !== "string" || typeof record.config_fingerprint !== "string" || typeof record.input_fingerprint !== "string") throw new Error("runtime config binding helper returned malformed output"); + return { workspace_id: record.workspace_id, config_fingerprint: record.config_fingerprint, input_fingerprint: record.input_fingerprint }; } private async publishSecure(workspaceId: string, revision: string, content: string, manifestBase: Record): Promise { @@ -227,17 +337,33 @@ export class WorkspaceRuntimeConfigLeaseFactory { workspace_revision: revision, config_hex: Buffer.from(content).toString("hex"), manifest_base: manifestBase }); if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("runtime config publish helper returned malformed output"); const result = value as Record; - if (Object.keys(result).sort().join(",") !== "dev,ino,manifest,manifestPath,manifest_sha256,path") throw new Error("runtime config publish helper returned malformed output"); + if (Object.keys(result).sort().join(",") !== "dev,ino,kind,manifest,manifestPath,manifest_sha256,path,protocol_version,workspace_id,workspace_revision") throw new Error("runtime config publish helper returned malformed output"); let expectedRoot = resolve(this.input.dataRoot); if (process.platform === "darwin" && (expectedRoot === "/var" || expectedRoot.startsWith("/var/") || expectedRoot === "/tmp" || expectedRoot.startsWith("/tmp/"))) expectedRoot = `/private${expectedRoot}`; const expectedPath = join(expectedRoot, "sessions", workspaceId, "preprocessing", "runtime-config", `${revision}.yaml`); const expectedManifestPath = join(expectedRoot, "sessions", workspaceId, "preprocessing", "runtime-config-manifests", `${revision}.json`); - if (result.path !== expectedPath || result.manifestPath !== expectedManifestPath - || typeof result.path !== "string" || !isAbsolute(result.path) || normalize(result.path) !== result.path - || typeof result.manifestPath !== "string" || !isAbsolute(result.manifestPath) || normalize(result.manifestPath) !== result.manifestPath - || result.workspace_id !== undefined || typeof result.manifest !== "string" - || typeof result.manifest_sha256 !== "string" || !/^[0-9a-f]{64}$/.test(result.manifest_sha256) + if (result.protocol_version !== 1 || result.kind !== "publication" + || result.workspace_id !== workspaceId || result.workspace_revision !== revision + || typeof result.path !== "string" || result.path !== expectedPath || !isAbsolute(result.path) || normalize(result.path) !== result.path + || typeof result.manifestPath !== "string" || result.manifestPath !== expectedManifestPath || !isAbsolute(result.manifestPath) || normalize(result.manifestPath) !== result.manifestPath + || typeof result.manifest !== "string" || typeof result.manifest_sha256 !== "string" || !/^[0-9a-f]{64}$/.test(result.manifest_sha256) + || digest(result.manifest) !== result.manifest_sha256 || typeof result.dev !== "number" || !Number.isSafeInteger(result.dev) || typeof result.ino !== "number" || !Number.isSafeInteger(result.ino)) throw new Error("runtime config publish helper returned malformed output"); + let manifest: unknown; + try { manifest = JSON.parse(result.manifest); } catch { throw new Error("runtime config publish helper returned malformed output"); } + if (!manifest || typeof manifest !== "object" || Array.isArray(manifest)) throw new Error("runtime config publish helper returned malformed output"); + const manifestRecord = manifest as Record; + if (manifestRecord.version !== 1 + || manifestRecord.workspace_id !== workspaceId || manifestRecord.workspace_revision !== revision + || manifestRecord.descriptor_git_blob !== manifestBase.descriptor_git_blob + || manifestRecord.descriptor_sha256 !== manifestBase.descriptor_sha256 + || manifestRecord.descriptor_dev !== manifestBase.descriptor_dev + || manifestRecord.descriptor_ino !== manifestBase.descriptor_ino + || manifestRecord.config_sha256 !== digest(content) + || !sameBinding(manifestRecord.config_dwh_binding, manifestBase.config_dwh_binding) + || manifestRecord.config_dev !== String(result.dev) || manifestRecord.config_ino !== String(result.ino)) { + throw new Error("runtime config publish helper returned malformed output"); + } return result as unknown as PublishedResponse; } @@ -271,8 +397,9 @@ export class WorkspaceRuntimeConfigLeaseFactory { snapshots_root: root, repository_root: repositoryRoot, workspace_revision: match[1], workspace_id: match[2], }) as Record; - const verifiedKeys = ["descriptor_dev", "descriptor_git_blob", "descriptor_ino", "git_source", "sha256", "snapshot_path", "source", "workspace_id", "workspace_revision"]; + const verifiedKeys = ["descriptor_dev", "descriptor_git_blob", "descriptor_ino", "git_source", "kind", "protocol_version", "sha256", "snapshot_path", "source", "workspace_id", "workspace_revision"]; if (Object.keys(verified).sort().join(",") !== verifiedKeys.join(",") + || verified.protocol_version !== 1 || verified.kind !== "verified_snapshot" || verified.workspace_id !== match[2] || verified.workspace_revision !== match[1] || typeof verified.source !== "string" || typeof verified.git_source !== "string" || verified.sha256 !== digest(verified.source) || verified.snapshot_path !== path diff --git a/backend/test/workspace-runtime-config-lease.test.ts b/backend/test/workspace-runtime-config-lease.test.ts index cfb0ef13..cc958915 100644 --- a/backend/test/workspace-runtime-config-lease.test.ts +++ b/backend/test/workspace-runtime-config-lease.test.ts @@ -4,7 +4,7 @@ import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; import { execFileSync } from "node:child_process"; import { createHash } from "node:crypto"; -import { WorkspaceRuntimeConfigLeaseFactory } from "../src/workspaces/runtime-config-lease.js"; +import { runBoundedHelper, WorkspaceRuntimeConfigLeaseFactory } from "../src/workspaces/runtime-config-lease.js"; import { parse } from "yaml"; import { parseWorkspaceYaml, serializeWorkspaceYaml } from "../src/workspaces/schema.js"; @@ -72,6 +72,7 @@ function fixture(extraEnv: Record = {}) { const factoryInput = { dataRoot, runtimeSnapshotRoot: snapshots, harnessDir: harness, configPath, env: { + NODE_ENV: "test", THT_WS_ABC_DWH_TRANSPORT: "postgres_direct", THT_WS_ABC_DWH_HOST: "dwh", THT_WS_ABC_DWH_PORT: "5432", THT_WS_ABC_DWH_USER: "reader", THT_WS_ABC_DWH_PASSWORD_FILE: secret, ...extraEnv, @@ -258,10 +259,12 @@ test("runtime config symlink replacement is refused", async () => { function realHarnessBinding(config: string): Record { const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py"); const python = join(process.cwd(), "..", "harness", ".venv", "bin", "python"); - return JSON.parse(execFileSync(python, [helper], { + const result = JSON.parse(execFileSync(python, [helper], { cwd: join(process.cwd(), "..", "harness"), encoding: "utf8", input: JSON.stringify({ protocol_version: 1, action: "binding", config_hex: Buffer.from(config).toString("hex") }), })); + const { protocol_version: _protocol, kind: _kind, ...binding } = result; + return binding; } test("explicit installation overlay is canonical and has one real harness binding", async () => { @@ -609,3 +612,18 @@ else: expect(output.trim()).toBe("rejected"); } finally { rmSync(f.root, { recursive: true, force: true }); } }); + + +test("bounded helper settles early stdin close and a never-reading child without crashing", async () => { + const options = (action: string, payload: string, timeoutMs = 200) => ({ + cwd: tmpdir(), env: process.env, action, payload, timeoutMs, + }); + await expect(runBoundedHelper(process.execPath, ["-e", "process.stdin.destroy(); setTimeout(() => {}, 1000)"], + options("early-close", "x".repeat(16 * 1024 * 1024)))).rejects.toThrow(); + const started = Date.now(); + const timer = new Promise((resolve) => setTimeout(resolve, 20)); + await expect(runBoundedHelper(process.execPath, ["-e", "setTimeout(() => {}, 10000)"], + options("never-read", "x".repeat(16 * 1024 * 1024), 80))).rejects.toThrow(/timed out|pipe|closed/i); + await timer; + expect(Date.now() - started).toBeLessThan(2_000); +}); diff --git a/harness/tests/test_runtime_config_lease_io.py b/harness/tests/test_runtime_config_lease_io.py index 8170f277..e9e6e868 100644 --- a/harness/tests/test_runtime_config_lease_io.py +++ b/harness/tests/test_runtime_config_lease_io.py @@ -1,10 +1,15 @@ """Focused unit coverage for the privileged runtime publication seam.""" +import io +import json import os +import subprocess +from pathlib import Path import pytest from tht import runtime_config_lease_io as lease_io +from tht.config import ConfigError, _read_runtime_config_source def _manifest() -> dict: @@ -82,3 +87,99 @@ def test_retry_reasserts_parent_durability_before_success(tmp_path, monkeypatch, events.clear() lease_io.publish(inp) assert events.index("config-parent") < events.index("manifest-parent") + + +def test_protocol_success_shapes_and_version_rejection(monkeypatch): + monkeypatch.setattr(lease_io, "binding", lambda _inp: { + "workspace_id": "abc", "config_fingerprint": "f", "input_fingerprint": "i", + }) + monkeypatch.setattr(lease_io.sys, "stdin", io.StringIO( + '{"protocol_version":1,"action":"binding","config_hex":""}' + )) + success = io.StringIO() + monkeypatch.setattr(lease_io.sys, "stdout", success) + lease_io.main() + assert set(json.loads(success.getvalue())) == { + "protocol_version", "kind", "workspace_id", "config_fingerprint", "input_fingerprint", + } + + # An old/new protocol mismatch must be rejected before action dispatch. + monkeypatch.setattr(lease_io.sys, "stdin", io.StringIO( + '{"protocol_version":999,"action":"binding","config_hex":""}' + )) + failure = io.StringIO() + monkeypatch.setattr(lease_io.sys, "stdout", failure) + with pytest.raises(SystemExit): + lease_io.main() + assert json.loads(failure.getvalue())["error"] == "unsupported runtime config protocol" + + +def test_secure_runtime_reader_binds_path_manifest_and_ignores_legacy_fd(monkeypatch, tmp_path): + revision = "a" * 40 + inp = { + "data_root": str(tmp_path / "data"), "workspace_id": "abc", + "workspace_revision": revision, "config_hex": b"runtime: true\n".hex(), + "manifest_base": { + "workspace_id": "abc", "workspace_revision": revision, + "descriptor_git_blob": "b" * 40, "descriptor_sha256": "c" * 64, + "descriptor_dev": "1", "descriptor_ino": "2", + "config_dwh_binding": {"workspace_id": "abc", "config_fingerprint": "e", "input_fingerprint": "f"}, + }, + } + result = lease_io.publish(inp) + path = Path(result["path"]) + monkeypatch.setenv("THT_RUNTIME_CONFIG_MANIFEST_SHA256", result["manifest_sha256"]) + monkeypatch.setenv("THT_CONFIG_MANIFEST_FD", "999") + monkeypatch.setenv("THT_CONFIG_MANIFEST_SHA256", "0" * 64) + source, manifest = _read_runtime_config_source(path) + assert source == "runtime: true\n" + assert manifest is not None and manifest["workspace_id"] == "abc" + + monkeypatch.setenv("THT_RUNTIME_CONFIG_MANIFEST_SHA256", "0" * 64) + with pytest.raises(ConfigError): + _read_runtime_config_source(path) + with pytest.raises(ConfigError): + _read_runtime_config_source(path.with_name("not-canonical.yaml")) + + +@pytest.mark.parametrize("object_kind", ["tree", "blob", "tag"]) +def test_verified_snapshot_rejects_non_commit_object(tmp_path, object_kind): + repo = tmp_path / "repo" + (repo / "workspaces").mkdir(parents=True) + subprocess.run(["git", "init", "--initial-branch=main"], cwd=repo, check=True, stdout=subprocess.DEVNULL) + subprocess.run(["git", "config", "user.name", "Fixture"], cwd=repo, check=True) + subprocess.run(["git", "config", "user.email", "fixture@example.invalid"], cwd=repo, check=True) + descriptor = "workspace:\n schema_version: 3\n id: abc\n" + (repo / "workspaces" / "abc.yaml").write_text(descriptor) + subprocess.run(["git", "add", "."], cwd=repo, check=True) + subprocess.run(["git", "commit", "-m", "fixture"], cwd=repo, check=True, stdout=subprocess.DEVNULL) + if object_kind == "tree": + revision = subprocess.check_output(["git", "rev-parse", "HEAD^{tree}"], cwd=repo, text=True).strip() + elif object_kind == "blob": + revision = subprocess.check_output(["git", "rev-parse", "HEAD:workspaces/abc.yaml"], cwd=repo, text=True).strip() + else: + subprocess.run(["git", "tag", "-a", "v1", "-m", "tag"], cwd=repo, check=True) + revision = subprocess.check_output(["git", "rev-parse", "refs/tags/v1^{tag}"], cwd=repo, text=True).strip() + + snapshots = tmp_path / "snapshots" / revision + snapshots.mkdir(parents=True, mode=0o700) + (tmp_path / "snapshots").chmod(0o700) + files = {"abc.yaml": descriptor, "abc.env.example": "# fixture\n", "abc.md": "# fixture\n"} + for name, content in files.items(): + target = snapshots / name + target.write_text(content) + target.chmod(0o400) + snapshot = { + "head": revision, + "revisions": [{"id": "abc", "commit": revision, "blob": "b" * 40, + "snapshotPath": f"{tmp_path / 'snapshots'}/{revision}/abc.yaml"}], + "files": {name: __import__("hashlib").sha256(content.encode()).hexdigest() for name, content in files.items()}, + } + manifest = snapshots / "snapshot.json" + manifest.write_text(json.dumps(snapshot)) + manifest.chmod(0o400) + with pytest.raises(RuntimeError, match="exact commit|unavailable"): + lease_io.verified_snapshot({ + "snapshots_root": str(tmp_path / "snapshots"), "repository_root": str(repo), + "workspace_revision": revision, "workspace_id": "abc", + }) diff --git a/harness/tht/runtime_config_lease_io.py b/harness/tht/runtime_config_lease_io.py index f7885008..fc6b0dd6 100644 --- a/harness/tht/runtime_config_lease_io.py +++ b/harness/tht/runtime_config_lease_io.py @@ -492,6 +492,10 @@ def publish(inp: dict) -> dict: # failed after the no-replace publication on an earlier invocation. publication_fsync(mandir, "manifest-parent") return { + "protocol_version": 1, + "kind": "publication", + "workspace_id": wid, + "workspace_revision": rev, "path": f"{canonical}/sessions/{wid}/preprocessing/runtime-config/{name}", "manifestPath": f"{canonical}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}", "manifest": mb.decode(), @@ -630,6 +634,8 @@ def verified_snapshot(inp: dict) -> dict: if blob != record.get("blob"): fail("workspace Git descriptor identity mismatch") return { + "protocol_version": 1, + "kind": "verified_snapshot", "workspace_id": wid, "workspace_revision": rev, "source": source.decode(), @@ -663,7 +669,10 @@ def binding(inp: dict) -> dict: stream.write(raw) path = Path(stream.name) try: - return config_dwh_binding(load_config(path)) + result = config_dwh_binding(load_config(path)) + if not isinstance(result, dict) or set(result) != {"workspace_id", "config_fingerprint", "input_fingerprint"}: + fail("runtime config binding returned malformed output") + return result finally: try: path.unlink() @@ -690,6 +699,7 @@ def main() -> None: result = verified_snapshot(inp) else: result = binding(inp) + result = {"protocol_version": 1, "kind": "binding", **result} print(json.dumps(result)) except Exception as e: # noqa: BLE001 print(json.dumps({"error": str(e)}))