mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-06 06:52:20 +02:00
fix(runtime): contain and supervise detached Bash jobs (ODY-145)
This commit is contained in:
+40
-1
@@ -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
|
||||
pass
|
||||
|
||||
@@ -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 ───────────────────────────────────────────
|
||||
|
||||
@@ -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)"
|
||||
|
||||
+115
-70
@@ -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.
|
||||
|
||||
|
||||
+38
-16
@@ -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
|
||||
|
||||
@@ -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)))
|
||||
@@ -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")):
|
||||
|
||||
+12
-9
@@ -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
|
||||
|
||||
@@ -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")
|
||||
@@ -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}
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user