mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-08 07:52:20 +02:00
fix(runtime): enforce functional containment after closing execution bypasses (ODY-141)
This commit is contained in:
@@ -388,7 +388,7 @@ def _contained_command(
|
||||
comparable — one is a mount namespace, the other is a regex — and choosing
|
||||
the second silently means an uncontained host execution reads in the
|
||||
transcript exactly like a contained one. The fallback still happens under
|
||||
report-only mode, which is what ships; the difference is that it is now
|
||||
an explicit report-only diagnostic mode; the difference is that it is now
|
||||
recorded in the result.
|
||||
|
||||
:raises containment.ContainmentUnavailable: filesystem containment could
|
||||
@@ -399,7 +399,7 @@ def _contained_command(
|
||||
content, cwd, chdir=chdir, interpreter_prefix=interpreter_prefix,
|
||||
)
|
||||
# The probe's filesystem answer and the wrapper's None/not-None answer rest
|
||||
# on the same condition (`not IS_WINDOWS and which("bwrap")`), so they agree
|
||||
# on the same functional namespace probe, so they agree
|
||||
# by construction. `wrapped` is still what decides, because it is what
|
||||
# actually runs: a probe that said yes to a wrapper that declined would be
|
||||
# the same false claim in the other direction.
|
||||
@@ -449,41 +449,9 @@ def _wrap_workspace_namespace(
|
||||
bubblewrap namespace preserves that public contract for each concurrent
|
||||
agent without creating a process-global /workspace symlink.
|
||||
"""
|
||||
if IS_WINDOWS or not shutil.which("bwrap"):
|
||||
if IS_WINDOWS or not containment._bwrap_available():
|
||||
return None
|
||||
args = [
|
||||
"bwrap", "--die-with-parent", "--new-session", "--tmpfs", "/",
|
||||
"--dir", "/usr", "--ro-bind", "/usr", "/usr",
|
||||
"--symlink", "usr/bin", "/bin",
|
||||
"--symlink", "usr/lib", "/lib",
|
||||
"--symlink", "usr/lib64", "/lib64",
|
||||
"--symlink", "usr/bin", "/sbin",
|
||||
"--dir", "/etc", "--ro-bind", "/etc", "/etc",
|
||||
# Read-only, not read-write. These two binds exist so a command can
|
||||
# *read* host material it legitimately needs — a dataset under /mnt, a
|
||||
# dotfile under /home. Binding them writable gave back most of what
|
||||
# the namespace was for: a Linux host with working bubblewrap running
|
||||
# this argv reaches outside the workspace and writes to the user's home
|
||||
# directory, measured rather than inferred. The workspace bind below is
|
||||
# the one writable path, which is what "workspace confinement" means.
|
||||
"--dir", "/home", "--ro-bind", "/home", "/home",
|
||||
"--dir", "/mnt", "--ro-bind", "/mnt", "/mnt",
|
||||
"--dir", "/tmp", "--tmpfs", "/tmp",
|
||||
"--dev-bind", "/dev", "/dev", "--proc", "/proc",
|
||||
"--dir", WORKSPACE_MOUNT, "--bind", cwd, WORKSPACE_MOUNT,
|
||||
]
|
||||
# The workspace stays writable at its real host path as well as at
|
||||
# /workspace. A command can carry the absolute host path: BashTool's own
|
||||
# /tmp redirect rewrites `/tmp/` to `<agent_cwd()>/.tmp/` before the
|
||||
# namespace is built, so the command reaching bwrap already names the real
|
||||
# path. Before /home and /mnt became read-only those writes landed only
|
||||
# because the workspace happened to sit under one of them. Binding the
|
||||
# workspace itself is the narrow version of what that accident provided:
|
||||
# the same directory by either name, and nothing else writable.
|
||||
real_cwd = os.path.realpath(cwd)
|
||||
if real_cwd not in _NAMESPACE_RESERVED_DESTS and len(real_cwd.split(os.sep)) >= 3:
|
||||
args.extend(_namespace_dir_chain(real_cwd))
|
||||
args.extend(("--bind", real_cwd, real_cwd))
|
||||
readonly = [path for path in ("/home", "/mnt") if os.path.isdir(path)]
|
||||
# setup-python installs interpreters under /opt, and local CI virtualenvs
|
||||
# can live under /tmp. Those paths are hidden by the private root/tmpfs.
|
||||
# Expose only the active interpreter environment, read-only, so Python
|
||||
@@ -513,9 +481,14 @@ def _wrap_workspace_namespace(
|
||||
and os.path.isdir(prefix)
|
||||
and has_environment_layout
|
||||
):
|
||||
args.extend(_namespace_dir_chain(prefix))
|
||||
args.extend(("--ro-bind", prefix, prefix))
|
||||
args.extend(("--chdir", chdir, "/bin/bash", "-lc", content))
|
||||
readonly.append(prefix)
|
||||
spec = containment.ContainmentSpec(
|
||||
workspace=cwd, env={}, wall_clock_s=DEFAULT_BASH_TIMEOUT,
|
||||
readonly_extra=tuple(readonly),
|
||||
)
|
||||
args = containment._bwrap_prefix(spec)
|
||||
args[-1] = chdir
|
||||
args.extend(("/bin/bash", "-lc", content))
|
||||
return shlex.join(args)
|
||||
|
||||
|
||||
@@ -524,7 +497,8 @@ def _owned_spec(cwd: str, env: Optional[dict], timeout: int, readonly_extra: tup
|
||||
readonly = []
|
||||
for prefix in (sys.prefix, sys.base_prefix):
|
||||
prefix = os.path.realpath(prefix)
|
||||
if not _namespace_visible_without_bind(prefix) and prefix not in _NAMESPACE_RESERVED_DESTS:
|
||||
visible = any(prefix == root or prefix.startswith(root + os.sep) for root in ("/usr", "/etc"))
|
||||
if not visible and prefix not in _NAMESPACE_RESERVED_DESTS:
|
||||
readonly.append(prefix)
|
||||
return containment.agent_spec(
|
||||
cwd, dict(os.environ if env is None else env), timeout,
|
||||
@@ -551,8 +525,12 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
|
||||
except containment.ContainmentUnavailable as exc:
|
||||
return containment.unavailable_tool_result(exc, tool=tool)
|
||||
except (OSError, RuntimeError, ValueError) as exc:
|
||||
boundary = grant.to_dict() if grant else {}
|
||||
boundary["executed"] = bool(getattr(exc, "containment_executed", False))
|
||||
if not getattr(exc, "containment_established", False):
|
||||
boundary.update(contained=False, enforced=[])
|
||||
return {"error": f"{tool}: execution failed: {exc}", "exit_code": 1,
|
||||
"containment": grant.to_dict() if grant else {"contained": False, "executed": False}}
|
||||
"containment": boundary}
|
||||
|
||||
boundary = result.grant.to_dict()
|
||||
boundary["executed"] = True
|
||||
|
||||
+209
-87
@@ -39,11 +39,10 @@ deployment detail.
|
||||
* :data:`MODE_REPORT_ONLY` — the same shortfall is recorded on the grant as
|
||||
``unenforced_required``, logged once, and the command runs.
|
||||
|
||||
The shipped default is report-only. On macOS and in the shipped Docker image
|
||||
there is no ``bwrap``, so enforcing filesystem containment by default would turn
|
||||
every ``bash`` call into a refusal the moment this module is wired up. Starting
|
||||
report-only makes that landing observable instead of breaking, and flipping the
|
||||
constant is reversible in a way that breaking every host is not.
|
||||
The shipped default is enforcing. Hosts without functional namespaces refuse
|
||||
native agent execution requiring filesystem and process-tree containment.
|
||||
Report-only remains an explicit internal diagnostic posture, never a tool or
|
||||
deployment setting. Networking is inherited unless the spec requests isolation.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -54,6 +53,7 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import select
|
||||
import signal
|
||||
import subprocess
|
||||
import sys
|
||||
@@ -81,10 +81,8 @@ logger = logging.getLogger(__name__)
|
||||
MODE_ENFORCING = "enforcing"
|
||||
MODE_REPORT_ONLY = "report_only"
|
||||
|
||||
#: Ship report-only; see the module docstring for why this is the reversible
|
||||
#: direction. Flip to MODE_ENFORCING to make an unestablishable required
|
||||
#: dimension refuse the command instead of reporting on it.
|
||||
CONTAINMENT_MODE = MODE_REPORT_ONLY
|
||||
#: Required containment must be established before model-controlled code runs.
|
||||
CONTAINMENT_MODE = MODE_ENFORCING
|
||||
|
||||
|
||||
# ── Dimensions ──────────────────────────────────────────────────────────────
|
||||
@@ -296,7 +294,23 @@ class Mechanism:
|
||||
|
||||
|
||||
def _bwrap_available() -> bool:
|
||||
return not IS_WINDOWS and bool(shutil.which("bwrap"))
|
||||
if IS_WINDOWS:
|
||||
return False
|
||||
executable = shutil.which("bwrap")
|
||||
if not executable:
|
||||
return False
|
||||
try:
|
||||
# Binary installation says nothing about namespace permissions (notably
|
||||
# under Docker's normal security profile). Only trusted probe code runs.
|
||||
probe = subprocess.run(
|
||||
[executable, "--die-with-parent", "--unshare-pid", "--ro-bind", "/", "/",
|
||||
"--proc", "/proc", "--dev", "/dev", "/bin/true"],
|
||||
stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
|
||||
timeout=3, check=False,
|
||||
)
|
||||
return probe.returncode == 0
|
||||
except (OSError, subprocess.SubprocessError):
|
||||
return False
|
||||
|
||||
|
||||
def _posix_group_available() -> bool:
|
||||
@@ -353,6 +367,7 @@ def _rlimit_dimensions(spec: ContainmentSpec) -> set[str]:
|
||||
provided.add(MEMORY)
|
||||
if (
|
||||
spec.max_processes is not None
|
||||
and os.geteuid() != 0 # RLIMIT_NPROC does not limit root.
|
||||
and _rlimit_fits("RLIMIT_NPROC", spec.max_processes)
|
||||
):
|
||||
provided.add(PROCESS_COUNT)
|
||||
@@ -361,7 +376,7 @@ def _rlimit_dimensions(spec: ContainmentSpec) -> set[str]:
|
||||
|
||||
def _bwrap_provides(spec: ContainmentSpec) -> frozenset[str]:
|
||||
# bwrap gives the private root and the workspace bind (filesystem), a new
|
||||
# session plus --die-with-parent (process_tree), and --unshare-net when the
|
||||
# PID namespace plus --die-with-parent (process_tree), and --unshare-net when the
|
||||
# spec asked for no network. The wall clock and the resource limits are
|
||||
# ours either way, applied to the bwrap process itself so its descendants
|
||||
# inherit them.
|
||||
@@ -372,16 +387,15 @@ def _bwrap_provides(spec: ContainmentSpec) -> frozenset[str]:
|
||||
|
||||
|
||||
def _posix_group_provides(spec: ContainmentSpec) -> frozenset[str]:
|
||||
# A process group plus setsid makes the kill authoritative and the wall
|
||||
# clock real for the whole tree. It says nothing about the filesystem: a
|
||||
# cwd is not a boundary.
|
||||
return frozenset({PROCESS_TREE, WALL_CLOCK} | _rlimit_dimensions(spec))
|
||||
# Groups support escalating teardown, but a descendant can call setsid()
|
||||
# and escape. They cannot truthfully establish process-tree containment.
|
||||
return frozenset({WALL_CLOCK} | _rlimit_dimensions(spec))
|
||||
|
||||
|
||||
def _windows_provides(spec: ContainmentSpec) -> frozenset[str]:
|
||||
# taskkill /T /F walks the child tree, which is the Windows equivalent of
|
||||
# signalling a group. There is no setrlimit and no namespace.
|
||||
return frozenset({PROCESS_TREE, WALL_CLOCK})
|
||||
# taskkill supports teardown, but is not a Job Object preventing escaped
|
||||
# descendants. No filesystem or process-tree containment is established.
|
||||
return frozenset({WALL_CLOCK})
|
||||
|
||||
|
||||
#: Strongest first. Selection walks this in order and stops at the first
|
||||
@@ -447,7 +461,7 @@ def _write_record(grant: ContainmentGrant) -> None:
|
||||
"id": grant.id,
|
||||
"owner": grant.owner,
|
||||
"manager_pid": os.getpid(),
|
||||
"manager_token": process_ownership.start_token(os.getpid()),
|
||||
"manager_token": process_ownership.capture(os.getpid())["start_token"],
|
||||
"mechanism": grant.mechanism,
|
||||
"mode": grant.mode,
|
||||
"workspace": grant.workspace,
|
||||
@@ -784,13 +798,16 @@ def unavailable_tool_result(exc: ContainmentUnavailable, *, tool: str) -> dict[s
|
||||
|
||||
|
||||
# ── Launch plumbing ─────────────────────────────────────────────────────────
|
||||
def _dir_chain(path: str) -> list[str]:
|
||||
def _dir_chain(path: str, mounted: tuple[str, ...] = ()) -> list[str]:
|
||||
"""``--dir`` args for every ancestor of ``path`` inside the private root.
|
||||
|
||||
bwrap mounts into a tmpfs root, so the destination's parents have to exist
|
||||
before the bind. Stops at the mount points the argv already creates.
|
||||
"""
|
||||
args: list[str] = []
|
||||
roots = ("/usr", "/etc", *mounted)
|
||||
if any(path == root or path.startswith(root + os.sep) for root in roots):
|
||||
return args
|
||||
parents: list[str] = []
|
||||
parent = os.path.dirname(path)
|
||||
while parent not in ("/", "", "/tmp", "/etc", "/usr", WORKSPACE_MOUNT):
|
||||
@@ -811,7 +828,7 @@ def _bwrap_prefix(spec: ContainmentSpec) -> list[str]:
|
||||
workspace is named by the spec, as ``readonly_extra`` or ``writable_extra``.
|
||||
"""
|
||||
args = [
|
||||
"bwrap", "--die-with-parent", "--new-session",
|
||||
"bwrap", "--die-with-parent", "--new-session", "--unshare-pid",
|
||||
"--tmpfs", "/",
|
||||
"--dir", "/usr", "--ro-bind", "/usr", "/usr",
|
||||
"--symlink", "usr/bin", "/bin",
|
||||
@@ -820,21 +837,22 @@ def _bwrap_prefix(spec: ContainmentSpec) -> list[str]:
|
||||
"--symlink", "usr/bin", "/sbin",
|
||||
"--dir", "/etc", "--ro-bind", "/etc", "/etc",
|
||||
"--dir", "/tmp", "--tmpfs", "/tmp",
|
||||
"--dev-bind", "/dev", "/dev", "--proc", "/proc",
|
||||
"--dir", WORKSPACE_MOUNT, "--bind", spec.workspace, WORKSPACE_MOUNT,
|
||||
"--dev", "/dev", "--proc", "/proc",
|
||||
]
|
||||
mounted: tuple[str, ...] = ()
|
||||
for flag, paths in (("--ro-bind", spec.readonly_extra), ("--bind", spec.writable_extra)):
|
||||
for path in sorted(paths, key=lambda value: (value.count(os.sep), value)):
|
||||
args.extend(_dir_chain(path, mounted))
|
||||
args.extend((flag, path, path))
|
||||
mounted += (path,)
|
||||
# Mount workspace last so a read-only ancestor never hides its writable bind.
|
||||
args.extend(("--dir", WORKSPACE_MOUNT, "--bind", spec.workspace, WORKSPACE_MOUNT))
|
||||
# Preserve absolute workspace paths in generated scripts without exposing
|
||||
# a writable parent directory.
|
||||
workspace = os.path.realpath(spec.workspace)
|
||||
if workspace not in _RESERVED_BIND_DESTS and workspace not in {"/usr", "/etc"}:
|
||||
args.extend(_dir_chain(workspace))
|
||||
args.extend(_dir_chain(workspace, mounted))
|
||||
args.extend(("--bind", workspace, workspace))
|
||||
for path in spec.readonly_extra:
|
||||
args.extend(_dir_chain(path))
|
||||
args.extend(("--ro-bind", path, path))
|
||||
for path in spec.writable_extra:
|
||||
args.extend(_dir_chain(path))
|
||||
args.extend(("--bind", path, path))
|
||||
if spec.network == NETWORK_NONE:
|
||||
args.append("--unshare-net")
|
||||
args.extend(("--chdir", WORKSPACE_MOUNT))
|
||||
@@ -871,7 +889,8 @@ def _rlimit_preexec(grant: ContainmentGrant) -> Optional[Callable[[], None]]:
|
||||
return _apply
|
||||
|
||||
|
||||
def _launch_argv(grant: ContainmentGrant, command: Any, *, argv: bool) -> list[str]:
|
||||
def _launch_argv(grant: ContainmentGrant, command: Any, *, argv: bool,
|
||||
ready_marker: Optional[str] = None) -> list[str]:
|
||||
spec = grant.spec
|
||||
if argv:
|
||||
parts = [str(part) for part in command]
|
||||
@@ -894,6 +913,13 @@ def _launch_argv(grant: ContainmentGrant, command: Any, *, argv: bool) -> list[s
|
||||
)
|
||||
parts = [shell, "-c", text]
|
||||
if grant.mechanism == "bubblewrap":
|
||||
if ready_marker is not None:
|
||||
# This trusted wrapper runs only after all bwrap setup succeeds.
|
||||
# A launch-time bind/security failure must not claim containment.
|
||||
parts = ["/bin/sh", "-c",
|
||||
'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
|
||||
return parts
|
||||
|
||||
@@ -905,8 +931,8 @@ def _spawn_kwargs(grant: ContainmentGrant) -> dict[str, Any]:
|
||||
# reach it, and teardown walks the tree with taskkill /T.
|
||||
kwargs["creationflags"] = getattr(subprocess, "CREATE_NEW_PROCESS_GROUP", 0x00000200)
|
||||
return kwargs
|
||||
# setsid is what makes PROCESS_TREE real: without it a timeout kill reaches
|
||||
# the wrapper shell and nothing it backgrounded.
|
||||
# A separate group makes ordinary tree teardown possible. The PID
|
||||
# namespace, not setsid, prevents descendants from escaping containment.
|
||||
kwargs["start_new_session"] = True
|
||||
preexec = _rlimit_preexec(grant)
|
||||
if preexec is not None:
|
||||
@@ -978,61 +1004,101 @@ async def run(
|
||||
"backend does not own; it cannot be run locally"
|
||||
)
|
||||
missing = frozenset(grant.spec.required) - frozenset(grant.enforced)
|
||||
if missing and grant.mode == MODE_ENFORCING:
|
||||
raise ContainmentUnavailable(missing, grant.mechanism)
|
||||
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()
|
||||
missing |= frozenset(grant.spec.required) - provided
|
||||
overclaimed = frozenset(grant.enforced) - provided
|
||||
if (missing or overclaimed) and grant.mode == MODE_ENFORCING:
|
||||
raise ContainmentUnavailable(missing | overclaimed, grant.mechanism)
|
||||
|
||||
spec = grant.spec
|
||||
launch = _launch_argv(grant, command, argv=argv)
|
||||
marker = uuid.uuid4().hex if grant.mechanism == "bubblewrap" else None
|
||||
try:
|
||||
proc = await asyncio.create_subprocess_exec(
|
||||
launch = _launch_argv(grant, command, argv=argv, ready_marker=marker)
|
||||
spawning = asyncio.create_task(asyncio.create_subprocess_exec(
|
||||
*launch,
|
||||
stdin=asyncio.subprocess.PIPE if stdin is not None else asyncio.subprocess.DEVNULL,
|
||||
stdin=asyncio.subprocess.PIPE if stdin is not None or marker is not None else asyncio.subprocess.DEVNULL,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
cwd=spec.workspace,
|
||||
env=dict(spec.env),
|
||||
**_spawn_kwargs(grant),
|
||||
)
|
||||
))
|
||||
try:
|
||||
proc = await asyncio.shield(spawning)
|
||||
except asyncio.CancelledError:
|
||||
# Cancellation must not detach an OS spawn already in progress.
|
||||
# Recover its handle before propagating cancellation to the caller.
|
||||
try:
|
||||
proc = await _complete_cleanup(spawning, propagate_cancel=False)
|
||||
except Exception:
|
||||
release(grant, grace_s=0)
|
||||
else:
|
||||
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
|
||||
except BaseException:
|
||||
# Acquisition can precede a failed or cancelled spawn. A grant without
|
||||
# a child must not become a permanent restart orphan.
|
||||
release(grant, grace_s=0)
|
||||
if "proc" not in locals():
|
||||
release(grant, grace_s=0)
|
||||
raise
|
||||
from src.agent_runtime.journal import mark_operation_started
|
||||
mark_operation_started("subprocess", pid=proc.pid)
|
||||
# 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.
|
||||
pgid = None if IS_WINDOWS else (_pgid_of(proc.pid) or proc.pid)
|
||||
pgid = None if IS_WINDOWS else proc.pid
|
||||
live = replace(grant, pid=proc.pid, pgid=pgid)
|
||||
# The start token is what makes this record signallable by a *later*
|
||||
# process. Without it a restart reaper holds a pid and no way to tell
|
||||
# whether the pid is still this child or something the kernel has since
|
||||
# handed to a stranger; see src/process_ownership.py.
|
||||
_update_record(
|
||||
grant.id,
|
||||
pid=proc.pid,
|
||||
pgid=pgid,
|
||||
started_at=time.time(),
|
||||
start_token=process_ownership.capture(proc.pid)["start_token"],
|
||||
)
|
||||
try:
|
||||
_update_record(grant.id, pid=proc.pid, pgid=pgid, started_at=time.time(),
|
||||
start_token=None if marker is not None else process_ownership.capture(proc.pid)["start_token"],
|
||||
containment_ready=marker is None, execution_started=marker is None)
|
||||
from src.agent_runtime.journal import mark_operation_started
|
||||
mark_operation_started("subprocess", pid=proc.pid)
|
||||
except BaseException:
|
||||
await _complete_cleanup(_release_awaited(live, proc), propagate_cancel=False)
|
||||
raise
|
||||
|
||||
out_buf: list[str] = []
|
||||
err_buf: list[str] = []
|
||||
out_budget = [int(spec.max_output_bytes)]
|
||||
err_budget = [int(spec.max_output_bytes)]
|
||||
started = time.time()
|
||||
readers = [
|
||||
asyncio.create_task(_drain(proc.stdout, out_buf, out_budget, output_cb)),
|
||||
asyncio.create_task(_drain(proc.stderr, err_buf, err_budget, output_cb)),
|
||||
]
|
||||
readers = [asyncio.create_task(_drain(proc.stderr, err_buf, err_budget, output_cb))]
|
||||
ready = marker is None
|
||||
execution_started = marker is None
|
||||
async def _wait() -> None:
|
||||
nonlocal ready, execution_started
|
||||
if marker is not None:
|
||||
expected = (marker + "\n").encode("ascii")
|
||||
try:
|
||||
receipt = await proc.stdout.readexactly(len(expected))
|
||||
except (asyncio.IncompleteReadError, OSError):
|
||||
receipt = b""
|
||||
if receipt != expected:
|
||||
raise ContainmentUnavailable(spec.required, grant.mechanism)
|
||||
# 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"])
|
||||
if hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal"):
|
||||
try:
|
||||
proc._ody_pidfd = os.pidfd_open(proc.pid)
|
||||
except OSError:
|
||||
raise ContainmentUnavailable(spec.required, grant.mechanism) from None
|
||||
ready = True
|
||||
_update_record(grant.id, containment_ready=True, execution_started=True)
|
||||
execution_started = True
|
||||
proc.stdin.write(expected)
|
||||
await proc.stdin.drain()
|
||||
readers.append(asyncio.create_task(_drain(proc.stdout, out_buf, out_budget, output_cb)))
|
||||
# Pipe backpressure is execution time too. Feeding a child that never
|
||||
# reads stdin must remain inside the same timeout/cancellation scope.
|
||||
if stdin is not None and proc.stdin is not None:
|
||||
if (stdin is not None or marker is not None) and proc.stdin is not None:
|
||||
try:
|
||||
proc.stdin.write(stdin)
|
||||
await proc.stdin.drain()
|
||||
if stdin is not None:
|
||||
proc.stdin.write(stdin)
|
||||
await proc.stdin.drain()
|
||||
except (BrokenPipeError, ConnectionResetError):
|
||||
pass
|
||||
finally:
|
||||
@@ -1052,34 +1118,42 @@ async def run(
|
||||
progress_task = asyncio.create_task(_progress()) if progress_cb else None
|
||||
timed_out = False
|
||||
outcome: Optional[ReleaseOutcome] = None
|
||||
async def _finish() -> ReleaseOutcome:
|
||||
try:
|
||||
return await _release_awaited(live, proc)
|
||||
finally:
|
||||
if progress_task is not None:
|
||||
progress_task.cancel()
|
||||
try:
|
||||
await progress_task
|
||||
except (asyncio.CancelledError, Exception):
|
||||
pass
|
||||
for task in readers:
|
||||
try:
|
||||
await asyncio.wait_for(asyncio.shield(task), timeout=1)
|
||||
except (asyncio.TimeoutError, Exception):
|
||||
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)
|
||||
try:
|
||||
try:
|
||||
await asyncio.wait_for(_wait(), timeout=spec.wall_clock_s)
|
||||
except asyncio.TimeoutError:
|
||||
if not ready:
|
||||
raise ContainmentUnavailable(spec.required, grant.mechanism)
|
||||
timed_out = True
|
||||
outcome = await _release_awaited(live, proc)
|
||||
except asyncio.CancelledError:
|
||||
await _release_awaited(live, proc)
|
||||
raise
|
||||
except BaseException as exc:
|
||||
exc.containment_established = ready
|
||||
exc.containment_executed = execution_started
|
||||
raise
|
||||
finally:
|
||||
if progress_task is not None:
|
||||
progress_task.cancel()
|
||||
try:
|
||||
await progress_task
|
||||
except (asyncio.CancelledError, Exception):
|
||||
pass
|
||||
for task in readers:
|
||||
try:
|
||||
await asyncio.wait_for(task, timeout=1)
|
||||
except (asyncio.TimeoutError, asyncio.CancelledError, Exception):
|
||||
task.cancel()
|
||||
|
||||
if not timed_out:
|
||||
# Even a clean exit goes through teardown: a command that backgrounded
|
||||
# something leaves the group populated, and leaving it running is the
|
||||
# leak this boundary exists to close.
|
||||
outcome = await _release_awaited(live, proc)
|
||||
else:
|
||||
# Clean exit, partial initialization, timeout, and cancellation share
|
||||
# the same teardown. Repeated cancellation cannot skip escalation.
|
||||
outcome = await _complete_cleanup(_finish())
|
||||
if timed_out:
|
||||
_update_record(grant.id, timed_out=True)
|
||||
|
||||
return ContainmentResult(
|
||||
@@ -1093,6 +1167,23 @@ async def run(
|
||||
)
|
||||
|
||||
|
||||
async def _complete_cleanup(awaitable, *, propagate_cancel: bool = True):
|
||||
"""Finish ownership cleanup despite further cancellation, then propagate it."""
|
||||
task = asyncio.ensure_future(awaitable)
|
||||
cancelled = False
|
||||
while not task.done():
|
||||
try:
|
||||
await asyncio.shield(task)
|
||||
except asyncio.CancelledError:
|
||||
if task.cancelled():
|
||||
raise
|
||||
cancelled = True
|
||||
result = task.result()
|
||||
if cancelled and propagate_cancel:
|
||||
raise asyncio.CancelledError
|
||||
return result
|
||||
|
||||
|
||||
# ── release ─────────────────────────────────────────────────────────────────
|
||||
# core.platform_compat.kill_process_tree delegates here as well. Native tools,
|
||||
# detached jobs and compatibility callers share escalation and death probes.
|
||||
@@ -1429,32 +1520,63 @@ async def _release_awaited(
|
||||
has entirely exited — turning every timeout into a false "survivors"
|
||||
report.
|
||||
"""
|
||||
writer = getattr(proc, "stdin", None)
|
||||
if writer is not None:
|
||||
writer.close()
|
||||
if hasattr(writer, "wait_closed"):
|
||||
try:
|
||||
await asyncio.wait_for(writer.wait_closed(), timeout=1)
|
||||
except (OSError, asyncio.TimeoutError):
|
||||
pass
|
||||
if IS_WINDOWS:
|
||||
return release(grant, grace_s=grace_s)
|
||||
|
||||
pid, pgid = grant.pid, grant.pgid
|
||||
_signal_tree(pid, pgid, signal.SIGTERM)
|
||||
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)
|
||||
def gone():
|
||||
if pidfd is not None:
|
||||
return bool(select.select([pidfd], [], [], 0)[0])
|
||||
return _tree_gone(pid, pgid)
|
||||
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:
|
||||
_signal_tree(pid, pgid, sig)
|
||||
send(signal.SIGTERM)
|
||||
try:
|
||||
await asyncio.wait_for(proc.wait(), timeout=max(grace_s, 0.05))
|
||||
except (asyncio.TimeoutError, ProcessLookupError):
|
||||
pass
|
||||
deadline = time.monotonic() + max(grace_s, 0.0)
|
||||
while time.monotonic() < deadline and not _tree_gone(pid, pgid):
|
||||
while time.monotonic() < deadline and not gone():
|
||||
await asyncio.sleep(_DEATH_POLL_S)
|
||||
|
||||
escalated = False
|
||||
if not _tree_gone(pid, pgid):
|
||||
if not gone():
|
||||
escalated = True
|
||||
_signal_tree(pid, pgid, signal.SIGKILL)
|
||||
send(signal.SIGKILL)
|
||||
try:
|
||||
await asyncio.wait_for(proc.wait(), timeout=1.0)
|
||||
except (asyncio.TimeoutError, ProcessLookupError):
|
||||
pass
|
||||
deadline = time.monotonic() + 1.0
|
||||
while time.monotonic() < deadline and not _tree_gone(pid, pgid):
|
||||
while time.monotonic() < deadline and not gone():
|
||||
await asyncio.sleep(_DEATH_POLL_S)
|
||||
|
||||
outcome = _outcome_for(grant, dead=_tree_gone(pid, pgid), escalated=escalated)
|
||||
outcome = _outcome_for(grant, dead=gone(), escalated=escalated)
|
||||
_finish_release(grant, outcome)
|
||||
return outcome
|
||||
|
||||
|
||||
@@ -60,8 +60,8 @@ async def supervise(payload: dict) -> None:
|
||||
output, code = f"background execution failed: {type(exc).__name__}: {exc}\n", 1
|
||||
report = {"containment": grant.to_dict(), "teardown": record.get("release") or {"dead": False},
|
||||
"output_truncated": False}
|
||||
report["containment"]["executed"] = bool(record.get("pid"))
|
||||
if not record.get("pid"):
|
||||
report["containment"]["executed"] = bool(record.get("execution_started"))
|
||||
if not record.get("containment_ready"):
|
||||
report["containment"].update(contained=False, enforced=[])
|
||||
if isinstance(exc, containment.ContainmentUnavailable):
|
||||
report.update(containment.unavailable_tool_result(exc, tool="bash"))
|
||||
|
||||
@@ -39,7 +39,7 @@ Token granularity, stated because it bounds the guarantee
|
||||
======== ============================= ===============
|
||||
Host Source Resolution
|
||||
======== ============================= ===============
|
||||
Linux ``/proc/<pid>/stat`` field 22 ~10 ms (1 tick)
|
||||
Linux boot ID + stat field 22 ~10 ms (1 tick)
|
||||
macOS ``ps -o lstart=`` 1 s
|
||||
Windows ``GetProcessTimes`` 100 ns
|
||||
======== ============================= ===============
|
||||
@@ -145,9 +145,18 @@ def _procfs_token(pid: int) -> Optional[str]:
|
||||
_, _, rest = raw.rpartition(")")
|
||||
fields = rest.split()
|
||||
try:
|
||||
return f"procfs:{fields[_PROC_STAT_STARTTIME_INDEX]}"
|
||||
ticks = fields[_PROC_STAT_STARTTIME_INDEX]
|
||||
except IndexError:
|
||||
raise InspectionUnavailable(f"/proc/{pid}/stat (unexpected layout)") from None
|
||||
try:
|
||||
boot = (PROC_ROOT / "sys/kernel/random/boot_id").read_text(encoding="ascii").strip()
|
||||
except OSError as exc:
|
||||
raise InspectionUnavailable(f"boot identity ({exc})") from exc
|
||||
if not boot:
|
||||
raise InspectionUnavailable("boot identity (empty)")
|
||||
# A persisted PID/start-tick pair can recur after reboot. Bind it to the
|
||||
# boot as well; older receipts cannot authorize a signal on a new boot.
|
||||
return f"procfs:{boot}:{ticks}"
|
||||
|
||||
|
||||
def _ps_token(pid: int) -> Optional[str]:
|
||||
|
||||
@@ -232,15 +232,18 @@ def reap_legacy_agent_tmux() -> Dict[str, Any]:
|
||||
if process_ownership.verify(pid, identities[pid]) != process_ownership.OWNED:
|
||||
continue
|
||||
spec = containment.ContainmentSpec(workspace=os.getcwd(), env={}, wall_clock_s=1,
|
||||
required=frozenset({containment.PROCESS_TREE}))
|
||||
required=frozenset())
|
||||
grant = containment.ContainmentGrant(
|
||||
id=uuid.uuid4().hex[:12], mechanism="process_group", workspace=spec.workspace,
|
||||
enforced=frozenset({containment.PROCESS_TREE}), degraded=(), unenforced_required=(),
|
||||
enforced=frozenset(), degraded=(containment.FILESYSTEM, containment.PROCESS_TREE), unenforced_required=(),
|
||||
owner=f"legacy-tmux:{session_id}", mode=containment.MODE_ENFORCING,
|
||||
spec=spec, pid=pid, pgid=containment._pgid_of(pid),
|
||||
)
|
||||
containment._write_record(grant)
|
||||
containment._update_record(grant.id, lifetime="cleanup", start_token=identities[pid])
|
||||
receipt = containment._load_records().get(grant.id, {})
|
||||
if receipt.get("start_token") != identities[pid] or receipt.get("lifetime") != "cleanup":
|
||||
raise RuntimeError("legacy tmux cleanup receipt was not persisted")
|
||||
tracked.append((grant, identities[pid]))
|
||||
dead = True
|
||||
for grant, token in tracked:
|
||||
|
||||
Reference in New Issue
Block a user