fix: harden runtime config helper protocol and lifecycle

This commit is contained in:
2026-08-11 13:03:37 +02:00
parent ec92f7f994
commit 11cc8628cf
4 changed files with 290 additions and 34 deletions
+158 -31
View File
@@ -54,6 +54,10 @@ interface SnapshotIdentity {
descriptorIno: string; descriptorIno: string;
} }
interface PublishedResponse { interface PublishedResponse {
protocol_version: 1;
kind: "publication";
workspace_id: string;
workspace_revision: string;
path: string; path: string;
manifestPath: string; manifestPath: string;
manifest: 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"); } 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<unknown> {
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<typeof setTimeout> | undefined;
let closeResolve: () => void = () => undefined;
const closePromise = new Promise<void>((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<void>((resolveWait) => setTimeout(resolveWait, 1_000))]);
};
return new Promise<unknown>((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<string, unknown>; const b = right as Record<string, unknown>;
return a.workspace_id === b.workspace_id && a.config_fingerprint === b.config_fingerprint && a.input_fingerprint === b.input_fingerprint;
}
export class WorkspaceRuntimeConfigLeaseFactory { export class WorkspaceRuntimeConfigLeaseFactory {
private readonly env: NodeJS.ProcessEnv; private readonly env: NodeJS.ProcessEnv;
private readonly secretRoots: readonly string[]; 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 executable = existsSync(python) ? python : existsSync(projectPython) ? projectPython : (process.env.PYTHON ?? "python3");
const helperArgs = existsSync(modulePath) ? [modulePath] : ["-m", "tht.runtime_config_lease_io"]; const helperArgs = existsSync(modulePath) ? [modulePath] : ["-m", "tht.runtime_config_lease_io"];
const env = { ...this.env }; 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)) { 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(":"); env.PYTHONPATH = [this.input.harnessDir, dirname(dirname(modulePath)), env.PYTHONPATH].filter(Boolean).join(":");
const payload = JSON.stringify({ protocol_version: 1, action, ...extra }); const payload = JSON.stringify({ protocol_version: 1, action, ...extra });
const timeoutMs = 10_000; return runBoundedHelper(executable, helperArgs, {
return await new Promise((resolveResult, reject) => { cwd: this.input.harnessDir, env, action, payload,
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);
}); });
} }
private async computeBinding(content: string): Promise<Record<string, string>> { private async computeBinding(content: string): Promise<Record<string, string>> {
const value = await this.helper("binding", { config_hex: Buffer.from(content).toString("hex") }); 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<string, unknown>; const record = value as Record<string, unknown>;
if (Object.values(record).some((v) => typeof v !== "string")) throw new Error("runtime config binding helper returned malformed output"); if (record.protocol_version !== 1 || record.kind !== "binding"
return record as Record<string, string>; || 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<string, unknown>): Promise<PublishedResponse> { private async publishSecure(workspaceId: string, revision: string, content: string, manifestBase: Record<string, unknown>): Promise<PublishedResponse> {
@@ -227,17 +337,33 @@ export class WorkspaceRuntimeConfigLeaseFactory {
workspace_revision: revision, config_hex: Buffer.from(content).toString("hex"), manifest_base: manifestBase }); 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"); if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("runtime config publish helper returned malformed output");
const result = value as Record<string, unknown>; const result = value as Record<string, unknown>;
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); let expectedRoot = resolve(this.input.dataRoot);
if (process.platform === "darwin" && (expectedRoot === "/var" || expectedRoot.startsWith("/var/") || expectedRoot === "/tmp" || expectedRoot.startsWith("/tmp/"))) expectedRoot = `/private${expectedRoot}`; 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 expectedPath = join(expectedRoot, "sessions", workspaceId, "preprocessing", "runtime-config", `${revision}.yaml`);
const expectedManifestPath = join(expectedRoot, "sessions", workspaceId, "preprocessing", "runtime-config-manifests", `${revision}.json`); const expectedManifestPath = join(expectedRoot, "sessions", workspaceId, "preprocessing", "runtime-config-manifests", `${revision}.json`);
if (result.path !== expectedPath || result.manifestPath !== expectedManifestPath if (result.protocol_version !== 1 || result.kind !== "publication"
|| typeof result.path !== "string" || !isAbsolute(result.path) || normalize(result.path) !== result.path || result.workspace_id !== workspaceId || result.workspace_revision !== revision
|| typeof result.manifestPath !== "string" || !isAbsolute(result.manifestPath) || normalize(result.manifestPath) !== result.manifestPath || typeof result.path !== "string" || result.path !== expectedPath || !isAbsolute(result.path) || normalize(result.path) !== result.path
|| result.workspace_id !== undefined || typeof result.manifest !== "string" || typeof result.manifestPath !== "string" || result.manifestPath !== expectedManifestPath || !isAbsolute(result.manifestPath) || normalize(result.manifestPath) !== result.manifestPath
|| typeof result.manifest_sha256 !== "string" || !/^[0-9a-f]{64}$/.test(result.manifest_sha256) || 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"); || 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<string, unknown>;
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; return result as unknown as PublishedResponse;
} }
@@ -271,8 +397,9 @@ export class WorkspaceRuntimeConfigLeaseFactory {
snapshots_root: root, repository_root: repositoryRoot, snapshots_root: root, repository_root: repositoryRoot,
workspace_revision: match[1], workspace_id: match[2], workspace_revision: match[1], workspace_id: match[2],
}) as Record<string, unknown>; }) as Record<string, unknown>;
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(",") 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] || verified.workspace_id !== match[2] || verified.workspace_revision !== match[1]
|| typeof verified.source !== "string" || typeof verified.git_source !== "string" || typeof verified.source !== "string" || typeof verified.git_source !== "string"
|| verified.sha256 !== digest(verified.source) || verified.snapshot_path !== path || verified.sha256 !== digest(verified.source) || verified.snapshot_path !== path
@@ -4,7 +4,7 @@ import { tmpdir } from "node:os";
import { dirname, join } from "node:path"; import { dirname, join } from "node:path";
import { execFileSync } from "node:child_process"; import { execFileSync } from "node:child_process";
import { createHash } from "node:crypto"; 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 { parse } from "yaml";
import { parseWorkspaceYaml, serializeWorkspaceYaml } from "../src/workspaces/schema.js"; import { parseWorkspaceYaml, serializeWorkspaceYaml } from "../src/workspaces/schema.js";
@@ -72,6 +72,7 @@ function fixture(extraEnv: Record<string, string> = {}) {
const factoryInput = { const factoryInput = {
dataRoot, runtimeSnapshotRoot: snapshots, harnessDir: harness, configPath, dataRoot, runtimeSnapshotRoot: snapshots, harnessDir: harness, configPath,
env: { env: {
NODE_ENV: "test",
THT_WS_ABC_DWH_TRANSPORT: "postgres_direct", THT_WS_ABC_DWH_HOST: "dwh", 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_PORT: "5432", THT_WS_ABC_DWH_USER: "reader",
THT_WS_ABC_DWH_PASSWORD_FILE: secret, ...extraEnv, 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<string, string> { function realHarnessBinding(config: string): Record<string, string> {
const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py"); const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py");
const python = join(process.cwd(), "..", "harness", ".venv", "bin", "python"); 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", cwd: join(process.cwd(), "..", "harness"), encoding: "utf8",
input: JSON.stringify({ protocol_version: 1, action: "binding", config_hex: Buffer.from(config).toString("hex") }), 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 () => { test("explicit installation overlay is canonical and has one real harness binding", async () => {
@@ -609,3 +612,18 @@ else:
expect(output.trim()).toBe("rejected"); expect(output.trim()).toBe("rejected");
} finally { rmSync(f.root, { recursive: true, force: true }); } } 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<void>((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);
});
@@ -1,10 +1,15 @@
"""Focused unit coverage for the privileged runtime publication seam.""" """Focused unit coverage for the privileged runtime publication seam."""
import io
import json
import os import os
import subprocess
from pathlib import Path
import pytest import pytest
from tht import runtime_config_lease_io as lease_io from tht import runtime_config_lease_io as lease_io
from tht.config import ConfigError, _read_runtime_config_source
def _manifest() -> dict: def _manifest() -> dict:
@@ -82,3 +87,99 @@ def test_retry_reasserts_parent_durability_before_success(tmp_path, monkeypatch,
events.clear() events.clear()
lease_io.publish(inp) lease_io.publish(inp)
assert events.index("config-parent") < events.index("manifest-parent") 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",
})
+11 -1
View File
@@ -492,6 +492,10 @@ def publish(inp: dict) -> dict:
# failed after the no-replace publication on an earlier invocation. # failed after the no-replace publication on an earlier invocation.
publication_fsync(mandir, "manifest-parent") publication_fsync(mandir, "manifest-parent")
return { return {
"protocol_version": 1,
"kind": "publication",
"workspace_id": wid,
"workspace_revision": rev,
"path": f"{canonical}/sessions/{wid}/preprocessing/runtime-config/{name}", "path": f"{canonical}/sessions/{wid}/preprocessing/runtime-config/{name}",
"manifestPath": f"{canonical}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}", "manifestPath": f"{canonical}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}",
"manifest": mb.decode(), "manifest": mb.decode(),
@@ -630,6 +634,8 @@ def verified_snapshot(inp: dict) -> dict:
if blob != record.get("blob"): if blob != record.get("blob"):
fail("workspace Git descriptor identity mismatch") fail("workspace Git descriptor identity mismatch")
return { return {
"protocol_version": 1,
"kind": "verified_snapshot",
"workspace_id": wid, "workspace_id": wid,
"workspace_revision": rev, "workspace_revision": rev,
"source": source.decode(), "source": source.decode(),
@@ -663,7 +669,10 @@ def binding(inp: dict) -> dict:
stream.write(raw) stream.write(raw)
path = Path(stream.name) path = Path(stream.name)
try: 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: finally:
try: try:
path.unlink() path.unlink()
@@ -690,6 +699,7 @@ def main() -> None:
result = verified_snapshot(inp) result = verified_snapshot(inp)
else: else:
result = binding(inp) result = binding(inp)
result = {"protocol_version": 1, "kind": "binding", **result}
print(json.dumps(result)) print(json.dumps(result))
except Exception as e: # noqa: BLE001 except Exception as e: # noqa: BLE001
print(json.dumps({"error": str(e)})) print(json.dumps({"error": str(e)}))