fix(runtime): verify namespace init death and record Wave 3-S validation

This commit is contained in:
Alexandre Teixeira
2026-10-01 23:17:01 +01:00
parent 0dec375631
commit 57143a14ba
7 changed files with 720 additions and 37 deletions
+218 -36
View File
@@ -215,6 +215,8 @@ class ContainmentGrant:
#: it outlives the leader's pid: the leader can exit while the processes it
#: backgrounded keep running in the same group.
pgid: Optional[int] = None
namespace_pid: Optional[int] = None
namespace_start_token: Optional[str] = None
@property
def contained(self) -> bool:
@@ -476,6 +478,8 @@ def _write_record(grant: ContainmentGrant) -> None:
"external": grant.external,
"pid": grant.pid,
"pgid": grant.pgid,
"namespace_pid": grant.namespace_pid,
"namespace_start_token": grant.namespace_start_token,
"acquired_at": time.time(),
"released_at": None,
"release": None,
@@ -827,8 +831,11 @@ def _bwrap_prefix(spec: ContainmentSpec) -> list[str]:
namespace was for. Anything a command legitimately needs outside the
workspace is named by the spec, as ``readonly_extra`` or ``writable_extra``.
"""
executable = shutil.which("bwrap")
if not executable:
raise ContainmentUnavailable(spec.required, "bubblewrap")
args = [
"bwrap", "--die-with-parent", "--new-session", "--unshare-pid",
os.path.abspath(executable), "--die-with-parent", "--new-session", "--unshare-pid",
"--tmpfs", "/",
"--dir", "/usr", "--ro-bind", "/usr", "/usr",
"--symlink", "usr/bin", "/bin",
@@ -890,7 +897,7 @@ def _rlimit_preexec(grant: ContainmentGrant) -> Optional[Callable[[], None]]:
def _launch_argv(grant: ContainmentGrant, command: Any, *, argv: bool,
ready_marker: Optional[str] = None) -> list[str]:
ready_marker: Optional[str] = None, info_fd: Optional[int] = None) -> list[str]:
spec = grant.spec
if argv:
parts = [str(part) for part in command]
@@ -920,7 +927,8 @@ def _launch_argv(grant: ContainmentGrant, command: Any, *, argv: bool,
'printf "%s\\n" "$1"; IFS= read -r ody_ack || exit 125; '
'[ "$ody_ack" = "$1" ] || exit 125; shift; exec "$@"',
"ody-boundary", ready_marker, *parts]
return _bwrap_prefix(spec) + parts
info_args = ["--info-fd", str(info_fd)] if info_fd is not None else []
return _bwrap_prefix(spec) + info_args + parts
return parts
@@ -1003,6 +1011,8 @@ async def run(
"containment: an external-bridge grant describes execution this "
"backend does not own; it cannot be run locally"
)
if (_load_records().get(grant.id) or {}).get("released_at"):
raise ValueError("containment: a released grant cannot execute again")
missing = frozenset(grant.spec.required) - frozenset(grant.enforced)
mechanism = next((item for item in MECHANISMS if item.name == grant.mechanism), None)
provided = mechanism.provides(grant.spec) if mechanism is not None else frozenset()
@@ -1013,8 +1023,15 @@ async def run(
spec = grant.spec
marker = uuid.uuid4().hex if grant.mechanism == "bubblewrap" else None
info_read = info_write = None
try:
launch = _launch_argv(grant, command, argv=argv, ready_marker=marker)
if marker is not None:
info_read, info_write = os.pipe()
os.set_blocking(info_read, False)
launch = _launch_argv(grant, command, argv=argv, ready_marker=marker, info_fd=info_write)
spawn_kwargs = _spawn_kwargs(grant)
if info_write is not None:
spawn_kwargs["pass_fds"] = (info_write,)
spawning = asyncio.create_task(asyncio.create_subprocess_exec(
*launch,
stdin=asyncio.subprocess.PIPE if stdin is not None or marker is not None else asyncio.subprocess.DEVNULL,
@@ -1022,7 +1039,7 @@ async def run(
stderr=asyncio.subprocess.PIPE,
cwd=spec.workspace,
env=dict(spec.env),
**_spawn_kwargs(grant),
**spawn_kwargs,
))
try:
proc = await asyncio.shield(spawning)
@@ -1034,6 +1051,10 @@ async def run(
except Exception:
release(grant, grace_s=0)
else:
if info_write is not None:
os.close(info_write)
info_write = None
proc._ody_info_read = info_read
live = replace(grant, pid=proc.pid, pgid=None if IS_WINDOWS else proc.pid)
await _complete_cleanup(_release_awaited(live, proc), propagate_cancel=False)
raise
@@ -1041,6 +1062,12 @@ async def run(
if "proc" not in locals():
release(grant, grace_s=0)
raise
finally:
if info_write is not None:
os.close(info_write)
if "proc" not in locals() and info_read is not None:
os.close(info_read)
proc._ody_info_read = info_read
# start_new_session makes the child its own group leader, so the group id
# is the child's pid. Captured here rather than at teardown: once the leader
# exits, getpgid can no longer tell us which group its children are in.
@@ -1069,7 +1096,7 @@ async def run(
ready = marker is None
execution_started = marker is None
async def _wait() -> None:
nonlocal ready, execution_started
nonlocal ready, execution_started, live
if marker is not None:
expected = (marker + "\n").encode("ascii")
try:
@@ -1078,9 +1105,14 @@ async def run(
receipt = b""
if receipt != expected:
raise ContainmentUnavailable(spec.required, grant.mechanism)
await _capture_namespace_identity(proc)
live = replace(live, namespace_pid=proc._ody_namespace_pid,
namespace_start_token=proc._ody_namespace_token)
# The trusted child is waiting for acknowledgment, so model code
# cannot exit/recycle the leader before we record its identity.
_update_record(grant.id, start_token=process_ownership.capture(proc.pid)["start_token"])
_update_record(grant.id, namespace_pid=live.namespace_pid,
namespace_start_token=live.namespace_start_token)
if hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal"):
try:
proc._ody_pidfd = os.pidfd_open(proc.pid)
@@ -1135,9 +1167,7 @@ async def run(
out_budget[0] = -1
task.cancel()
await asyncio.gather(task, return_exceptions=True)
pidfd = getattr(proc, "_ody_pidfd", None)
if pidfd is not None:
os.close(pidfd)
_close_process_handles(proc)
try:
try:
await asyncio.wait_for(_wait(), timeout=spec.wall_clock_s)
@@ -1184,6 +1214,41 @@ async def _complete_cleanup(awaitable, *, propagate_cancel: bool = True):
return result
async def _capture_namespace_identity(proc) -> None:
"""Read bwrap's trusted init identity before acknowledging model execution."""
fd = getattr(proc, "_ody_info_read", None)
if fd is None:
return
data = bytearray()
try:
while True:
try:
chunk = os.read(fd, 4096)
except BlockingIOError:
await asyncio.sleep(.01)
continue
if not chunk:
break
data.extend(chunk)
if len(data) > 4096:
raise ValueError("oversized namespace identity")
info = json.loads(data)
pid = int(info["child-pid"])
if pid <= 0 or pid == os.getpid():
raise ValueError("invalid namespace init identity")
proc._ody_namespace_pid = pid
proc._ody_namespace_token = process_ownership.capture(pid)["start_token"]
if hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal"):
proc._ody_namespace_pidfd = os.pidfd_open(pid)
elif not proc._ody_namespace_token:
raise ValueError("namespace init identity cannot be inspected")
except (OSError, ValueError, KeyError, TypeError) as exc:
raise ContainmentUnavailable(frozenset({PROCESS_TREE}), "bubblewrap") from exc
finally:
os.close(fd)
proc._ody_info_read = None
# ── release ─────────────────────────────────────────────────────────────────
# core.platform_compat.kill_process_tree delegates here as well. Native tools,
# detached jobs and compatibility callers share escalation and death probes.
@@ -1368,6 +1433,65 @@ def _ownership_gate(
def release(grant: ContainmentGrant, *, grace_s: float = 2.0,
start_token: Optional[str] = None, require_identity: bool = False) -> ReleaseOutcome:
"""Release owner and recorded namespace init; report death only for both."""
record = _load_records().get(grant.id, {})
previous = record.get("release") or {}
if record.get("released_at") and previous.get("dead"):
# Death belongs to the completed grant, not the current occupant of a
# reused PID slot. Repeated release must never signal it again.
return ReleaseOutcome(dead=True, escalated=bool(previous.get("escalated")),
mechanism=previous.get("mechanism", grant.mechanism),
ownership=previous.get("ownership"))
namespace_pid = grant.namespace_pid or record.get("namespace_pid")
namespace_token = grant.namespace_start_token or record.get("namespace_start_token")
owner = _release_owner(grant, grace_s=grace_s, start_token=start_token,
require_identity=require_identity, _record_release=False)
outcome = owner
if namespace_pid:
try:
namespace_pid = int(namespace_pid)
if namespace_pid <= 0 or namespace_pid == os.getpid():
raise ValueError("invalid namespace init")
except (TypeError, ValueError):
namespace = ReleaseOutcome(dead=False, escalated=False,
ownership=process_ownership.UNVERIFIABLE)
else:
verdict = process_ownership.verify(namespace_pid, namespace_token) if pid_alive(namespace_pid) else process_ownership.GONE
if verdict in (process_ownership.GONE, process_ownership.FOREIGN):
# Reusing PID 1's host slot proves its original namespace has
# completed death; never signal its new occupant.
namespace = ReleaseOutcome(dead=True, escalated=False, ownership=verdict)
else:
target = replace(grant, id=grant.id + ":namespace", mechanism="process_group",
pid=namespace_pid, pgid=None, namespace_pid=None,
namespace_start_token=None)
namespace_fd = None
try:
if hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal"):
try:
namespace_fd = os.pidfd_open(namespace_pid)
except OSError:
pass
namespace = _release_owner(target, grace_s=grace_s, start_token=namespace_token,
require_identity=True, _record_release=False,
_pidfd=namespace_fd)
finally:
if namespace_fd is not None:
os.close(namespace_fd)
outcome = replace(owner, dead=owner.dead and namespace.dead,
escalated=owner.escalated or namespace.escalated,
survivors=tuple(dict.fromkeys((*owner.survivors, *namespace.survivors))))
elif grant.mechanism == "bubblewrap" and record.get("execution_started"):
# A pre-upgrade receipt lacks proof of namespace completion. Keep it
# visible rather than declaring a potentially blocked tree dead.
outcome = replace(owner, dead=False, ownership=process_ownership.UNVERIFIABLE)
_finish_release(grant, outcome)
return outcome
def _release_owner(grant: ContainmentGrant, *, grace_s: float = 2.0,
start_token: Optional[str] = None, require_identity: bool = False,
_record_release: bool = True, _pidfd: Optional[int] = None) -> ReleaseOutcome:
"""Authoritative teardown: signal the group, escalate, then verify.
Returns whether the tree is **observed** gone. A caller must not record a
@@ -1380,6 +1504,10 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0,
a zombie still belongs to its process group, so the group probe would
otherwise report a tree that is already gone.
"""
def finish(target, outcome):
if _record_release:
_finish_release(target, outcome)
pid, pgid = grant.pid, grant.pgid
# A grant that carries its own pid belongs to the process holding it: this
# caller launched the child and no identity question arises. A grant whose
@@ -1405,16 +1533,28 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0,
if pgid is not None and pgid <= 0:
pgid = None
grant = replace(grant, pid=pid or None, pgid=pgid)
def gone():
if _pidfd is not None:
return bool(select.select([_pidfd], [], [], 0)[0])
return _tree_gone(pid, pgid, reap=True)
def send(sig):
if _pidfd is not None:
try:
signal.pidfd_send_signal(_pidfd, sig)
except OSError:
pass
else:
_signal_tree(pid, pgid, sig)
if not pid and not _group_present(pgid):
outcome = _outcome_for(grant, dead=True, escalated=False)
_finish_release(grant, outcome)
finish(grant, outcome)
return outcome
if recovered or require_identity:
refusal = _ownership_gate(grant, pid, pgid, token)
if refusal is not None:
_finish_release(grant, refusal)
finish(grant, refusal)
return refusal
if IS_WINDOWS:
@@ -1431,36 +1571,36 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0,
while time.monotonic() < deadline and pid_alive(pid):
time.sleep(_DEATH_POLL_S)
outcome = _outcome_for(grant, dead=not pid_alive(pid), escalated=True)
_finish_release(grant, outcome)
finish(grant, outcome)
return outcome
if _tree_gone(pid, pgid, reap=True):
if gone():
outcome = _outcome_for(grant, dead=True, escalated=False)
_finish_release(grant, outcome)
finish(grant, outcome)
return outcome
_signal_tree(pid, pgid, signal.SIGTERM)
send(signal.SIGTERM)
escalated = False
deadline = time.monotonic() + max(grace_s, 0.0)
while time.monotonic() < deadline and not _tree_gone(pid, pgid, reap=True):
while time.monotonic() < deadline and not gone():
time.sleep(_DEATH_POLL_S)
if not _tree_gone(pid, pgid, reap=True):
if not gone():
escalated = True
if recovered or require_identity:
refusal = _ownership_gate(grant, pid, pgid, token)
if refusal is not None:
_finish_release(grant, refusal)
finish(grant, refusal)
return refusal
_signal_tree(pid, pgid, signal.SIGKILL)
send(signal.SIGKILL)
# SIGKILL cannot be caught, so a short verification window is enough.
# Anything still here is out of our reach — a zombie whose parent is
# not us, or a pid we never owned.
deadline = time.monotonic() + 1.0
while time.monotonic() < deadline and not _tree_gone(pid, pgid, reap=True):
while time.monotonic() < deadline and not gone():
time.sleep(_DEATH_POLL_S)
outcome = _outcome_for(grant, dead=_tree_gone(pid, pgid, reap=True), escalated=escalated)
_finish_release(grant, outcome)
outcome = _outcome_for(grant, dead=gone(), escalated=escalated)
finish(grant, outcome)
return outcome
@@ -1506,6 +1646,25 @@ def reap_record(record: Mapping[str, Any], *, grace_s: float = 2.0) -> ReleaseOu
async def _release_awaited(
grant: ContainmentGrant,
proc: "asyncio.subprocess.Process",
*, grace_s: float = 2.0,
) -> ReleaseOutcome:
try:
return await _release_awaited_impl(grant, proc, grace_s=grace_s)
finally:
_close_process_handles(proc)
def _close_process_handles(proc) -> None:
for name in ("_ody_info_read", "_ody_pidfd", "_ody_namespace_pidfd"):
fd = getattr(proc, name, None)
if fd is not None:
os.close(fd)
setattr(proc, name, None)
async def _release_awaited_impl(
grant: ContainmentGrant,
proc: "asyncio.subprocess.Process",
*,
@@ -1531,30 +1690,51 @@ async def _release_awaited(
if IS_WINDOWS:
return release(grant, grace_s=grace_s)
if getattr(proc, "_ody_info_read", None) is not None:
try:
await asyncio.wait_for(_capture_namespace_identity(proc), timeout=1)
except (ContainmentUnavailable, asyncio.TimeoutError):
pass # Setup never reached acknowledgment; no model code ran.
pid, pgid = grant.pid, grant.pgid
if grant.mechanism == "bubblewrap" and proc.returncode is not None:
# Namespace-owner death already destroyed the namespace. Its numeric
# PID/PGID may now be reused; no further signal is necessary or safe.
await proc.wait()
outcome = _outcome_for(grant, dead=True, escalated=False)
_finish_release(grant, outcome)
return outcome
pidfd = getattr(proc, "_ody_pidfd", None)
namespace_pid = getattr(proc, "_ody_namespace_pid", None)
namespace_token = getattr(proc, "_ody_namespace_token", None)
namespace_fd = getattr(proc, "_ody_namespace_pidfd", None)
if namespace_pid:
_update_record(grant.id, namespace_pid=namespace_pid, namespace_start_token=namespace_token)
def namespace_gone():
if namespace_fd is not None:
return bool(select.select([namespace_fd], [], [], 0)[0])
if namespace_pid:
if not pid_alive(namespace_pid):
return True
return process_ownership.verify(namespace_pid, namespace_token) in (
process_ownership.GONE, process_ownership.FOREIGN,
)
return True
def gone():
if pidfd is not None:
return bool(select.select([pidfd], [], [], 0)[0])
return _tree_gone(pid, pgid)
owner_gone = bool(select.select([pidfd], [], [], 0)[0])
elif grant.mechanism == "bubblewrap":
owner_gone = proc.returncode is not None
else:
owner_gone = _tree_gone(pid, pgid)
return owner_gone and namespace_gone()
def send(sig):
if pidfd is not None:
try:
# Killing the bwrap owner destroys its private PID namespace,
# including descendants that changed session/group. A pidfd
# keeps a recycled numeric PID out of the signal path.
signal.pidfd_send_signal(pidfd, sig)
except OSError:
pass
else:
elif proc.returncode is None or grant.mechanism != "bubblewrap":
_signal_tree(pid, pgid, sig)
if namespace_fd is not None:
try:
signal.pidfd_send_signal(namespace_fd, sig)
except OSError:
pass
elif namespace_pid and process_ownership.verify(namespace_pid, namespace_token) == process_ownership.OWNED:
_signal_tree(namespace_pid, None, sig)
send(signal.SIGTERM)
try:
await asyncio.wait_for(proc.wait(), timeout=max(grace_s, 0.05))
@@ -1577,6 +1757,8 @@ async def _release_awaited(
await asyncio.sleep(_DEATH_POLL_S)
outcome = _outcome_for(grant, dead=gone(), escalated=escalated)
if not namespace_gone():
outcome = replace(outcome, survivors=tuple(dict.fromkeys((*outcome.survivors, namespace_pid))))
_finish_release(grant, outcome)
return outcome
@@ -1588,4 +1770,4 @@ def _finish_release(grant: ContainmentGrant, outcome: ReleaseOutcome) -> None:
# Deliberately NOT released: the record stays active so a reaper sees it
# again. A record claiming teardown it did not achieve is the defect
# this reverses.
_update_record(grant.id, release=outcome.to_dict())
_update_record(grant.id, released_at=None, release=outcome.to_dict())