diff --git a/agent/deadline.py b/agent/deadline.py index 19d9d5fed9..cbc01da260 100644 --- a/agent/deadline.py +++ b/agent/deadline.py @@ -16,6 +16,7 @@ from __future__ import annotations import asyncio import contextvars +from contextlib import contextmanager import faulthandler import logging import os @@ -339,13 +340,73 @@ def run_bounded_sync( # --- Whole-tree process termination -------------------------------------------- +@contextmanager +def _process_tree_snapshot(pid: int, *, hard_kill: bool): + """Stop each hard-kill target before discovering its children: a running + parent can fork after psutil builds its PID map and escape the final signal. + Resume anything we stopped if signalling fails. Graceful signals never stop + their recipients, since their handlers must remain able to run. + """ + descendants = [] + stopped = [] + try: + try: + import psutil + root = psutil.Process(pid) + descendants = root.children(recursive=True) + if hard_kill: + pending = [root] + seen = {os.getpid()} + known = {process.pid: process for process in descendants} + stop_deadline = time.monotonic() + 1.0 + while pending: + for process in pending: + if process.pid in seen: + continue + seen.add(process.pid) + if time.monotonic() >= stop_deadline: + raise TimeoutError("process tree did not stop before snapshot deadline") + try: + status = process.status() + if status in (psutil.STATUS_ZOMBIE, psutil.STATUS_DEAD): + continue + if status != psutil.STATUS_STOPPED: + process.suspend() + stopped.append(process) + while process.status() != psutil.STATUS_STOPPED: + if time.monotonic() >= stop_deadline: + raise TimeoutError("process did not stop before snapshot deadline") + time.sleep(0.001) + except psutil.NoSuchProcess: + continue + # Rescan after stopping the discovered generation. A child + # may have forked while that generation was being stopped. + for process in root.children(recursive=True): + known.setdefault(process.pid, process) + descendants = list(known.values()) + pending = [process for process in descendants if process.pid not in seen] + except Exception: + # Preserve the existing best-effort group fallback when discovery or + # stopping is unavailable; never strand a successfully stopped target. + logger.debug("kill_process_tree: snapshot incomplete for pid %s", pid, exc_info=True) + yield descendants + finally: + for process in stopped: + try: + process.resume() + except Exception: + logger.debug("kill_process_tree: target already gone or resume refused", exc_info=True) + + def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: """Terminate ``pid`` and all its descendants, portably; True when anything was signalled. - Windows: ``taskkill /F /T`` (``sig`` ignored). POSIX: descendants are snapshotted via psutil - BEFORE signalling (once the parent dies they reparent and a parent walk finds nothing), then - the process group is signalled when ``pid`` leads one, and every snapshotted descendant - individually — which also reaches ``setsid`` children. ``sig`` defaults to ``SIGKILL``.""" + Windows: ``taskkill /F /T`` (``sig`` ignored). POSIX: snapshot descendants via + psutil; for SIGKILL, stop and rescan the live tree so a concurrent fork cannot + escape a stale snapshot. Signal identity-checked descendants before their + parent, then its group when ``pid`` leads one. Stopping is best-effort with a + bounded wait; unavailable psutil still leaves process-group cleanup. Other + signals do not suspend recipients. ``sig`` defaults to ``SIGKILL``.""" if sys.platform == "win32": try: from hermes_cli._subprocess_compat import windows_hide_flags @@ -367,37 +428,33 @@ def kill_process_tree(pid: int, *, sig: Optional[int] = None) -> bool: if sig is None: sig = _signal.SIGKILL - try: - import psutil - descendants = psutil.Process(int(pid)).children(recursive=True) - except Exception: - # Already gone, or psutil unavailable — the group signal still covers same-session descendants. - descendants = [] - - signalled = False - try: - # getpgid→killpg has an inherent TOCTOU shared by every killpg site; the psutil - # sweep below is identity-aware (PID + create time) and does not. - pgid = os.getpgid(pid) - except (ProcessLookupError, PermissionError, OSError): - pgid = None - try: - if pgid is not None and pgid == pid: - # pid leads its own group (the == check avoids signalling the caller's group). - os.killpg(pgid, sig) # windows-footgun: ok — POSIX-only branch (win32 returns above) - else: - os.kill(pid, sig) - signalled = True - except ProcessLookupError: - pass - except (PermissionError, OSError): - logger.debug("kill_process_tree: signal failed for pid %s", pid, exc_info=True) - - for child in descendants: + with _process_tree_snapshot(int(pid), hard_kill=sig == _signal.SIGKILL) as descendants: + signalled = False + # Signal descendants while their ownership ancestry is still observable. + # Frozen hard-kill targets cannot fork during this bottom-up teardown. + for child in reversed(descendants): + try: + if child.is_running(): + child.send_signal(sig) + signalled = True + except Exception: + continue try: - if child.is_running(): # identity-aware: recycled PIDs skipped - child.send_signal(sig) - signalled = True - except Exception: - continue - return signalled + # getpgid→killpg has an inherent TOCTOU shared by every killpg site; the psutil + # sweep below is identity-aware (PID + create time) and does not. + pgid = os.getpgid(pid) + except (ProcessLookupError, PermissionError, OSError): + pgid = None + try: + if pgid is not None and pgid == pid: + # pid leads its own group (the == check avoids signalling the caller's group). + os.killpg(pgid, sig) # windows-footgun: ok — POSIX-only branch (win32 returns above) + else: + os.kill(pid, sig) + signalled = True + except ProcessLookupError: + pass + except (PermissionError, OSError): + logger.debug("kill_process_tree: signal failed for pid %s", pid, exc_info=True) + + return signalled diff --git a/evals/cron_timeout_fork_race.py b/evals/cron_timeout_fork_race.py new file mode 100644 index 0000000000..2e3c7bd549 --- /dev/null +++ b/evals/cron_timeout_fork_race.py @@ -0,0 +1,93 @@ +"""Linux live cron race probe: real PID map, synchronized fork, owned-only cleanup. + +Run with the repository's Python and a disposable HERMES_HOME. No provider calls. +The scheduling hook retains psutil's real snapshot; it never invents process state. +""" +import ctypes +import json +import os +from pathlib import Path +import sys +import tempfile +import time + + +def main(): + sys.path.insert(0, str(Path(__file__).resolve().parents[1])) + import psutil + + # Confine adoption to this disposable probe, not the test runner or gateway. + if ctypes.CDLL(None, use_errno=True).prctl(36, 1, 0, 0, 0) != 0: + raise OSError(ctypes.get_errno(), "PR_SET_CHILD_SUBREAPER") + with tempfile.TemporaryDirectory(prefix="cron-owned-fork-") as root: + os.environ["HERMES_HOME"] = root + os.environ["HERMES_CRON_SCRIPT_TIMEOUT"] = "2" + from cron.scheduler_script import _run_job_script + + scripts = Path(root, "scripts") + scripts.mkdir() + release, marker = Path(root, "release"), Path(root, "pids.json") + script = scripts / "spawner.py" + script.write_text( + "import subprocess,sys,time,os,json\nfrom pathlib import Path\n" + "def spawn():\n" + " return subprocess.Popen([sys.executable,'-c','import time;time.sleep(15)']," + "start_new_session=True,stdin=subprocess.DEVNULL,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL)\n" + "first=spawn()\n" + f"Path({str(marker)!r}).write_text(json.dumps([os.getpid(),first.pid]))\n" + f"while not Path({str(release)!r}).exists():time.sleep(0.001)\n" + "late=spawn()\n" + f"Path({str(marker)!r}).write_text(json.dumps([os.getpid(),first.pid,late.pid]))\n" + "time.sleep(15)\n", + encoding="utf-8", + ) + original = psutil._ppid_map + schedules = [] + + def scheduling_map(): + mapping = original() + release.touch() + deadline = time.monotonic() + 2 + while time.monotonic() < deadline: + try: + ids = json.loads(marker.read_text(encoding="utf-8")) + if len(ids) == 3 or psutil.Process(ids[0]).status() == psutil.STATUS_STOPPED: + break + except (FileNotFoundError, ValueError, psutil.NoSuchProcess): + pass + time.sleep(0.001) + schedules.append(len(mapping)) + return mapping + + psutil._ppid_map = scheduling_map + try: + start = time.monotonic() + result = _run_job_script(str(script), workdir=root) + elapsed = time.monotonic() - start + finally: + psutil._ppid_map = original + states = {} + owned = psutil.Process().children(recursive=True) + try: + ids = json.loads(marker.read_text(encoding="utf-8")) + for pid in ids: + try: + states[pid] = psutil.Process(pid).status() + except psutil.NoSuchProcess: + states[pid] = "gone" + print(json.dumps({"result": result, "elapsed": elapsed, "states": states, + "real_snapshot_sizes": schedules}), flush=True) + assert result[1].startswith("Script timed out after 2s:"), result + assert len(ids) >= 2, "pre-existing detached child did not start" + assert all(state in ("gone", psutil.STATUS_ZOMBIE) for state in states.values()), states + finally: + for process in owned: + try: + process.kill() + process.wait(timeout=5) + except psutil.NoSuchProcess: + pass + + +if __name__ == "__main__": + main() diff --git a/tests/agent/test_deadline_fork_race.py b/tests/agent/test_deadline_fork_race.py new file mode 100644 index 0000000000..1a21ea0f88 --- /dev/null +++ b/tests/agent/test_deadline_fork_race.py @@ -0,0 +1,66 @@ +"""Real-process deadline regression; the probe owns and reaps its whole subtree.""" +import os +from pathlib import Path +import subprocess +import signal +import sys +import time + +import pytest + + +@pytest.mark.linux_only +def test_cron_timeout_closes_the_descendant_snapshot_fork_window(tmp_path): + probe = Path(__file__).resolve().parents[2] / "evals" / "cron_timeout_fork_race.py" + result = subprocess.run( + [sys.executable, str(probe)], + env={**os.environ, "HOME": str(tmp_path), "HERMES_HOME": str(tmp_path)}, + stdin=subprocess.DEVNULL, capture_output=True, text=True, timeout=35, + ) + assert result.returncode == 0, result.stdout + result.stderr + + +@pytest.mark.linux_only +@pytest.mark.parametrize("initially_stopped", [False, True]) +def test_refused_hard_kill_preserves_the_targets_original_run_state(monkeypatch, initially_stopped): + import psutil + from agent.deadline import kill_process_tree + + proc = subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(30)"], + stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, + start_new_session=True, + ) + process = psutil.Process(proc.pid) + real_kill, real_killpg = os.kill, os.killpg + + def refuse_owned_kill(pid, sig): + assert pid == proc.pid + if sig == signal.SIGKILL: + raise PermissionError("owned probe: hard kill refused") + return real_kill(pid, sig) + + def refuse_owned_group(pgid, sig): + assert pgid == proc.pid + if sig == signal.SIGKILL: + raise PermissionError("owned probe: hard kill refused") + return real_killpg(pgid, sig) + + def wait_for_state(stopped): + deadline = time.monotonic() + 5 + while (process.status() == psutil.STATUS_STOPPED) != stopped and time.monotonic() < deadline: + time.sleep(0.005) + assert (process.status() == psutil.STATUS_STOPPED) == stopped + + try: + if initially_stopped: + process.suspend() + wait_for_state(True) + with monkeypatch.context() as patcher: + patcher.setattr(os, "kill", refuse_owned_kill) + patcher.setattr(os, "killpg", refuse_owned_group) + assert kill_process_tree(proc.pid) is False + wait_for_state(initially_stopped) + finally: + proc.kill() + proc.wait(timeout=5) diff --git a/tests/cron/test_cron_script.py b/tests/cron/test_cron_script.py index 4cc001aade..b4e2c360b3 100644 --- a/tests/cron/test_cron_script.py +++ b/tests/cron/test_cron_script.py @@ -782,7 +782,9 @@ class TestScriptTimeoutTreeKill: ok, out = sched_script._run_job_script( str(scripts_dir / "spawner.py"), workdir=str(cron_env) ) - assert not ok, f"script should have timed out, got {out!r}" + assert not ok and out.startswith("Script timed out after 2s:"), ( + f"expected the timeout path, got success={ok}, output={out!r}" + ) deadline = time.monotonic() + 5 gpid = None diff --git a/website/docs/developer-guide/cron-internals.md b/website/docs/developer-guide/cron-internals.md index 968af066cd..d68f676ed9 100644 --- a/website/docs/developer-guide/cron-internals.md +++ b/website/docs/developer-guide/cron-internals.md @@ -231,6 +231,16 @@ The script timeout defaults to 3600 seconds (1 hour). `_get_script_timeout()` re This timeout bounds the **pre-run script only**, not the agent. Skill-based / LLM-driven jobs run on a separate *inactivity*-based budget (`HERMES_CRON_TIMEOUT`, default 600s of idle time, `0` = unlimited) — they can run for hours as long as they keep calling tools or streaming tokens, and are only killed after the configured idle period with no activity. Scripts are dispatched to a persistent thread pool (not held under the tick lock), so a long-running script does not block other due jobs from firing. +On timeout or ownership cancellation, `cron.scheduler_script` uses the shared +`agent.deadline.kill_process_tree` hard-kill path. On POSIX it briefly stops and +rescans the live tree before signalling descendants and their parent, including +children in separate sessions with no inherited output pipes. This closes the +fork-after-snapshot race. The stop wait is bounded; discovery or permission +failures still use best-effort group cleanup, not a sandbox guarantee. Any target +stopped by cleanup is resumed if termination fails; already-stopped targets keep +their original state. Explicit graceful signals do not suspend their recipients. +Windows continues to use `taskkill /F /T`. + ### Provider Recovery `run_job()` passes the user's configured fallback providers and credential pool into the `AIAgent` instance: