Files
ThothII/harness/tht/runtime_config_lease_io.py
T

701 lines
30 KiB
Python

"""Small privileged filesystem seam for durable runtime configuration publication."""
from __future__ import annotations
import ctypes
import fcntl
import hashlib
import json
import os
import platform
import stat
import subprocess
import sys
import tempfile
import time
from pathlib import Path
def fail(msg: str) -> None:
raise RuntimeError(msg)
def _canonical_root(root: str) -> str:
# Darwin exposes /tmp and /var as symlink aliases. Linux does not: rewriting
# these paths there would redirect valid installations to a different root.
if sys.platform == "darwin":
if root == "/tmp" or root.startswith("/tmp/"):
return "/private" + root
if root == "/var" or root.startswith("/var/"):
return "/private" + root
return root
def _identity(st: os.stat_result, path: str) -> dict[str, str]:
return {"path": path, "dev": str(st.st_dev), "ino": str(st.st_ino),
"mode": format(stat.S_IMODE(st.st_mode), "o"), "uid": str(st.st_uid)}
def _rename_noreplace(src: str, dst: str, directory_fd: int, kind: str) -> None:
"""Atomically rename *src* to *dst* without replacing an existing entry.
``link`` is deliberately not used here: the state-file protocol requires a
rename, and a hard-link publication leaves a second name visible during a
crash. Unsupported platforms fail closed rather than silently weakening the
protocol. The env seam is intentionally narrow so fault tests can exercise
every publication rename without monkey-patching the privileged process.
"""
if os.environ.get("THT_RUNTIME_CONFIG_RENAME_FAIL") in {"1", kind}:
fail("runtime config rename failed")
libc = ctypes.CDLL(None, use_errno=True)
src_b = os.fsencode(src)
dst_b = os.fsencode(dst)
if sys.platform == "darwin":
fn = getattr(libc, "renameatx_np", None)
if fn is None:
fail("runtime config no-replace rename is unavailable")
fn.argtypes = [ctypes.c_int, ctypes.c_char_p, ctypes.c_int, ctypes.c_char_p, ctypes.c_uint]
fn.restype = ctypes.c_int
# RENAME_EXCL is the Darwin no-overwrite operation.
rc = fn(directory_fd, src_b, directory_fd, dst_b, 0x00000004)
elif sys.platform.startswith("linux"):
# renameat2(2), RENAME_NOREPLACE. syscall numbers are stable for the
# supported Linux architectures; an unavailable syscall fails closed.
number = {"x86_64": 316, "aarch64": 276, "arm64": 276}.get(platform.machine())
if number is None or not hasattr(libc, "syscall"):
fail("runtime config no-replace rename is unavailable")
libc.syscall.argtypes = [ctypes.c_long, ctypes.c_int, ctypes.c_char_p,
ctypes.c_int, ctypes.c_char_p, ctypes.c_uint]
libc.syscall.restype = ctypes.c_long
rc = libc.syscall(number, directory_fd, src_b, directory_fd, dst_b, 1)
else:
fail("runtime config no-replace rename is unavailable")
if rc != 0:
err = ctypes.get_errno()
if err == 17:
raise FileExistsError(err, os.strerror(err), dst)
raise OSError(err, os.strerror(err), dst)
def directory_identities(root: str, paths: list[tuple[str, int]]) -> list[dict[str, str]]:
"""Return the ordered, no-follow identity chain bound by a publication.
``paths`` contains canonical absolute directory paths paired with already-open
descriptors. Prefix components are stat'ed without following symlinks; the
terminal workspace-owned components are additionally represented by fstat on
the descriptors opened by ``walk``.
"""
canonical = _canonical_root(root)
root_parts = [part for part in Path(canonical).parts if part not in ("", "/")]
entries: list[dict[str, str]] = []
current = "/"
st = os.stat("/", follow_symlinks=False)
entries.append(_identity(st, current))
for part in root_parts:
current = (current.rstrip("/") + "/" + part) if current != "/" else "/" + part
st = os.stat(current, follow_symlinks=False)
if stat.S_ISLNK(st.st_mode) or not stat.S_ISDIR(st.st_mode):
fail("runtime config directory is not trusted")
entries.append(_identity(st, current))
for path, fd in paths:
cpath = _canonical_root(path)
st = os.fstat(fd)
if not stat.S_ISDIR(st.st_mode) or stat.S_IMODE(st.st_mode) != 0o700 or st.st_uid != os.getuid():
fail("runtime config directory is not trusted")
# Keep one ordered entry per path. Existing prefixes are left in place.
if not any(item["path"] == cpath for item in entries):
entries.append(_identity(st, cpath))
else:
for item in entries:
if item["path"] == cpath:
if item["dev"] != str(st.st_dev) or item["ino"] != str(st.st_ino):
fail("runtime config directory changed")
break
return entries
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:
"""Open one directory component without following a replaced entry.
The pre-open lstat and post-open fstat identity check is required on Darwin,
where O_NOFOLLOW has historically been unavailable for directory openat.
mkdir races are resolved by opening and validating the winner.
"""
flags = os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | os.O_NOFOLLOW
while True:
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")
fd = os.open(name, flags, dir_fd=parent)
try:
opened = os.fstat(fd)
if (opened.st_dev != entry.st_dev or opened.st_ino != entry.st_ino
or not stat.S_ISDIR(opened.st_mode)):
fail("runtime config directory changed during open")
return fd
except BaseException:
os.close(fd)
raise
except FileNotFoundError:
if not create:
raise
try:
os.mkdir(name, 0o700, dir_fd=parent)
except FileExistsError:
# Another publisher won creation. Re-enter the identity-checked
# open path instead of exposing EEXIST to the caller.
continue
if parent is not None:
os.fsync(parent)
# Re-open through the same no-follow and identity checks.
continue
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:
"""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")
# 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 sys.platform == "darwin" and (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:
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
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 publication_fsync(fd: int, stage: str) -> None:
"""Fsync one publication boundary, with an explicit test-only fault seam.
The seam is inert unless the caller opts in with the exact stage name. It is
intentionally kept here (rather than in TypeScript) so the real helper's
durable ordering can be exercised without weakening production behaviour.
"""
if os.environ.get("THT_RUNTIME_CONFIG_FSYNC_FAIL") == stage:
fail(f"runtime config fsync failed at {stage}")
os.fsync(fd)
def read_all(fd: int, limit: int = 16 * 1024 * 1024) -> bytes:
if limit < 0:
fail("runtime config file is too large")
os.lseek(fd, 0, os.SEEK_SET)
if limit == 0:
if os.read(fd, 1):
fail("runtime config file is too large")
return b""
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:
# The bounded read above cannot observe an additional byte when it
# lands exactly on the ceiling; probe once before accepting it.
if os.read(fd, 1):
fail("runtime config file is too large")
return b"".join(chunks)
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", "descriptor_dev", "descriptor_ino", "config_sha256",
"config_dwh_binding", "config_dev", "config_ino", "config_size", "config_mode",
"config_uid", "config_nlink", "directory_identities",
}
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 ("descriptor_dev", "descriptor_ino", "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")
identities = value.get("directory_identities")
if (not isinstance(identities, list) or not identities or
any(not isinstance(item, dict) or set(item) != {"path", "dev", "ino", "mode", "uid"}
or not isinstance(item["path"], str) or not os.path.isabs(item["path"])
or any(not isinstance(item[key], str) or not item[key].isdigit()
for key in ("dev", "ino", "mode", "uid"))
or item["mode"] == "0"
for item in identities)):
fail("runtime config manifest is invalid")
if len({item["path"] for item in identities}) != len(identities):
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")
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)
canonical = _canonical_root(root)
directory_manifest = directory_identities(root, [
(f"{canonical}/sessions", sessions),
(f"{canonical}/sessions/{wid}", ws),
(f"{canonical}/sessions/{wid}/preprocessing", prep),
(f"{canonical}/sessions/{wid}/preprocessing/runtime-config", cfgdir),
(f"{canonical}/sessions/{wid}/preprocessing/runtime-config-manifests", mandir),
])
# The retained preprocessing directory is the single cross-process lock seam.
# No pathname lock file is created in the workspace layout.
deadline = time.monotonic() + 2.0
while True:
try:
fcntl.flock(prep, fcntl.LOCK_EX | fcntl.LOCK_NB)
break
except BlockingIOError:
if time.monotonic() >= deadline:
# Never let a wedged publisher block its caller indefinitely. The
# Node boundary turns this stable conflict into a bounded failure.
fail("runtime config publication is busy")
time.sleep(0.01)
try:
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 = read_all(fd)
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)
publication_fsync(fd, "config-file")
try:
_rename_noreplace(stage, name, cfgdir, "config")
except FileExistsError:
# A concurrent equal publisher may already have won. It is
# accepted only after reopening and comparing its bytes below.
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 read_all(fd) != content:
fail("same-revision runtime configuration changed")
finally:
os.close(fd)
publication_fsync(cfgdir, "config-parent")
# A prior invocation may have reported success after publishing the
# entry but before its parent fsync. Re-establish that durability
# boundary before making the manifest durable.
publication_fsync(cfgdir, "config-parent")
# Identity is deliberately recorded after final no-replace publication.
got = current(cfgdir, name, 0o400)
assert got
fd, s = got
os.close(fd)
manifest = {"version": 1, **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),
"directory_identities": directory_manifest,
}
)
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 = 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:
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)
publication_fsync(fd, "manifest-file")
try:
_rename_noreplace(stage, mname, mandir, "manifest")
except FileExistsError:
# A no-replace loser is successful only after validating the
# durable winner byte-for-byte and against the strict schema.
winner = current(mandir, mname, 0o600)
if winner is None:
fail("runtime config manifest publication raced")
wfd, _ = winner
try:
existing = read_all(wfd)
finally:
os.close(wfd)
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")
finally:
os.close(fd)
try:
os.unlink(stage, dir_fd=mandir)
except FileNotFoundError:
pass
publication_fsync(mandir, "manifest-parent")
# As with the config directory, retries must repair a boundary that
# failed after the no-replace publication on an earlier invocation.
publication_fsync(mandir, "manifest-parent")
return {
"path": f"{canonical}/sessions/{wid}/preprocessing/runtime-config/{name}",
"manifestPath": f"{canonical}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}",
"manifest": mb.decode(),
"manifest_sha256": hashlib.sha256(mb).hexdigest(),
"dev": s.st_dev,
"ino": s.st_ino,
}
finally:
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:
descriptor_info = 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)
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")
if not isinstance(records, list) or not isinstance(files, dict) or not records:
fail("workspace snapshot integrity check failed")
record_by_id: dict[str, dict] = {}
for item in records:
if not isinstance(item, dict) or set(item) != {"id", "commit", "blob", "snapshotPath"}:
fail("workspace snapshot integrity check failed")
item_id = item.get("id")
if not isinstance(item_id, str) or not safe_id(item_id) or item_id in record_by_id:
fail("workspace snapshot integrity check failed")
if item.get("commit") != rev or item.get("snapshotPath") != f"{root}/{rev}/{item_id}.yaml":
fail("workspace snapshot integrity check failed")
if not isinstance(item.get("blob"), str) or not safe_rev(item["blob"]):
fail("workspace snapshot integrity check failed")
record_by_id[item_id] = item
expected_names = {name for item_id in record_by_id for name in (f"{item_id}.yaml", f"{item_id}.env.example", f"{item_id}.md")}
if set(files) != expected_names or any(not isinstance(v, str) or not __import__("re").fullmatch(r"[0-9a-f]{64}", v) for v in files.values()):
fail("workspace snapshot integrity check failed")
# Verify every immutable file declared by snapshot.json, not just the selected
# descriptor. This prevents extra records/files from smuggling a second state.
for filename in sorted(expected_names):
f = os.open(filename, os.O_RDONLY | os.O_NOFOLLOW, dir_fd=rdir)
try:
read_regular(f, 0o400)
actual = hashlib.sha256(read_all(f)).hexdigest()
finally:
os.close(f)
if actual != files[filename]:
fail("workspace snapshot integrity check failed")
record = record_by_id.get(wid)
expected_path = f"{root}/{rev}/{wid}.yaml"
if record is None or manifest.get("head") != rev or record.get("snapshotPath") != expected_path:
fail("workspace snapshot integrity check failed")
if files.get(f"{wid}.yaml") != hashlib.sha256(source).hexdigest():
fail("workspace snapshot integrity check failed")
repo = inp.get("repository_root")
if not isinstance(repo, str) or not os.path.isabs(repo):
fail("invalid repository root")
# Git replacement refs and ambient repository/config variables are attacker
# controlled process state. Snapshot identity must be the raw object named by
# the commit, with fixed Git configuration and repository boundaries.
# Keep no inherited GIT_* controls at all (including GIT_CONFIG_PARAMETERS,
# alternates, and repository path overrides), then add only fixed semantics.
git_env = {key: value for key, value in os.environ.items() if not key.startswith("GIT_")}
git_env.update({"GIT_NO_REPLACE_OBJECTS": "1", "GIT_CONFIG_NOSYSTEM": "1",
"GIT_CONFIG_GLOBAL": os.devnull, "GIT_CONFIG_SYSTEM": os.devnull})
try:
# A 40-hex object name is not necessarily a commit (trees and blobs are
# valid Git objects and also accept the <object>:path syntax). Require
# the raw object itself to be a commit, with replacement/config controls
# disabled, before reading any descriptor bytes.
object_type = subprocess.check_output(
["git", "--no-replace-objects", "-C", repo, "cat-file", "-t", rev],
stderr=subprocess.DEVNULL, text=True, timeout=5, env=git_env,
).strip()
resolved_commit = subprocess.check_output(
["git", "--no-replace-objects", "-C", repo, "rev-parse", f"{rev}^{{commit}}"],
stderr=subprocess.DEVNULL, text=True, timeout=5, env=git_env,
).strip()
if object_type != "commit" or resolved_commit != rev:
fail("workspace Git revision is not an exact commit")
blob = subprocess.check_output(
["git", "--no-replace-objects", "-C", repo, "rev-parse", f"{rev}:workspaces/{wid}.yaml"],
stderr=subprocess.DEVNULL, text=True, timeout=5, env=git_env,
).strip()
git_source = subprocess.check_output(
["git", "--no-replace-objects", "-C", repo, "show", f"{rev}:workspaces/{wid}.yaml"],
stderr=subprocess.DEVNULL, timeout=5, env=git_env,
)
except (OSError, subprocess.SubprocessError):
fail("workspace Git revision is unavailable")
try:
git_text = git_source.decode("utf-8")
except UnicodeDecodeError:
fail("workspace Git descriptor identity mismatch")
if blob != record.get("blob"):
fail("workspace Git descriptor identity mismatch")
return {
"workspace_id": wid,
"workspace_revision": rev,
"source": source.decode(),
"git_source": git_text,
"sha256": hashlib.sha256(source).hexdigest(),
"descriptor_git_blob": record.get("blob"),
"descriptor_dev": descriptor_info.st_dev,
"descriptor_ino": descriptor_info.st_ino,
"snapshot_path": f"{root}/{rev}/{wid}.yaml",
}
def binding(inp: dict) -> dict:
# Runtime handoff variables are capabilities, never helper input. Remove
# inherited values before load_config can inspect its environment.
for key in tuple(os.environ):
if key.startswith(("THT_RUNTIME_CONFIG_", "THT_CONFIG_")):
os.environ.pop(key, None)
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)
if not isinstance(inp, dict) or inp.get("protocol_version") != 1:
fail("unsupported runtime config protocol")
action = inp.get("action")
request_keys = {
"publish": {"protocol_version", "action", "data_root", "workspace_id", "workspace_revision", "config_hex", "manifest_base"},
"verified-snapshot": {"protocol_version", "action", "snapshots_root", "repository_root", "workspace_revision", "workspace_id"},
"binding": {"protocol_version", "action", "config_hex"},
}
if action not in request_keys or set(inp) != request_keys[action]:
fail("invalid runtime config request")
if action == "publish":
result = publish(inp)
elif action == "verified-snapshot":
result = verified_snapshot(inp)
else:
result = binding(inp)
print(json.dumps(result))
except Exception as e: # noqa: BLE001
print(json.dumps({"error": str(e)}))
raise SystemExit(1)
if __name__ == "__main__":
main()