fix(tui-gateway): the hard-exit kill owns every foreground spawn and never blocks
Re-review findings on the foreground-process registry: - Spawn vs exit races. A child spawned but not yet registered when the hard-exit kill ran, or a command launched after the kill took its snapshot, survived under init. The hard-exit kill now raises a one-way exit fence and waits (bounded) for spawns already past it to register; a spawn refuses once the fence is up, and one that registers after a timed-out wait kills itself. - The immediate kill fell back to proc.kill() for every handle. On Modal/Daytona/Vercel that is a blocking SDK cancel (an 8s cancel made _hard_exit take 8s). Popen handles are still killed inline (killpg, never blocks); every other handle's kill runs on a daemon thread under one shared 0.5s deadline. - The registry lock was a plain Lock: a signal landing on a thread that held it deadlocked the exit. It is now an RLock (via a Condition) and the hard-exit path only takes it with a timeout, falling back to a lock-free copy. - The kanban worker's SIGALRM deadman os._exit()ed without the kill; it now goes through it.
This commit is contained in:
@@ -371,6 +371,13 @@ def _install_single_query_signal_handlers(cli):
|
||||
from cli import _arm_exit_watchdog_on_shutdown_signal, _flush_logging_and_stdio, _flush_one_shot_session_store, _interrupt_agent_for_signal
|
||||
import signal as _signal
|
||||
|
||||
def _kill_foreground_and_exit(*_):
|
||||
# The worker's command runs in its own process group: SIGKILL it or it outlives os._exit.
|
||||
with suppress(Exception):
|
||||
from tools.environments.base import kill_live_foreground_processes
|
||||
kill_live_foreground_processes(now=True)
|
||||
os._exit(0)
|
||||
|
||||
def _signal_handler_q(signum, frame):
|
||||
logger.debug("Received signal %s in single-query mode", signum)
|
||||
_arm_exit_watchdog_on_shutdown_signal() # covers wedges in the unwind below
|
||||
@@ -390,7 +397,7 @@ def _install_single_query_signal_handlers(cli):
|
||||
if os.environ.get("HERMES_KANBAN_TASK"):
|
||||
with suppress(Exception):
|
||||
if hasattr(_signal, "SIGALRM"):
|
||||
_signal.signal(_signal.SIGALRM, lambda *_: os._exit(0))
|
||||
_signal.signal(_signal.SIGALRM, _kill_foreground_and_exit)
|
||||
_signal.alarm(5)
|
||||
with suppress(Exception):
|
||||
# Durable flush FIRST: memory-provider shutdown inside _run_cleanup can issue aux-LLM calls,
|
||||
@@ -399,12 +406,8 @@ def _install_single_query_signal_handlers(cli):
|
||||
# store here or the worker's turn (and its usage deltas) never become durable (#88583 /
|
||||
# #50881 class). Best-effort under the SIGALRM deadman above.
|
||||
_flush_one_shot_session_store(cli)
|
||||
# The worker's command runs in its own process group: SIGKILL it or it outlives os._exit.
|
||||
with suppress(Exception):
|
||||
from tools.environments.base import kill_live_foreground_processes
|
||||
kill_live_foreground_processes(now=True)
|
||||
_flush_logging_and_stdio()
|
||||
os._exit(0)
|
||||
_kill_foreground_and_exit()
|
||||
raise KeyboardInterrupt()
|
||||
with suppress(Exception): # restricted environments
|
||||
for _name in ("SIGINT", "SIGTERM", "SIGHUP"):
|
||||
|
||||
@@ -652,6 +652,14 @@ def _isolate_hermes_home(_hermetic_environment):
|
||||
return None
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_foreground_exit_fence():
|
||||
"""A test that drives a hard-exit path raises the one-way foreground-spawn fence; lower it after."""
|
||||
yield
|
||||
if (base := sys.modules.get("tools.environments.base")) is not None:
|
||||
base._exit_fenced = False
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _neutralize_kanban_memory_guard(request, monkeypatch):
|
||||
"""Pin the kanban dispatcher's memory guard to "no data" for every test.
|
||||
|
||||
@@ -242,3 +242,75 @@ def test_exit_cleanup_kills_foreground_command_still_running(monkeypatch):
|
||||
with contextlib.suppress(Exception):
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGKILL)
|
||||
env.cleanup()
|
||||
|
||||
|
||||
# Child for the hard-exit race tests: it really os._exit()s right after the kill, like the watchdogs.
|
||||
_HARD_EXIT_RACE_CHILD = r"""
|
||||
import os, sys, threading
|
||||
from tools.environments import base
|
||||
from tools.environments.local import LocalEnvironment
|
||||
scenario, cmd = sys.argv[1], sys.argv[2]
|
||||
env = LocalEnvironment(cwd=os.getcwd())
|
||||
if scenario == "spawn_before_publish":
|
||||
spawned, release = threading.Event(), threading.Event()
|
||||
real_run_bash = LocalEnvironment._run_bash
|
||||
def gated(self, command, **kw):
|
||||
proc = real_run_bash(self, command, **kw)
|
||||
if cmd in command:
|
||||
spawned.set()
|
||||
release.wait(10) # barrier: the child exists, it is not published yet
|
||||
return proc
|
||||
LocalEnvironment._run_bash = gated
|
||||
threading.Thread(target=env.execute, args=(cmd,), kwargs={"timeout": 600}, daemon=True).start()
|
||||
assert spawned.wait(20)
|
||||
threading.Timer(0.2, release.set).start() # registration lands while the killer is in flight
|
||||
base.kill_live_foreground_processes(now=True)
|
||||
else: # launch_after_fence
|
||||
base.kill_live_foreground_processes(now=True)
|
||||
t = threading.Thread(target=env.execute, args=(cmd,), kwargs={"timeout": 600}, daemon=True)
|
||||
t.start()
|
||||
t.join(3)
|
||||
os._exit(0)
|
||||
"""
|
||||
|
||||
|
||||
@pytest.mark.skipif(os.name == "nt", reason="POSIX process groups")
|
||||
@pytest.mark.live_system_guard_bypass # a red run must reap survivors reparented to init
|
||||
@pytest.mark.parametrize("scenario", ["spawn_before_publish", "launch_after_fence"])
|
||||
def test_hard_exit_leaves_no_foreground_survivor_around_the_spawn(scenario, tmp_path):
|
||||
"""A hard exit must own every foreground child: one spawned but not yet registered when the kill
|
||||
runs, and one launched after the kill took its snapshot (the exit fence refuses it)."""
|
||||
import sys
|
||||
|
||||
import psutil
|
||||
|
||||
cmd = f"sleep {35000 + os.getpid() % 1000}.{len(scenario)}"
|
||||
repo = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
r = subprocess.run([sys.executable, "-c", _HARD_EXIT_RACE_CHILD, scenario, cmd], cwd=str(tmp_path),
|
||||
env={**os.environ, "PYTHONPATH": repo}, capture_output=True, text=True, timeout=120)
|
||||
assert r.returncode == 0, r.stderr[-2000:]
|
||||
time.sleep(0.3)
|
||||
survivors = [p for p in psutil.process_iter(["cmdline"]) if cmd in " ".join(p.info["cmdline"] or [])]
|
||||
for p in survivors:
|
||||
with contextlib.suppress(psutil.Error):
|
||||
p.kill()
|
||||
assert not survivors, f"{scenario}: {[' '.join(p.info['cmdline']) for p in survivors]} outlived the hard exit"
|
||||
|
||||
|
||||
def test_hard_exit_kill_never_blocks_on_a_slow_remote_cancel(monkeypatch):
|
||||
"""Modal/Daytona/Vercel cancel through a blocking SDK call; the hard-exit kill runs just before
|
||||
os._exit, so it must give those one short deadline instead of waiting them out."""
|
||||
from tools.environments import base
|
||||
from tools.environments.base_output import _ThreadedProcessHandle
|
||||
|
||||
class _SdkEnv: # the kill every SDK backend inherits: proc.kill() -> cancel_fn
|
||||
_kill_process = base.BaseEnvironment._kill_process
|
||||
_force_kill_process = base.BaseEnvironment._force_kill_process
|
||||
|
||||
cancelled = threading.Event()
|
||||
handle = _ThreadedProcessHandle(lambda: ("", 0), cancel_fn=lambda: (cancelled.set(), time.sleep(8)))
|
||||
monkeypatch.setitem(base._live_foreground, id(handle), (_SdkEnv(), handle))
|
||||
t0 = time.monotonic()
|
||||
base.kill_live_foreground_processes(now=True)
|
||||
elapsed = time.monotonic() - t0
|
||||
assert cancelled.is_set() and elapsed < 1.0, f"blocked {elapsed:.2f}s"
|
||||
|
||||
@@ -11,6 +11,7 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import shlex
|
||||
import subprocess
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
@@ -56,21 +57,75 @@ _activity_callback_local = threading.local()
|
||||
# own session/process group, so a host that exits mid-command (TUI client gone, SIGTERM) would
|
||||
# orphan the whole tree; the process-exit funnel ``cleanup_all_environments`` kills them.
|
||||
_live_foreground: dict[int, tuple["BaseEnvironment", "ProcessHandle"]] = {}
|
||||
_live_foreground_lock = threading.Lock()
|
||||
# Reentrant, and the hard-exit path only ever takes it with a timeout: a signal handler can run
|
||||
# on a thread that already holds it.
|
||||
_live_foreground_cond = threading.Condition(threading.RLock())
|
||||
_exit_fenced = False # one-way, set by the hard-exit kill: no foreground command spawns after it
|
||||
_spawns_in_flight = 0 # past the fence check, child maybe alive, not yet in _live_foreground
|
||||
_HARD_KILL_BUDGET_S = 0.5
|
||||
|
||||
|
||||
def _enter_foreground_spawn() -> bool:
|
||||
global _spawns_in_flight
|
||||
with _live_foreground_cond:
|
||||
if _exit_fenced:
|
||||
return False
|
||||
_spawns_in_flight += 1
|
||||
return True
|
||||
|
||||
|
||||
def _leave_foreground_spawn(env: "BaseEnvironment", spawned) -> bool:
|
||||
"""Publish ``spawned`` (None: the spawn failed); True when the exit fence went up meanwhile."""
|
||||
global _spawns_in_flight
|
||||
with _live_foreground_cond:
|
||||
_spawns_in_flight -= 1
|
||||
if spawned is not None:
|
||||
_live_foreground[id(spawned)] = (env, spawned)
|
||||
_live_foreground_cond.notify_all()
|
||||
return _exit_fenced
|
||||
|
||||
|
||||
def _quiet_kill(kill: Callable, proc) -> None:
|
||||
try:
|
||||
kill(proc)
|
||||
except Exception:
|
||||
logger.debug("exit-time kill of a foreground command failed", exc_info=True)
|
||||
|
||||
|
||||
def kill_live_foreground_processes(*, now: bool = False) -> int:
|
||||
"""Kill every in-flight foreground command's process tree; returns how many were signalled.
|
||||
|
||||
``now=True`` is for a caller about to ``os._exit``: the graceful kill TERMs, waits and only then
|
||||
KILLs, so a SIGTERM-ignoring command outlives a hard exit that lands inside that window."""
|
||||
with _live_foreground_lock:
|
||||
live = list(_live_foreground.values())
|
||||
for env, proc in live:
|
||||
KILLs, so a SIGTERM-ignoring command outlives a hard exit that lands inside that window. It also
|
||||
raises the exit fence and waits for spawns already past it to register, so no command started
|
||||
around the snapshot survives, and it never blocks past ``_HARD_KILL_BUDGET_S``: SDK cancels
|
||||
(Modal, Daytona, Vercel) run on daemon threads under that one deadline."""
|
||||
global _exit_fenced
|
||||
if not now:
|
||||
with _live_foreground_cond:
|
||||
live = list(_live_foreground.values())
|
||||
for env, proc in live:
|
||||
_quiet_kill(env._kill_process, proc)
|
||||
return len(live)
|
||||
deadline = time.monotonic() + _HARD_KILL_BUDGET_S
|
||||
_exit_fenced = True
|
||||
if _live_foreground_cond.acquire(timeout=_HARD_KILL_BUDGET_S):
|
||||
try:
|
||||
(env._force_kill_process if now else env._kill_process)(proc)
|
||||
except Exception:
|
||||
logger.debug("exit-time kill of a foreground command failed", exc_info=True)
|
||||
_live_foreground_cond.wait_for(lambda: _spawns_in_flight == 0, max(0.0, deadline - time.monotonic()))
|
||||
live = list(_live_foreground.values())
|
||||
finally:
|
||||
_live_foreground_cond.release()
|
||||
else: # the holder is stuck under our signal: a lock-free copy beats hanging the exit
|
||||
live = list(_live_foreground.values())
|
||||
remote = []
|
||||
for env, proc in live:
|
||||
if isinstance(proc, subprocess.Popen): # killpg/kill: never blocks
|
||||
_quiet_kill(env._force_kill_process, proc)
|
||||
else:
|
||||
remote.append(threading.Thread(target=_quiet_kill, args=(env._force_kill_process, proc), daemon=True))
|
||||
remote[-1].start()
|
||||
for t in remote:
|
||||
t.join(max(0.0, deadline - time.monotonic()))
|
||||
return len(live)
|
||||
|
||||
|
||||
@@ -565,17 +620,23 @@ class BaseEnvironment(ABC):
|
||||
def _spawn_and_wait() -> dict:
|
||||
if parent_activity_cb is not None:
|
||||
set_activity_callback(parent_activity_cb)
|
||||
spawned = self._run_bash(wrapped, login=login, timeout=effective_timeout, stdin_data=effective_stdin)
|
||||
if not _enter_foreground_spawn():
|
||||
return {"output": "[host is exiting: command not started]", "returncode": 130}
|
||||
spawned = None
|
||||
try:
|
||||
spawned = self._run_bash(wrapped, login=login, timeout=effective_timeout, stdin_data=effective_stdin)
|
||||
finally:
|
||||
fenced = _leave_foreground_spawn(self, spawned)
|
||||
proc_holder.append(spawned)
|
||||
with _live_foreground_lock:
|
||||
_live_foreground[id(spawned)] = (self, spawned)
|
||||
if fenced: # the hard-exit kill may have stopped waiting for us before we registered
|
||||
self._force_kill_process(spawned)
|
||||
try:
|
||||
return self._wait_for_process(
|
||||
spawned, timeout=effective_timeout, bounded_capture=bounded_capture,
|
||||
watch_interrupt_tid=parent_tid,
|
||||
**({"yield_handler": yield_handler} if yield_handler is not None else {}))
|
||||
finally:
|
||||
with _live_foreground_lock:
|
||||
with _live_foreground_cond:
|
||||
_live_foreground.pop(id(spawned), None)
|
||||
|
||||
def _on_timeout() -> None:
|
||||
|
||||
Reference in New Issue
Block a user