From 7892f7f650b6feb845667c27544aaadb51a8bf33 Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Thu, 1 Oct 2026 21:23:21 +0100 Subject: [PATCH] fix(runtime): contain and supervise detached Bash jobs (ODY-145) --- core/atomic_io.py | 41 ++++- core/platform_compat.py | 6 +- src/agent_tools/bg_job_tools.py | 3 + src/bg_jobs.py | 185 +++++++++++++-------- src/containment.py | 54 ++++-- src/containment_worker.py | 77 +++++++++ src/process_reaper.py | 7 + src/tool_execution.py | 21 ++- tests/test_background_containment.py | 126 ++++++++++++++ tests/test_bg_job_tools.py | 8 +- tests/test_native_execution_containment.py | 11 ++ 11 files changed, 438 insertions(+), 101 deletions(-) create mode 100644 src/containment_worker.py create mode 100644 tests/test_background_containment.py diff --git a/core/atomic_io.py b/core/atomic_io.py index 831b90848..f2f18b409 100644 --- a/core/atomic_io.py +++ b/core/atomic_io.py @@ -16,9 +16,48 @@ from __future__ import annotations import json import os import uuid +import functools +import threading from typing import Any, Optional +_STORE_LOCKS: dict[str, threading.RLock] = {} +_STORE_LOCKS_GUARD = threading.Lock() + + +def store_transaction(path_factory): + """Serialize a JSON read/modify/write across runtime threads and processes.""" + def decorate(function): + @functools.wraps(function) + def locked(*args, **kwargs): + path = os.path.abspath(str(path_factory())) + ".lock" + with _STORE_LOCKS_GUARD: + lock = _STORE_LOCKS.setdefault(path, threading.RLock()) + with lock: + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "a+b") as handle: + if os.name == "nt": + import msvcrt + if os.fstat(handle.fileno()).st_size == 0: + handle.write(b"0") + handle.flush() + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_LOCK, 1) + else: + import fcntl + fcntl.flock(handle, fcntl.LOCK_EX) + try: + return function(*args, **kwargs) + finally: + if os.name == "nt": + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) + else: + fcntl.flock(handle, fcntl.LOCK_UN) + return locked + return decorate + + def atomic_write_json(path: str, data: Any, *, indent: Optional[int] = None) -> None: """Atomically persist `data` as JSON at `path`. @@ -64,4 +103,4 @@ def atomic_write_text(path: str, text: str) -> None: try: os.unlink(tmp) except OSError: - pass \ No newline at end of file + pass diff --git a/core/platform_compat.py b/core/platform_compat.py index be6d30b87..75122c92c 100644 --- a/core/platform_compat.py +++ b/core/platform_compat.py @@ -124,7 +124,7 @@ def pid_alive(pid: Optional[int]) -> bool: return True # EPERM and other inspection failures are not ESRCH. -def kill_process_tree(pid: Optional[int]): +def kill_process_tree(pid: Optional[int], *, start_token=None, pgid=None, require_identity=False): """Use the runtime's shared escalating teardown and return verified death. Callers retaining durable PIDs must validate their recorded identity before @@ -140,9 +140,9 @@ def kill_process_tree(pid: Optional[int]): id="", mechanism="windows_tree" if IS_WINDOWS else "process_group", workspace=spec.workspace, enforced=frozenset(), degraded=(), unenforced_required=(), owner="compatibility", mode=containment.CONTAINMENT_MODE, - spec=spec, pid=int(pid), pgid=containment._pgid_of(int(pid)), + spec=spec, pid=int(pid), pgid=pgid or containment._pgid_of(int(pid)), ) - return containment.release(grant) + return containment.release(grant, start_token=start_token, require_identity=require_identity) # ── Shell / executable resolution ─────────────────────────────────────────── diff --git a/src/agent_tools/bg_job_tools.py b/src/agent_tools/bg_job_tools.py index a29e813cc..692f459a8 100644 --- a/src/agent_tools/bg_job_tools.py +++ b/src/agent_tools/bg_job_tools.py @@ -87,6 +87,9 @@ class ManageBgJobsTool: if rec.get("status") != "running": return {"output": f"Job `{job_id}` already {_status_label(rec)}; nothing to kill.", "exit_code": 0} killed = bg_jobs.kill(job_id) + if not killed or not killed.get("killed"): + return {"error": f"Could not verify termination of background job `{job_id}`.", + "exit_code": 1, "teardown": (killed or {}).get("teardown")} return {"output": f"Killed background job `{job_id}` ({(killed or {}).get('command', '').splitlines()[0][:80]}).", "exit_code": 0} out = rec.get("output") or "(no output yet)" diff --git a/src/bg_jobs.py b/src/bg_jobs.py index 7d1d44cf5..ec0d9b828 100644 --- a/src/bg_jobs.py +++ b/src/bg_jobs.py @@ -22,18 +22,16 @@ from __future__ import annotations import json import os -import shlex +import sys import subprocess import time import uuid from pathlib import Path from typing import Any, Dict, List, Optional -from core.atomic_io import atomic_write_json +from core.atomic_io import atomic_write_json, store_transaction from core.platform_compat import ( detached_popen_kwargs, - find_bash, - git_bash_path, kill_process_tree, pid_alive, ) @@ -53,6 +51,7 @@ _MAX_OUTPUT_CHARS = 16000 # files) is kept before pruning, so neither the store nor data/bg_jobs/ grows # without bound. The agent has already consumed the result by then. _RETENTION_S = 3600 # 1 hour after follow-up +_LIVE_PROCS: dict[int, subprocess.Popen] = {} def _load() -> Dict[str, Dict[str, Any]]: @@ -79,66 +78,50 @@ def _pid_alive(pid: Optional[int]) -> bool: return pid_alive(pid) +@store_transaction(lambda: _STORE) def launch(command: str, session_id: str, cwd: Optional[str] = None, - max_runtime_s: int = DEFAULT_MAX_RUNTIME_S) -> Dict[str, Any]: + max_runtime_s: int = DEFAULT_MAX_RUNTIME_S, env: Optional[dict] = None) -> Dict[str, Any]: """Launch `command` detached. Returns the job record (status='running'). - Output + the final exit code are written to files so status survives a - server restart. The process is put in its own session (setsid) so it - outlives the request/stream that started it. + A trusted detached supervisor owns the shared containment runner, output, + wall clock and exit metadata, independently of the request/server lifetime. """ _JOBS_DIR.mkdir(parents=True, exist_ok=True) job_id = uuid.uuid4().hex[:12] log_path = _JOBS_DIR / f"{job_id}.log" exit_path = _JOBS_DIR / f"{job_id}.exit" - # The user command goes in its OWN script file, run as a child `bash`. This - # is what isolates it: an `exit` inside it only ends that child (so the - # wrapper still records the exit code), and — unlike textually wrapping the - # command in `( … )` — the wrapper can't be broken by an unbalanced paren or - # a trailing line-continuation in the command. `$?` is the child's real - # exit status. - bash = find_bash() - if bash: - # POSIX, or Windows with Git Bash/WSL. The user command goes in its OWN - # script file, run as a child `bash` — an `exit` inside it only ends - # that child (so the wrapper still records the exit code), and an - # unbalanced paren / trailing line-continuation in the command can't - # break the wrapper. `$?` is the child's real exit status. Paths are - # emitted as POSIX (forward-slash) + shell-quoted so Git Bash on Windows - # handles drive paths and spaces correctly. - cmd_path = _JOBS_DIR / f"{job_id}.cmd.sh" - cmd_path.write_text(command + "\n", encoding="utf-8") - lp, xp, cp = (shlex.quote(git_bash_path(p)) for p in (log_path, exit_path, cmd_path)) - script_path = _JOBS_DIR / f"{job_id}.sh" - script_path.write_text( - f"bash {cp} > {lp} 2>&1\n" - f"echo $? > {xp}\n", - encoding="utf-8", - ) - argv = [bash, str(script_path)] - else: - # Windows without any bash installed: cmd.exe wrapper. The command runs - # in its own child .cmd so %ERRORLEVEL% is the command's real exit code. - child_path = _JOBS_DIR / f"{job_id}.child.cmd" - child_path.write_text("@echo off\r\n" + command + "\r\n", encoding="utf-8") - script_path = _JOBS_DIR / f"{job_id}.cmd" - script_path.write_text( - "@echo off\r\n" - f'call "{child_path}" > "{log_path}" 2>&1\r\n' - f'echo %ERRORLEVEL%> "{exit_path}"\r\n', - encoding="utf-8", - ) - argv = [os.environ.get("ComSpec", "cmd.exe"), "/c", str(script_path)] - - proc = subprocess.Popen( - argv, - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - stdin=subprocess.DEVNULL, - cwd=cwd or None, - **detached_popen_kwargs(), # detach from the request lifecycle (setsid / DETACHED_PROCESS) - ) + from src import containment + from src.agent_tools.subprocess_tools import _owned_spec, _replace_workspace_alias + spec = _owned_spec(cwd or os.getcwd(), env, max_runtime_s) + grant = containment.acquire(spec, owner=f"bg:{session_id}") + bounded_command = command + if containment.FILESYSTEM not in grant.enforced: + bounded_command = _replace_workspace_alias(command, grant.workspace) + result_path = _JOBS_DIR / f"{job_id}.result.json" + payload = { + "store_path": str(containment._store_path().resolve()), + "grant": {**grant.to_dict(), "owner": grant.owner}, + "spec": { + "workspace": spec.workspace, "env": dict(spec.env), "wall_clock_s": spec.wall_clock_s, + "required": sorted(spec.required), "network": spec.network, + "readonly_extra": list(spec.readonly_extra), "writable_extra": list(spec.writable_extra), + "max_output_bytes": spec.max_output_bytes, + }, + "command": bounded_command, "log_path": str(log_path.resolve()), + "result_path": str(result_path.resolve()), "exit_path": str(exit_path.resolve()), + } + try: + with open(log_path, "ab") as bootstrap_log: + proc = subprocess.Popen( + [sys.executable, str(Path(containment.__file__).with_name("containment_worker.py"))], + stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=bootstrap_log, + cwd=str(Path(containment.__file__).resolve().parent.parent), + **detached_popen_kwargs(), + ) + except BaseException: + containment.release(grant, grace_s=0) + raise rec = { "id": job_id, @@ -153,15 +136,31 @@ def launch(command: str, session_id: str, cwd: Optional[str] = None, "followed_up": False, # has the agent been re-invoked with the result? "log_path": str(log_path), "exit_path": str(exit_path), + "result_path": str(result_path), + "containment_id": grant.id, + "containment": {**grant.to_dict(), "contained": False, "enforced": [], "pending": True, "executed": False}, + "pgid": None if os.name == "nt" else proc.pid, # Identity, not just a slot. The pid above is reused by the kernel, and # this record outlives the process and the server; the token is what a # later run compares before it signals anything. See # src/process_ownership.py. "start_token": process_ownership.capture(proc.pid)["start_token"], } - jobs = _load() - jobs[job_id] = rec - _save(jobs) + try: + containment._update_record(grant.id, lifetime="background", supervisor_pid=proc.pid, + supervisor_token=rec["start_token"]) + jobs = _load() + jobs[job_id] = rec + _save(jobs) + # The supervisor cannot execute until the identity and job record are durable. + proc.stdin.write(json.dumps(payload).encode("utf-8")) + proc.stdin.close() + except BaseException: + kill_process_tree(proc.pid) + proc.wait(timeout=5) + containment.release(grant, grace_s=0) + raise + _LIVE_PROCS[proc.pid] = proc return rec @@ -194,10 +193,14 @@ def _prune(jobs: Dict[str, Dict[str, Any]], now: float) -> bool: return bool(stale) +@store_transaction(lambda: _STORE) def refresh() -> Dict[str, Dict[str, Any]]: """Reconcile every running job against disk. Marks done/failed (incl. timeout). Idempotent — safe to call from a poll loop. Returns the store.""" jobs = _load() + for pid, proc in list(_LIVE_PROCS.items()): + if proc.poll() is not None: + _LIVE_PROCS.pop(pid, None) changed = False now = time.time() for rec in jobs.values(): @@ -212,13 +215,24 @@ def refresh() -> Dict[str, Dict[str, Any]]: rec["exit_code"] = code rec["status"] = "done" if code == 0 else "failed" rec["ended_at"] = now + if rec.get("result_path"): + try: + report = json.loads(Path(rec["result_path"]).read_text(encoding="utf-8")) + rec.update(report) + except (OSError, ValueError): + rec["status"], rec["exit_code"] = "failed", 1 + rec["result_unavailable"] = True changed = True elif (now - rec.get("started_at", now)) > rec.get("max_runtime_s", DEFAULT_MAX_RUNTIME_S): # Runaway / stuck — reap it but STILL surface a follow-up. - _kill(rec.get("pid")) - rec["status"] = "failed" - rec["exit_code"] = -1 - rec["ended_at"] = now + outcome = _kill_record(rec) + rec["teardown"] = outcome.to_dict() + if outcome.dead: + rec["status"] = "failed" + rec["exit_code"] = -1 + rec["ended_at"] = now + else: + rec["kill_failed"] = True rec["timed_out"] = True changed = True elif not _pid_alive(rec.get("pid")) and not exit_path.exists(): @@ -236,9 +250,33 @@ def refresh() -> Dict[str, Dict[str, Any]]: return jobs -def _kill(pid: Optional[int]) -> None: +def _kill(pid: Optional[int], **kwargs): # Cross-platform process-tree teardown (POSIX killpg / Windows taskkill /T). - kill_process_tree(pid) + return kill_process_tree(pid, **kwargs) + + +def _kill_record(rec): + from src import containment + verdict = process_ownership.verify(rec.get("pid"), rec.get("start_token")) + if verdict in (process_ownership.FOREIGN, process_ownership.UNVERIFIABLE): + return containment.ReleaseOutcome(dead=False, escalated=False, ownership=verdict) + if rec.get("containment_id"): + record = containment._load_records().get(rec["containment_id"]) + if record and record.get("pid"): + outcome = containment.reap_record(record) + if not outcome.dead: + return outcome + outcome = _kill(rec.get("pid"), start_token=rec.get("start_token"), + pgid=rec.get("pgid"), require_identity=True) + proc = _LIVE_PROCS.get(rec.get("pid")) + if proc and outcome.dead: + proc.wait(timeout=5) + _LIVE_PROCS.pop(proc.pid, None) + if outcome.dead and rec.get("containment_id"): + record = containment._load_records().get(rec["containment_id"]) + if record and not record.get("pid"): + containment.reap_record(record) + return outcome def pending_followups() -> List[Dict[str, Any]]: @@ -249,6 +287,7 @@ def pending_followups() -> List[Dict[str, Any]]: if r.get("status") in ("done", "failed") and not r.get("followed_up")] +@store_transaction(lambda: _STORE) def mark_followed_up(job_id: str) -> None: jobs = _load() if job_id in jobs: @@ -269,6 +308,7 @@ def list_for_session(session_id: str) -> List[Dict[str, Any]]: return [r for r in refresh().values() if r.get("session_id") == session_id] +@store_transaction(lambda: _STORE) def kill(job_id: str) -> Optional[Dict[str, Any]]: """Terminate a running job's process tree and mark it killed. Returns the updated record, or None if the id is unknown. Idempotent: a job that already @@ -279,16 +319,21 @@ def kill(job_id: str) -> Optional[Dict[str, Any]]: if rec is None: return None if rec.get("status") == "running": - _kill(rec.get("pid")) - rec["status"] = "failed" - rec["exit_code"] = -1 - rec["ended_at"] = time.time() - rec["killed"] = True - rec["followed_up"] = True + outcome = _kill_record(rec) + rec["teardown"] = outcome.to_dict() + if outcome.dead: + rec["status"] = "failed" + rec["exit_code"] = -1 + rec["ended_at"] = time.time() + rec["killed"] = True + rec["followed_up"] = True + else: + rec["kill_failed"] = True _save(jobs) return rec +@store_transaction(lambda: _STORE) def disown_unverified() -> Dict[str, Any]: """Stop tracking running jobs whose process can no longer be proven ours. diff --git a/src/containment.py b/src/containment.py index b7213b173..3634a434f 100644 --- a/src/containment.py +++ b/src/containment.py @@ -49,6 +49,7 @@ constant is reversible in a way that breaking every host is not. from __future__ import annotations import asyncio +import codecs import json import logging import os @@ -63,7 +64,7 @@ from pathlib import Path, PurePosixPath from types import MappingProxyType from typing import Any, Awaitable, Callable, Mapping, Optional -from core.atomic_io import atomic_write_json +from core.atomic_io import atomic_write_json, store_transaction from core.platform_compat import IS_WINDOWS, find_bash, pid_alive from src import process_ownership @@ -439,6 +440,7 @@ def _prune(records: dict[str, dict[str, Any]]) -> dict[str, dict[str, Any]]: return kept +@store_transaction(lambda: _store_path()) def _write_record(grant: ContainmentGrant) -> None: records = _prune(_load_records()) records[grant.id] = { @@ -465,6 +467,7 @@ def _write_record(grant: ContainmentGrant) -> None: _save_records(records) +@store_transaction(lambda: _store_path()) def _update_record(grant_id: str, **fields: Any) -> None: records = _load_records() record = records.get(grant_id) @@ -487,6 +490,7 @@ def active_grants() -> list[dict[str, Any]]: ] +@store_transaction(lambda: _store_path()) def forget(grant_id: str) -> None: """Drop a record outright. For a reaper that has finished with it.""" records = _load_records() @@ -908,7 +912,7 @@ def _spawn_kwargs(grant: ContainmentGrant) -> dict[str, Any]: return kwargs -async def _drain(stream, buffer: list[str], budget: list[int]) -> None: +async def _drain(stream, buffer: list[str], budget: list[int], output_cb=None) -> None: """Read a stream to EOF, keeping at most ``budget[0]`` bytes. Reading past the cap and discarding is deliberate: stopping the read would @@ -921,20 +925,30 @@ async def _drain(stream, buffer: list[str], budget: list[int]) -> None: """ if stream is None: return + decoder = codecs.getincrementaldecoder("utf-8")(errors="replace") + def emit(text): + if not text: + return + buffer.append(text) + if output_cb: + try: + output_cb(text) + except OSError: + budget[0] = -1 + logger.warning("containment: output sink failed", exc_info=True) while True: line = await stream.read(65536) if not line: + emit(decoder.decode(b"", final=True)) break if budget[0] < 0: continue - if len(line) <= budget[0]: - budget[0] -= len(line) - buffer.append(line.decode("utf-8", errors="replace")) - continue - chunk = line[: budget[0]] + chunk = line[:budget[0]] if chunk: - buffer.append(chunk.decode("utf-8", errors="replace")) - budget[0] = -1 + emit(decoder.decode(chunk)) + if budget[0] < 0: + continue + budget[0] = budget[0] - len(line) if len(line) <= budget[0] else -1 async def run( @@ -944,6 +958,7 @@ async def run( argv: bool = False, stdin: Optional[bytes] = None, progress_cb: Optional[Callable[[dict], Awaitable[None]]] = None, + output_cb: Optional[Callable[[str], None]] = None, ) -> ContainmentResult: """Execute inside an existing grant. @@ -1006,8 +1021,8 @@ async def run( err_budget = [int(spec.max_output_bytes)] started = time.time() readers = [ - asyncio.create_task(_drain(proc.stdout, out_buf, out_budget)), - asyncio.create_task(_drain(proc.stderr, err_buf, err_budget)), + asyncio.create_task(_drain(proc.stdout, out_buf, out_budget, output_cb)), + asyncio.create_task(_drain(proc.stderr, err_buf, err_budget, output_cb)), ] async def _wait() -> None: # Pipe backpressure is execution time too. Feeding a child that never @@ -1027,7 +1042,8 @@ async def run( await asyncio.sleep(2.0) if progress_cb: try: - await progress_cb({"elapsed_s": round(time.time() - started, 1)}) + tail = "\n".join(("".join(out_buf) + "".join(err_buf)).splitlines()[-12:])[-8192:] + await progress_cb({"elapsed_s": round(time.time() - started, 1), "tail": tail}) except Exception: pass @@ -1257,7 +1273,8 @@ def _ownership_gate( ) -def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: +def release(grant: ContainmentGrant, *, grace_s: float = 2.0, + start_token: Optional[str] = None, require_identity: bool = False) -> ReleaseOutcome: """Authoritative teardown: signal the group, escalate, then verify. Returns whether the tree is **observed** gone. A caller must not record a @@ -1276,12 +1293,12 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: # pid had to be recovered from the durable store is the restart case, and # there the pid is a *claim* about a process this run never started. recovered = pid is None - token: Optional[str] = None + token: Optional[str] = start_token if pid is None or (pgid is None and not IS_WINDOWS): record = _load_records().get(grant.id) or {} pid = pid if pid is not None else record.get("pid") pgid = pgid if pgid is not None else record.get("pgid") - token = record.get("start_token") + token = token or record.get("start_token") try: pid = int(pid) if pid else 0 except (TypeError, ValueError): @@ -1301,7 +1318,7 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: _finish_release(grant, outcome) return outcome - if recovered: + if recovered or require_identity: refusal = _ownership_gate(grant, pid, pgid, token) if refusal is not None: _finish_release(grant, refusal) @@ -1336,6 +1353,11 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: time.sleep(_DEATH_POLL_S) if not _tree_gone(pid, pgid, reap=True): escalated = True + if recovered or require_identity: + refusal = _ownership_gate(grant, pid, pgid, token) + if refusal is not None: + _finish_release(grant, refusal) + return refusal _signal_tree(pid, pgid, 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 diff --git a/src/containment_worker.py b/src/containment_worker.py new file mode 100644 index 000000000..d0b412f13 --- /dev/null +++ b/src/containment_worker.py @@ -0,0 +1,77 @@ +"""Trusted detached supervisor; command execution stays in containment.run.""" +from __future__ import annotations + +import asyncio +import json +import signal +import sys +import types +from pathlib import Path + +# Launch by absolute script path, so a task workspace cannot shadow src. +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) +# This supervisor needs atomic I/O and platform primitives, not core's chat +# facade (auth, database, LLM startup). Keep that facade out of the detached +# process without changing the application's normal imports. +core_package = types.ModuleType("core") +core_package.__path__ = [str(Path(__file__).resolve().parent.parent / "core")] +sys.modules["core"] = core_package + +from core.atomic_io import atomic_write_json, atomic_write_text +from src import containment + + +async def supervise(payload: dict) -> None: + containment._store_path = lambda: Path(payload["store_path"]) + data = payload["spec"] + data["required"] = frozenset(data["required"]) + spec = containment.ContainmentSpec(**data) + info = payload["grant"] + grant = containment.ContainmentGrant( + id=info["id"], mechanism=info["mechanism"], workspace=spec.workspace, + enforced=frozenset(info["enforced"]), degraded=tuple(info["degraded"]), + unenforced_required=tuple(info["unenforced_required"]), owner=info["owner"], + mode=info["mode"], spec=spec, + ) + task = asyncio.current_task() + loop = asyncio.get_running_loop() + if sys.platform != "win32": + loop.add_signal_handler(signal.SIGTERM, task.cancel) + loop.add_signal_handler(signal.SIGINT, task.cancel) + try: + with open(payload["log_path"], "w", encoding="utf-8") as log: + def capture(text): + log.write(text) + log.flush() + result = await containment.run(grant, payload["command"], output_cb=capture) + output = "" + code = 124 if result.timed_out else result.exit_code + if not result.release or not result.release.dead: + code = 1 + report = {"containment": result.grant.to_dict(), + "teardown": result.release.to_dict() if result.release else {"dead": False}, + "output_truncated": result.output_truncated, + "timed_out": result.timed_out} + report["containment"]["executed"] = True + if result.output_truncated: + output = "\n…[output truncated by containment capture limit]…\n" + except BaseException as exc: + record = containment._load_records().get(grant.id, {}) + 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"].update(contained=False, enforced=[]) + if isinstance(exc, containment.ContainmentUnavailable): + report.update(containment.unavailable_tool_result(exc, tool="bash")) + if output: + with open(payload["log_path"], "a", encoding="utf-8") as log: + log.write(output) + atomic_write_json(payload["result_path"], report) + # Publish completion last: refresh must never see an exit without metadata. + atomic_write_text(payload["exit_path"], str(code if code is not None else 1)) + + +if __name__ == "__main__": + asyncio.run(supervise(json.load(sys.stdin))) diff --git a/src/process_reaper.py b/src/process_reaper.py index 8eddf5db2..706acaeee 100644 --- a/src/process_reaper.py +++ b/src/process_reaper.py @@ -71,6 +71,13 @@ def reap_containment_grants() -> Dict[str, Any]: containment.forget(grant_id) report["already_gone"] += 1 continue + if record.get("lifetime") == "background" and process_ownership.verify( + record.get("supervisor_pid"), record.get("supervisor_token"), + ) == process_ownership.OWNED: + # Detached jobs deliberately survive a server restart. Their + # supervisor owns the wall clock and teardown, independently. + report["background_kept"] = report.get("background_kept", 0) + 1 + continue verdict = process_ownership.verify_record(record) if verdict == process_ownership.GONE: if containment._group_present(record.get("pgid")): diff --git a/src/tool_execution.py b/src/tool_execution.py index b5976900a..fdfdddd6d 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -1187,6 +1187,10 @@ def _split_bg_marker(content: str): return False, content +def _agent_subprocess_env() -> dict: + return {**os.environ, "TERM": "xterm-256color", "COLUMNS": "120", "LINES": "40", "HOME": _AGENT_WORKDIR} + + async def _direct_fallback( tool: str, content: str, @@ -1197,13 +1201,7 @@ async def _direct_fallback( disabled_tools: Optional[set] = None, tool_policy: Optional[ToolPolicy] = None, ) -> Optional[Dict]: - _subproc_env = { - **os.environ, - "TERM": "xterm-256color", - "COLUMNS": "120", - "LINES": "40", - "HOME": _AGENT_WORKDIR, - } + _subproc_env = _agent_subprocess_env() try: ctx = { @@ -1663,7 +1661,11 @@ async def _execute_tool_block_impl( if _is_bg and _bg_cmd: from src import bg_jobs mark_dispatch() - rec = bg_jobs.launch(_bg_cmd, session_id=session_id, cwd=agent_cwd()) + from src import containment + try: + rec = bg_jobs.launch(_bg_cmd, session_id=session_id, cwd=agent_cwd(), env=_agent_subprocess_env()) + except containment.ContainmentUnavailable as exc: + return "bash (background): containment unavailable", containment.unavailable_tool_result(exc, tool="bash") # Only this server launch may seal detached-job authority; a # handler/bridge output carrying a job id is not a grant source. save_background_authority(rec["id"], active_request_authority()) @@ -1673,7 +1675,7 @@ async def _execute_tool_block_impl( "output": ( f"Started background job `{rec['id']}`. It is running detached; " f"do NOT wait for it or poll it. You will be automatically re-invoked " - f"with its full output when it finishes. Continue with other work, or " + f"with its captured output and any capture limit when it finishes. Continue with other work, or " f"end your turn now and resume when the result arrives. If the user " f"later asks to check progress or stop it, call the manage_bg_jobs " f"tool yourself (output or kill); do not tell them to run a tool " @@ -1681,6 +1683,7 @@ async def _execute_tool_block_impl( ), "exit_code": 0, "bg_job_id": rec["id"], + "containment": rec.get("containment"), } logger.info(f"Tool executed: {desc} -> bg job {rec['id']}") return desc, result diff --git a/tests/test_background_containment.py b/tests/test_background_containment.py new file mode 100644 index 000000000..52d01f483 --- /dev/null +++ b/tests/test_background_containment.py @@ -0,0 +1,126 @@ +"""Detached Bash uses the same boundary; restart and kill retain ownership.""" +import asyncio +import os +import time +from collections import namedtuple + +import pytest + +from src import bg_jobs, containment, process_ownership, process_reaper, tool_execution +from src.tool_execution import NO_TOOL_SECURITY_CONTEXT +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture +def jobs(tmp_path, monkeypatch): + monkeypatch.setattr(bg_jobs, "_JOBS_DIR", tmp_path / "jobs") + monkeypatch.setattr(bg_jobs, "_STORE", tmp_path / "jobs.json") + monkeypatch.setattr(containment, "_store_path", lambda: tmp_path / "grants.json") + monkeypatch.setattr(containment, "CONTAINMENT_MODE", containment.MODE_REPORT_ONLY) + monkeypatch.setattr(containment, "MECHANISMS", tuple(m for m in containment.MECHANISMS if m.name == "process_group")) + monkeypatch.setattr(tool_execution, "_owner_is_admin", lambda owner: True) + launched = [] + yield tmp_path, launched + for record in launched: + current = bg_jobs.get(record["id"]) + if current and current["status"] == "running": + bg_jobs.kill(record["id"]) + proc = bg_jobs._LIVE_PROCS.pop(record["pid"], None) + if proc: + proc.wait(timeout=8) + + +def finished(job_id): + deadline = time.monotonic() + 10 + while time.monotonic() < deadline: + record = bg_jobs.get(job_id) + if record["status"] != "running": + return record + time.sleep(0.03) + pytest.fail("background job did not finish") + + +def test_detached_execution_owns_boundary_and_reports_death(jobs): + path, launched = jobs + record = bg_jobs.launch("printf captured", "chat", cwd=str(path)) + launched.append(record) + result = finished(record["id"]) + assert result["output"] == "captured" + assert result["exit_code"] == 0 + assert result["containment"]["mechanism"] == "process_group" + assert result["containment"]["contained"] is False + assert result["teardown"]["dead"] is True + + +async def test_bg_marker_refuses_without_spawning_and_authority_still_gates(jobs, monkeypatch): + path, _ = jobs + monkeypatch.setattr(containment, "CONTAINMENT_MODE", containment.MODE_ENFORCING) + monkeypatch.setattr(bg_jobs.subprocess, "Popen", lambda *args, **kwargs: pytest.fail("uncontained bg spawn")) + block = namedtuple("Block", "tool_type content")("bash", "#!bg\nprintf unsafe") + execute = server_authorized_executor(tool_execution.execute_tool_block) + _, result = await execute(block, session_id="chat", owner="alice", workspace=str(path), + security_context=NO_TOOL_SECURITY_CONTEXT) + assert result["containment"]["executed"] is False + assert "bg_job_id" not in result + _, denied = await tool_execution.execute_tool_block( + block, session_id="chat", owner="alice", workspace=str(path), + security_context=NO_TOOL_SECURITY_CONTEXT, request_authority=None, + ) + assert denied["failure_kind"] == "request_authority_denied" + + +def test_detached_supervisor_enforces_timeout(jobs): + path, launched = jobs + record = bg_jobs.launch("sleep 60", "chat", cwd=str(path), max_runtime_s=1) + launched.append(record) + result = finished(record["id"]) + assert result["timed_out"] is True + assert result["teardown"]["dead"] is True + + +def test_restart_keeps_verified_background_supervisor(jobs): + path, launched = jobs + record = bg_jobs.launch("sleep 60", "chat", cwd=str(path)) + launched.append(record) + report = process_reaper.reap_containment_grants() + assert report["background_kept"] == 1 + killed = bg_jobs.kill(record["id"]) + assert killed["killed"] is True + assert killed["teardown"]["dead"] is True + + +def test_kill_never_marks_a_foreign_pid_killed(jobs, monkeypatch): + record = {"id": "stale", "status": "running", "pid": 12345, "start_token": "old", + "session_id": "chat", "started_at": time.time(), "exit_path": "missing"} + bg_jobs._save({"stale": record}) + monkeypatch.setattr(process_ownership, "verify", lambda *args: process_ownership.FOREIGN) + monkeypatch.setattr(bg_jobs, "_kill", lambda *args, **kwargs: pytest.fail("foreign process signalled")) + result = bg_jobs.kill("stale") + assert result["status"] == "running" + assert result.get("killed") is not True + assert result["teardown"]["dead"] is False + + +def test_running_detached_output_and_concurrent_grants_are_preserved(jobs): + path, launched = jobs + for number in range(3): + launched.append(bg_jobs.launch(f"printf job-{number}; sleep 0.3", "chat", cwd=str(path))) + for number, record in enumerate(launched): + assert finished(record["id"])["output"] == f"job-{number}" + grants = containment._load_records() + assert {record["containment_id"] for record in launched} <= grants.keys() + assert all(grants[record["containment_id"]]["release"]["dead"] for record in launched) + + +def test_detached_output_is_available_while_running(jobs): + path, launched = jobs + record = bg_jobs.launch("printf progress; sleep 5", "chat", cwd=str(path)) + launched.append(record) + deadline = time.monotonic() + 3 + while time.monotonic() < deadline: + current = bg_jobs.get(record["id"]) + if "progress" in current["output"]: + assert current["status"] == "running" + return + time.sleep(0.03) + pytest.fail("detached stdout was unavailable until completion") diff --git a/tests/test_bg_job_tools.py b/tests/test_bg_job_tools.py index a21fde88f..d2c035795 100644 --- a/tests/test_bg_job_tools.py +++ b/tests/test_bg_job_tools.py @@ -11,7 +11,7 @@ import time import pytest -from src import bg_jobs +from src import bg_jobs, containment, process_ownership from src.agent_tools.bg_job_tools import ManageBgJobsTool @@ -23,7 +23,11 @@ def store(tmp_path, monkeypatch): monkeypatch.setattr(bg_jobs, "_JOBS_DIR", jobs_dir) monkeypatch.setattr(bg_jobs, "_pid_alive", lambda pid: True) killed: list = [] - monkeypatch.setattr(bg_jobs, "_kill", lambda pid: killed.append(pid)) + monkeypatch.setattr(process_ownership, "verify", lambda *args: process_ownership.OWNED) + def fake_kill(pid, **kwargs): + killed.append(pid) + return containment.ReleaseOutcome(dead=True, escalated=False) + monkeypatch.setattr(bg_jobs, "_kill", fake_kill) return {"dir": jobs_dir, "killed": killed} diff --git a/tests/test_native_execution_containment.py b/tests/test_native_execution_containment.py index 9c86401f0..8303e5927 100644 --- a/tests/test_native_execution_containment.py +++ b/tests/test_native_execution_containment.py @@ -140,3 +140,14 @@ async def test_python_final_expression_and_opt_in_imports(native_boundary): }) assert result["output"] == "42" assert result["teardown"]["dead"] is True + + +async def test_capture_preserves_multibyte_text_across_chunks(native_boundary): + spec = containment.ContainmentSpec( + workspace=str(native_boundary), env=dict(os.environ), wall_clock_s=5, + required=frozenset({containment.PROCESS_TREE, containment.WALL_CLOCK}), max_output_bytes=200000, + ) + result = await containment.run(containment.acquire(spec, owner="unicode"), + [sys.executable, "-c", "import sys; sys.stdout.write('€' * 30000)"], argv=True) + assert result.stdout == "€" * 30000 + assert result.output_truncated is False