diff --git a/tools/process_registry.py b/tools/process_registry.py index 950329e38e..efbdd77a05 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -1,15 +1,11 @@ -""" -Process Registry -- in-memory registry for background processes spawned via -terminal(background=true): rolling 200KB output buffer, poll/log/wait/kill, -crash recovery via a JSON checkpoint, and session-scoped tracking for gateway -reset protection. - -Background processes execute THROUGH the environment interface -- nothing runs -on the host unless TERMINAL_ENV=local; for Docker/Singularity/Modal/Daytona/SSH -the command runs inside the sandbox. +"""Process Registry -- in-memory registry for background processes spawned via +terminal(background=true): rolling 200KB output buffer, poll/log/wait/kill, JSON +checkpoint for crash recovery, session-scoped tracking for gateway reset protection. +Nothing runs on the host unless TERMINAL_ENV=local; other backends run in their sandbox. """ import codecs +from contextlib import suppress import json import logging import os @@ -23,9 +19,8 @@ import uuid from pathlib import Path _IS_WINDOWS = platform.system() == "Windows" -# systemd transient scopes exist only on Linux. Gate every scope-path branch -# on this constant (not merely "not Windows") so macOS and other POSIX -# platforms provably never touch systemd code (#70716 cross-platform audit). +# systemd transient scopes exist only on Linux; gate every scope-path branch on this +# (not merely "not Windows") so macOS and other POSIX platforms never touch systemd. _IS_LINUX = platform.system() == "Linux" from tools.environments.local import _find_shell, _resolve_safe_cwd, _sanitize_subprocess_env from hermes_cli._subprocess_compat import windows_hide_flags @@ -38,45 +33,35 @@ from agent.redact import redact_sensitive_text logger = logging.getLogger(__name__) - -# Checkpoint file for crash recovery (gateway only) +# Crash-recovery checkpoint (gateway only) CHECKPOINT_PATH = get_hermes_home() / "processes.json" -# Limits -MAX_OUTPUT_CHARS = 200_000 # 200KB rolling output buffer -FINISHED_TTL_SECONDS = 1800 # Keep finished processes for 30 minutes -MAX_PROCESSES = 64 # Max concurrent tracked processes (LRU pruning) +MAX_OUTPUT_CHARS = 200_000 # rolling output buffer +FINISHED_TTL_SECONDS = 1800 # keep finished processes 30 minutes +MAX_PROCESSES = 64 # max tracked processes (LRU pruning) -# Watch pattern rate limiting — PER SESSION. At most ONE watch-match notification -# every WATCH_MIN_INTERVAL_SECONDS; a match inside the cooldown is dropped and counts -# as one strike per window. After WATCH_STRIKE_LIMIT consecutive strike windows the -# session's watch_patterns are permanently disabled and it falls back to -# notify_on_complete semantics (one notification when the process exits). +# Watch-pattern rate limiting, PER SESSION: one watch-match notification per +# WATCH_MIN_INTERVAL_SECONDS; a match inside the cooldown is dropped and counts as one +# strike per window; WATCH_STRIKE_LIMIT consecutive strike windows permanently disable +# watching and fall back to notify_on_complete semantics. WATCH_MIN_INTERVAL_SECONDS = 15 WATCH_STRIKE_LIMIT = 3 - -# Lifetime cap — independent of the strike counter. A pattern recurring at a cadence -# just above the cooldown never trips the strike limit yet forces a full-context agent -# turn every time; watch_patterns is documented as "ONLY for rare one-shot signals", -# so after this many delivered matches we disable it and fall back to notify_on_complete. +# Lifetime cap, independent of strikes: a pattern recurring just above the cooldown never +# strikes yet forces a full-context agent turn each time; watch_patterns is "ONLY for +# rare one-shot signals", so after this many deliveries fall back to notify_on_complete. WATCH_LIFETIME_MAX_HITS = 8 - -# Global circuit breaker across all sessions — secondary safety net so concurrent -# siblings can't collectively flood the user even when each is under its own cap. +# Global circuit breaker across all sessions so concurrent siblings can't collectively +# flood the user even when each is under its own cap. WATCH_GLOBAL_MAX_PER_WINDOW = 15 WATCH_GLOBAL_WINDOW_SECONDS = 10 WATCH_GLOBAL_COOLDOWN_SECONDS = 30 -# --------------------------------------------------------------------------- -# systemd cgroup isolation for gateway-spawned local executors -# --------------------------------------------------------------------------- -# Under a systemd gateway with MemoryMax, local background commands inherit the -# gateway's cgroup, so a memory-heavy executor can get the ENTIRE gateway killed by -# systemd-oomd. Wrapping the spawn in ``systemd-run --user --scope`` gives the worker -# its own transient cgroup. Usability is probed once (the binary can exist while the -# user D-Bus session is absent — system services, containers) and cached. - +# --- systemd cgroup isolation for gateway-spawned local executors ------------------ +# Under a systemd gateway with MemoryMax, local background commands inherit the gateway's +# cgroup, so a memory-heavy executor can get the ENTIRE gateway killed by systemd-oomd; +# ``systemd-run --user --scope`` gives the worker its own transient cgroup. Usability is +# probed once (binary present but user D-Bus absent in system services/containers). _SYSTEMD_SCOPE_AVAILABLE: Optional[bool] = None _SYSTEMD_SCOPE_PROBE_LOCK = threading.Lock() _SYSTEMD_SCOPE_PROBED_AT = 0.0 @@ -88,11 +73,9 @@ _WORKER_MEMORY_MAX_CAP_BYTES = 4 * 1024 * 1024 * 1024 def _worker_memory_max_bytes() -> int: """Finite per-worker cgroup limit that can never widen host risk. - ``TERMINAL_LOCAL_MEMORY_MAX_MB`` is honored only when it *tightens* the safe bound (min of the gateway's cgroup-v2 ``memory.max`` and half of physical RAM, - capped at 4 GiB), so an oversized override cannot exceed the enclosing slice. - """ + capped at 4 GiB), so an oversized override cannot exceed the enclosing slice.""" override_bound: Optional[int] = None override = os.getenv("TERMINAL_LOCAL_MEMORY_MAX_MB", "").strip() if override: @@ -106,90 +89,59 @@ def _worker_memory_max_bytes() -> int: logger.warning( "Ignoring invalid TERMINAL_LOCAL_MEMORY_MAX_MB=%r; " "expected an integer representing at least %d MiB", - override, - _MIN_WORKER_MEMORY_MAX_BYTES // (1024 * 1024), - ) - + override, _MIN_WORKER_MEMORY_MAX_BYTES // (1024 * 1024)) candidates: List[int] = [] - try: - for line in Path("/proc/self/cgroup").read_text(encoding="utf-8").splitlines(): - if line.startswith("0::"): - relative = line.partition("::")[2].lstrip("/") - raw_limit = ( - Path("/sys/fs/cgroup") / relative / "memory.max" - ).read_text(encoding="utf-8").strip() - if raw_limit.isdigit() and int(raw_limit) >= _MIN_WORKER_MEMORY_MAX_BYTES: - candidates.append(int(raw_limit)) - break - except (OSError, ValueError): - pass - - try: + with suppress(OSError, ValueError): + lines = Path("/proc/self/cgroup").read_text(encoding="utf-8").splitlines() + v2 = next((ln for ln in lines if ln.startswith("0::")), None) + if v2 is not None: + relative = v2.partition("::")[2].lstrip("/") + raw_limit = (Path("/sys/fs/cgroup") / relative / "memory.max").read_text(encoding="utf-8").strip() + if raw_limit.isdigit() and int(raw_limit) >= _MIN_WORKER_MEMORY_MAX_BYTES: + candidates.append(int(raw_limit)) + with suppress(OSError, ValueError, TypeError): physical_bytes = int(os.sysconf("SC_PHYS_PAGES")) * int(os.sysconf("SC_PAGE_SIZE")) - candidates.append(min( - _WORKER_MEMORY_MAX_CAP_BYTES, - max(_MIN_WORKER_MEMORY_MAX_BYTES, physical_bytes // 2), - )) - except (OSError, ValueError, TypeError): - pass - + candidates.append(min(_WORKER_MEMORY_MAX_CAP_BYTES, max(_MIN_WORKER_MEMORY_MAX_BYTES, physical_bytes // 2))) safe_bound = min(candidates) if candidates else _DEFAULT_WORKER_MEMORY_MAX_BYTES return min(override_bound, safe_bound) if override_bound else safe_bound def _systemd_scope_argv(binary: str, unit_name: str, *argv: str) -> List[str]: - """``systemd-run --user --scope`` command line shared by the probe and real spawns. - - ``--collect`` makes the transient scope self-clean after exit; ``--unit`` gives it - a recognisable name for ``systemctl --user status`` / journalctl. - """ + """``systemd-run --user --scope`` argv shared by the probe and real spawns. + ``--collect`` self-cleans the scope after exit; ``--unit`` names it for systemctl.""" return [ - binary, "--user", "--scope", "--quiet", - "--unit", unit_name, - "--collect", + binary, "--user", "--scope", "--quiet", "--unit", unit_name, "--collect", "--property", "MemoryAccounting=yes", "--property", f"MemoryMax={_worker_memory_max_bytes()}", "--property", "OOMPolicy=kill", - "--", - *argv, + "--", *argv, ] def _systemd_scope_cached() -> Optional[bool]: - """Cached probe verdict, or None when a (re)probe is due. - - A True verdict is permanent; a False one expires after - ``_SYSTEMD_SCOPE_FAILURE_TTL_SECONDS`` so a transient D-Bus outage isn't sticky. - """ - cached = _SYSTEMD_SCOPE_AVAILABLE - if cached is True: + """Cached probe verdict, or None when a (re)probe is due. True is permanent; False + expires after ``_SYSTEMD_SCOPE_FAILURE_TTL_SECONDS`` so a D-Bus blip isn't sticky.""" + if _SYSTEMD_SCOPE_AVAILABLE is True: return True - if cached is False and time.monotonic() - _SYSTEMD_SCOPE_PROBED_AT < _SYSTEMD_SCOPE_FAILURE_TTL_SECONDS: - return False - return None + stale = time.monotonic() - _SYSTEMD_SCOPE_PROBED_AT >= _SYSTEMD_SCOPE_FAILURE_TTL_SECONDS + return None if _SYSTEMD_SCOPE_AVAILABLE is None or stale else False def _systemd_run_user_scope_available() -> bool: - """Return True if ``systemd-run --user --scope`` can create a cgroup. - - ``shutil.which`` alone is insufficient: system-service deployments and containers - may lack the user D-Bus session bus even though the binary is on PATH, so every - spawn would fail with ``Failed to connect to user bus``. We run a cheap no-op - probe (``systemd-run --user --scope --unit=… -- /bin/true``) and cache the outcome. - """ + """True if ``systemd-run --user --scope`` can create a cgroup. + ``shutil.which`` alone is insufficient: system services and containers may lack + the user D-Bus bus even with the binary on PATH (every spawn would fail with + ``Failed to connect to user bus``), so a cheap ``/bin/true`` probe is run and cached.""" global _SYSTEMD_SCOPE_AVAILABLE, _SYSTEMD_SCOPE_PROBED_AT verdict = _systemd_scope_cached() if verdict is not None: return verdict - - # Double-checked locking keeps concurrent first-use spawns from observing a - # temporary False while the definitive probe is still in flight — such a race - # would launch the losing workload back inside the gateway cgroup. + # Double-checked locking: a concurrent first-use spawn must not observe a temporary + # False mid-probe, or it would launch back inside the gateway cgroup. with _SYSTEMD_SCOPE_PROBE_LOCK: verdict = _systemd_scope_cached() if verdict is not None: return verdict - available = False if _IS_LINUX: try: @@ -197,23 +149,19 @@ def _systemd_run_user_scope_available() -> bool: binary = shutil.which("systemd-run") if binary: - # A unique unit avoids collisions; the timeout bounds D-Bus. + # Unique unit avoids collisions; the timeout bounds D-Bus. probe_unit = f"hermes-probe-scope-{os.getpid()}-{uuid.uuid4().hex[:8]}" result = subprocess.run( - _systemd_scope_argv(binary, probe_unit, "/bin/true"), - capture_output=True, - timeout=3, + _systemd_scope_argv(binary, probe_unit, "/bin/true"), capture_output=True, timeout=3, ) available = result.returncode == 0 if not available: logger.debug( "systemd-run --user --scope probe failed (rc=%s): %s", - result.returncode, - (result.stderr or b"").decode("utf-8", "replace").strip(), + result.returncode, (result.stderr or b"").decode("utf-8", "replace").strip(), ) except Exception as exc: logger.debug("systemd-run --user --scope probe error: %s", exc) - _SYSTEMD_SCOPE_AVAILABLE = available _SYSTEMD_SCOPE_PROBED_AT = time.monotonic() return available @@ -221,79 +169,51 @@ def _systemd_run_user_scope_available() -> bool: def _is_supervised_gateway_process() -> bool: """Whether this process is the live, supervised Hermes gateway itself. - - Supervisor markers and ``_HERMES_GATEWAY`` are inherited by every descendant - (and importing ``gateway.run`` sets the latter), so also require ownership of - the live gateway PID file — transient scopes are for the gateway, not terminal - children or unrelated CLIs in the same supervised tree. - """ + Supervisor markers and ``_HERMES_GATEWAY`` are inherited by every descendant (and + importing ``gateway.run`` sets the latter), so also require ownership of the live + gateway PID file — scopes are for the gateway, not terminal children or CLIs.""" if os.environ.get("_HERMES_GATEWAY") != "1": return False - try: from gateway.restart import is_gateway_supervisor_process from gateway.status import get_running_pid - return ( - is_gateway_supervisor_process() - and get_running_pid(cleanup_stale=False) == os.getpid() - ) + return is_gateway_supervisor_process() and get_running_pid(cleanup_stale=False) == os.getpid() except Exception as exc: logger.debug("Could not verify supervised gateway process identity: %s", exc) return False -def _build_systemd_scope_argv( - shell_argv: List[str], - unit_suffix: str, -) -> List[str]: +def _build_systemd_scope_argv(shell_argv: List[str], unit_suffix: str) -> List[str]: """Wrap *shell_argv* in a ``systemd-run --user --scope`` invocation with its own memory accounting, so an OOM in the worker cannot kill the gateway cgroup.""" import shutil binary = shutil.which("systemd-run") if binary is None: - # Caller should have checked _systemd_run_user_scope_available(); - # guard anyway so we never pass None into Popen. + # Caller should have probed availability; never pass None into Popen anyway. return shell_argv return _systemd_scope_argv(binary, f"hermes-worker-{unit_suffix}", *shell_argv) def _stop_systemd_unit(unit_name: str) -> bool: """Stop a transient systemd user scope by unit name. - - Reaps the *entire* cgroup — catching double-forked descendants that survive a - plain PID signal because they were reparented to init inside the scope. - ``systemctl --user stop`` SIGTERMs every process in the cgroup and escalates to - SIGKILL after ``TimeoutStopSec``. - - Returns True if the unit was stopped (or was already gone), False if - ``systemctl`` is unavailable or the stop command failed. - """ + Reaps the *entire* cgroup — catching double-forked descendants reparented to init + inside the scope that survive a plain PID signal (SIGTERM all, SIGKILL after + ``TimeoutStopSec``). True if stopped or already gone; False if ``systemctl`` is + unavailable or the stop failed.""" import shutil binary = shutil.which("systemctl") if binary is None: return False try: - result = subprocess.run( - [binary, "--user", "stop", unit_name], - capture_output=True, - timeout=15, - ) + result = subprocess.run([binary, "--user", "stop", unit_name], capture_output=True, timeout=15) if result.returncode != 0: stderr = (result.stderr or b"").decode(errors="replace").strip() - stderr_lower = stderr.lower() - if any( - marker in stderr_lower - for marker in ("not loaded", "not found", "does not exist") - ): + if any(marker in stderr.lower() for marker in ("not loaded", "not found", "does not exist")): return True - logger.debug( - "systemctl --user stop %s exited %d: %s", - unit_name, result.returncode, - stderr, - ) + logger.debug("systemctl --user stop %s exited %d: %s", unit_name, result.returncode, stderr) return False return True except Exception as exc: @@ -312,32 +232,43 @@ def format_uptime_short(seconds: int) -> str: return f"{hours}h {mins}m" +def _not_found(session_id: str) -> dict: + return {"status": "not_found", "error": f"No process with ID {session_id}"} + + +def _output_tail(session: "ProcessSession", n: int) -> str: + """Last *n* chars of the session output with ANSI sequences stripped.""" + from tools.ansi_strip import strip_ansi + + return strip_ansi(session.output_buffer[-n:]) + + @dataclass class ProcessSession: """A tracked background process with output buffering.""" - id: str # Unique session ID ("proc_xxxxxxxxxxxx") - command: str # Original command string - task_id: str = "" # Task/sandbox isolation key - owner_task_id: str = "" # RAW spawning task id (e.g. "sa-..."); task_id is the - # CONTAINER key (may be collapsed by _resolve_container_task_id) - # so ownership checks must use this field - session_key: str = "" # Gateway session key (for reset protection) - pid: Optional[int] = None # OS process ID + id: str # "proc_xxxxxxxxxxxx" + command: str + task_id: str = "" # Task/sandbox isolation key (CONTAINER key, + # may be collapsed by _resolve_container_task_id) + owner_task_id: str = "" # RAW spawning task id ("sa-..."); ownership + # checks must use this, not task_id + session_key: str = "" # Gateway session key (reset protection) + pid: Optional[int] = None process: Optional[subprocess.Popen] = None # Popen handle (local only) - env_ref: Any = None # Reference to the environment object - cwd: Optional[str] = None # Working directory - started_at: float = 0.0 # time.time() of spawn (wall clock) + env_ref: Any = None # Environment object (sandbox spawns) + cwd: Optional[str] = None + started_at: float = 0.0 # time.time() of spawn host_start_time: Optional[int] = None # kernel start ticks (/proc//stat f22) — PID-reuse guard - exited: bool = False # Whether the process has finished - exit_code: Optional[int] = None # Exit code (None if still running) + exited: bool = False + exit_code: Optional[int] = None # None while running completion_reason: str = "exited" # exited|killed|lost|failed_start|already_exited termination_source: str = "" # process.kill|kill_all|backend_lost|failed_start - output_buffer: str = "" # Rolling output (last MAX_OUTPUT_CHARS) + output_buffer: str = "" # Rolling tail (last max_output_chars) max_output_chars: int = MAX_OUTPUT_CHARS - detached: bool = False # True if recovered from crash (no pipe) + detached: bool = False # Recovered from checkpoint (no pipe) pid_scope: str = "host" # "host" for local/PTY PIDs, "sandbox" for env-local PIDs - systemd_unit: str = "" # transient scope unit name when spawned under systemd-run (#70716) - # Watcher/notification metadata (persisted for crash recovery) + systemd_unit: str = "" # transient scope unit name when spawned under systemd-run + # Watcher/notification routing (persisted for crash recovery) watcher_platform: str = "" watcher_chat_id: str = "" watcher_user_id: str = "" @@ -345,25 +276,22 @@ class ProcessSession: watcher_thread_id: str = "" watcher_message_id: str = "" # Triggering message id — reply anchor for topic routing watcher_interval: int = 0 # 0 = no watcher configured - # Session-db id of the spawning conversation; lets the gateway drop completions - # whose session was closed at a user boundary (/new) instead of injecting them - # into the chat's NEW session. + # Session-db id of the spawning conversation; lets the gateway drop completions whose + # session was closed at a user boundary (/new) instead of injecting into the NEW one. parent_session_id: str = "" - notify_on_complete: bool = False # Queue agent notification on exit + notify_on_complete: bool = False # Queue agent notification on exit watch_patterns: List[str] = field(default_factory=list) _watch_hits: int = field(default=0, repr=False) # total matches delivered _watch_suppressed: int = field(default=0, repr=False) # matches dropped by rate limit _watch_disabled: bool = field(default=False, repr=False) # permanently killed after strike limit - # Per-session rate-limit state (see WATCH_* constants). A strike is a WINDOW with - # drops, not a dropped match. - _watch_last_emit_at: float = field(default=0.0, repr=False) + # Rate-limit window state (see WATCH_*). A strike is a WINDOW with drops, not a drop. _watch_cooldown_until: float = field(default=0.0, repr=False) _watch_strike_candidate: bool = field(default=False, repr=False) _watch_consecutive_strikes: int = field(default=0, repr=False) _completion_event: threading.Event = field(default_factory=threading.Event, repr=False) _lock: threading.Lock = field(default_factory=threading.Lock) _reader_thread: Optional[threading.Thread] = field(default=None, repr=False) - _pty: Any = field(default=None, repr=False) # ptyprocess handle (when use_pty=True) + _pty: Any = field(default=None, repr=False) # ptyprocess handle (use_pty=True) def append_output(self, text: str) -> None: """Append to the rolling output buffer under the session lock, keeping the tail.""" @@ -372,16 +300,26 @@ class ProcessSession: if len(self.output_buffer) > self.max_output_chars: self.output_buffer = self.output_buffer[-self.max_output_chars:] + def mark_exited(self, exit_code, reason: str = "exited", source: str = "") -> None: + """Record an exit. A kill that raced the observer already recorded its own + exit_code/reason; never overwrite it.""" + self.exited = True + if self.completion_reason != "killed": + self.exit_code = exit_code + self.completion_reason = reason + if source: + self.termination_source = source + +# Watcher routing fields, in event-dict key order (``watcher_`` on the session). +_WATCHER_ROUTE_KEYS = ("platform", "chat_id", "user_id", "user_name", "thread_id", "message_id") # Session fields persisted verbatim in the crash-recovery checkpoint (plus # ``session_id``; ``command`` is redacted and ``owner_task_id`` defaulted on write). _CHECKPOINT_FIELDS = ( "command", "pid", "pid_scope", "host_start_time", "systemd_unit", "cwd", "started_at", "task_id", "owner_task_id", "session_key", - "watcher_platform", "watcher_chat_id", "watcher_user_id", "watcher_user_name", - "watcher_thread_id", "watcher_message_id", "watcher_interval", - "parent_session_id", "notify_on_complete", "watch_patterns", -) + *(f"watcher_{k}" for k in _WATCHER_ROUTE_KEYS), "watcher_interval", + "parent_session_id", "notify_on_complete", "watch_patterns") _CHECKPOINT_DEFAULTS = { f.name: ([] if f.name == "watch_patterns" else f.default) for f in ProcessSession.__dataclass_fields__.values() @@ -391,30 +329,21 @@ _CHECKPOINT_DEFAULTS = { class ProcessRegistry: """In-memory registry of running and finished background processes. - Thread-safe: accessed from executor threads (terminal_tool, process handlers), - the gateway asyncio loop (watchers, reset checks) and the cleanup thread. - """ + the gateway asyncio loop (watchers, reset checks) and the cleanup thread.""" _SHELL_NOISE_SUBSTRINGS = ( - "bash: cannot set terminal process group", - "bash: no job control in this shell", - "no job control in this shell", - "cannot set terminal process group", - "tcsetattr: Inappropriate ioctl for device", - ) + "no job control in this shell", "cannot set terminal process group", + "tcsetattr: Inappropriate ioctl for device") def __init__(self): self._running: Dict[str, ProcessSession] = {} self._finished: Dict[str, ProcessSession] = {} self._lock = threading.Lock() - # Side-channel for check_interval watchers (gateway reads after agent run) self.pending_watchers: List[Dict[str, Any]] = [] - - # Unified queue for all background events (completion, watch_match, - # async_delegation...; distinguished by "type"). CLI process_loop and the - # gateway drain it after each agent turn to auto-trigger new turns. + # Unified queue for all background events (distinguished by "type"); the CLI + # process_loop and the gateway drain it after each agent turn to trigger new turns. import queue as _queue_mod self.completion_queue: _queue_mod.Queue = _queue_mod.Queue() # Rehydrate durable delegation completions once, at registry startup. @@ -423,23 +352,19 @@ class ProcessRegistry: restore_undelivered_completions(self.completion_queue) except Exception as exc: logger.warning("Could not restore async delegation completions: %s", exc) - - # Sessions whose completion the agent already consumed via wait()/read_log() - # — it has the output in hand, so drain loops AND gateway/tui watchers skip. + # Completions the agent already consumed via wait()/read_log() (output in + # hand): drain loops AND gateway/tui watchers skip them. self._completion_consumed: set = set() - # Sessions merely *observed* exited via poll(). poll() is read-only and must - # NOT mark consumed (a status check would suppress the watcher's autonomous - # delivery turn), but on the CLI the poll result is inline in the same turn, - # so drain_notifications() skips these to avoid a duplicate [SYSTEM: ...]; + # Sessions merely *observed* exited via poll(). poll() is read-only and must NOT + # mark consumed (a status check would suppress the watcher's autonomous delivery + # turn), but the CLI has the poll result inline in the same turn, so + # drain_notifications() skips these to avoid a duplicate [SYSTEM: ...]; # gateway/tui watchers deliberately ignore this set. self._poll_observed: set = set() - # Global watch-match circuit breaker across all sessions. self._global_watch_lock = threading.Lock() - self._global_watch_window_start: float = 0.0 - self._global_watch_window_hits: int = 0 - self._global_watch_tripped_until: float = 0.0 - self._global_watch_suppressed_during_trip: int = 0 + self._global_watch_window_start = self._global_watch_tripped_until = 0.0 + self._global_watch_window_hits = self._global_watch_suppressed_during_trip = 0 # Driver-installed sinks (desktop gateway): on_output(session, chunk) streams # live output from reader threads; on_close(session_or_none, process_id) drops # a read-only terminal tab without killing the process. @@ -455,137 +380,97 @@ class ProcessRegistry: return "\n".join(lines) def _emit_output(self, session: ProcessSession, chunk: str) -> None: - """Forward a freshly-read chunk to the live-output sink, if one is set. - Called from reader threads; never raise into the read loop.""" + """Forward a chunk to the live-output sink; called from reader threads, never raises.""" sink = self.on_output if sink is None or not chunk: return - try: + with suppress(Exception): sink(session, chunk) - except Exception: - pass def _check_watch_patterns(self, session: ProcessSession, new_text: str) -> None: """Scan a freshly-read chunk for watch patterns and queue notifications. - - Rate limiting per session (see the WATCH_* constants): one match per cooldown - window, a match inside the window is one strike, WATCH_STRIKE_LIMIT - consecutive strikes or WATCH_LIFETIME_MAX_HITS total deliveries disable - watching and promote the session to notify_on_complete. - """ + Per-session rate limiting (see WATCH_* constants): one match per cooldown + window, a match inside the window is one strike, WATCH_STRIKE_LIMIT consecutive + strikes or WATCH_LIFETIME_MAX_HITS total deliveries disable watching and + promote the session to notify_on_complete.""" if not session.watch_patterns or session._watch_disabled: return # Late chunks after the reader declared exit are post-exit noise; dropping them # avoids stale notifications minutes after the process ended. if session.exited: return - - # Scan new text line-by-line for pattern matches - matched_lines = [] - matched_pattern = None - for line in new_text.splitlines(): - for pat in session.watch_patterns: - if pat in line: - matched_lines.append(line.rstrip()) - if matched_pattern is None: - matched_pattern = pat - break # one match per line is enough - - if not matched_lines: + hits = [ # (first matching pattern, line) — one match per line + (next(p for p in session.watch_patterns if p in line), line.rstrip()) + for line in new_text.splitlines() if any(p in line for p in session.watch_patterns)] + if not hits: return - + matched_pattern = hits[0][0] + matched_lines = [line for _, line in hits] now = time.time() - should_disable = False - lifetime_exhausted = False with session._lock: - # Case 1: inside the cooldown — drop, count one strike per window, and - # disable + promote once the strike limit is hit. if session._watch_cooldown_until and now < session._watch_cooldown_until: + # Inside the cooldown: drop, count one strike per window, disable + + # promote once the strike limit is hit. session._watch_suppressed += len(matched_lines) - if not session._watch_strike_candidate: - # First drop in this window — count one strike. - session._watch_strike_candidate = True - session._watch_consecutive_strikes += 1 - if session._watch_consecutive_strikes >= WATCH_STRIKE_LIMIT: - session._watch_disabled = True - # Promote to notify_on_complete so the agent still gets - # exactly one notification when the process actually ends. - session.notify_on_complete = True - should_disable = True - return_early = True - else: - # Case 2: cooldown expired. A prior window with no drops resets the - # consecutive-strike counter (healthy cadence again). - if session._watch_cooldown_until and not session._watch_strike_candidate: - session._watch_consecutive_strikes = 0 - session._watch_strike_candidate = False - - # Emit the notification and start a new cooldown window. - session._watch_last_emit_at = now - session._watch_cooldown_until = now + WATCH_MIN_INTERVAL_SECONDS - session._watch_hits += 1 - suppressed = session._watch_suppressed - session._watch_suppressed = 0 - return_early = False - # Lifetime cap: this match is still delivered, but no further ones. - lifetime_exhausted = session._watch_hits >= WATCH_LIFETIME_MAX_HITS - if lifetime_exhausted: - session._watch_disabled = True - session.notify_on_complete = True - - if return_early: - if should_disable: - # Exactly one summary so the agent/user sees why things went quiet. - self.completion_queue.put({ - **self._watch_event_base(session), - "type": "watch_disabled", - "suppressed": session._watch_suppressed, - "message": ( - f"Watch patterns disabled for process {session.id} — " - f"{WATCH_STRIKE_LIMIT} consecutive rate-limit windows triggered " - f"(min spacing {WATCH_MIN_INTERVAL_SECONDS}s). " - f"Falling back to notify_on_complete semantics; you'll get " - f"exactly one notification when the process exits." - ), - }) - return - - # Trim matched output to a reasonable size + if session._watch_strike_candidate: + return + session._watch_strike_candidate = True + session._watch_consecutive_strikes += 1 + if session._watch_consecutive_strikes < WATCH_STRIKE_LIMIT: + return + session._watch_disabled = True + # Promote so the agent still gets exactly one notification on exit, + # plus exactly one summary so it sees why things went quiet. + session.notify_on_complete = True + self._emit_watch_disabled( + session, session._watch_suppressed, + f"{WATCH_STRIKE_LIMIT} consecutive rate-limit windows triggered " + f"(min spacing {WATCH_MIN_INTERVAL_SECONDS}s). ") + return + # Cooldown expired. A prior window with no drops resets the + # consecutive-strike counter (healthy cadence again). + if session._watch_cooldown_until and not session._watch_strike_candidate: + session._watch_consecutive_strikes = 0 + session._watch_strike_candidate = False + # Emit and start a new cooldown window. + session._watch_cooldown_until = now + WATCH_MIN_INTERVAL_SECONDS + session._watch_hits += 1 + suppressed = session._watch_suppressed + session._watch_suppressed = 0 + # Lifetime cap: this match is still delivered, but no further ones. + lifetime_exhausted = session._watch_hits >= WATCH_LIFETIME_MAX_HITS + if lifetime_exhausted: + session._watch_disabled = True + session.notify_on_complete = True output = "\n".join(matched_lines[:20]) if len(output) > 2000: output = output[:2000] + "\n...(truncated)" - - if not self._global_watch_admit(now): - # Even when the breaker drops the final match, still explain the silence. - if lifetime_exhausted: - self._emit_lifetime_watch_disabled(session) - return - - notification = { - **self._watch_event_base(session), - "type": "watch_match", - "pattern": matched_pattern, - "output": output, - "suppressed": suppressed, - } - _redact_process_result(notification) - self.completion_queue.put(notification) - + if self._global_watch_admit(now): + notification = { + **self._watch_event_base(session), + "type": "watch_match", + "pattern": matched_pattern, + "output": output, + "suppressed": suppressed, + } + _redact_process_result(notification) + self.completion_queue.put(notification) + # Even when the breaker drops the final match, still explain the silence. if lifetime_exhausted: - self._emit_lifetime_watch_disabled(session) + self._emit_watch_disabled( + session, 0, f"reached the lifetime cap of {WATCH_LIFETIME_MAX_HITS} delivered matches. ", + ) - def _emit_lifetime_watch_disabled(self, session: ProcessSession) -> None: - """Queue the watch_disabled summary for the lifetime-cap path.""" + def _emit_watch_disabled(self, session: ProcessSession, suppressed: int, why: str) -> None: + """Queue the one-shot watch_disabled summary (strike-limit or lifetime-cap path).""" self.completion_queue.put({ **self._watch_event_base(session), "type": "watch_disabled", - "suppressed": 0, + "suppressed": suppressed, "message": ( - f"Watch patterns disabled for process {session.id} — " - f"reached the lifetime cap of {WATCH_LIFETIME_MAX_HITS} delivered " - f"matches. Falling back to notify_on_complete semantics; you'll get " - f"exactly one notification when the process exits." - ), + f"Watch patterns disabled for process {session.id} — {why}" + f"Falling back to notify_on_complete semantics; you'll get " + f"exactly one notification when the process exits."), }) @staticmethod @@ -597,89 +482,58 @@ class ProcessRegistry: "task_id": session.task_id, "owner_task_id": session.owner_task_id or session.task_id, "command": session.command, - "platform": session.watcher_platform, - "chat_id": session.watcher_chat_id, - "user_id": session.watcher_user_id, - "user_name": session.watcher_user_name, - "thread_id": session.watcher_thread_id, - "message_id": session.watcher_message_id, + **{key: getattr(session, f"watcher_{key}") for key in _WATCHER_ROUTE_KEYS}, } @staticmethod def _global_watch_event(type_: str, message: str, **extra) -> dict: """Unaddressed (all-sessions) watch breaker event.""" return { - "session_id": "", - "session_key": "", - "command": "", - "type": type_, - **extra, + "session_id": "", "session_key": "", "command": "", "type": type_, **extra, "message": message, - "platform": "", - "chat_id": "", - "user_id": "", - "user_name": "", - "thread_id": "", + "platform": "", "chat_id": "", "user_id": "", "user_name": "", "thread_id": "", } def _global_watch_admit(self, now: float) -> bool: """True if this watch_match may pass the global breaker. - In cooldown: drop and count. Otherwise slide the rolling window; exceeding the cap trips the breaker for WATCH_GLOBAL_COOLDOWN_SECONDS with ONE - "tripped" summary, and the cooldown's end emits ONE "released" summary. - """ - release_msg = None + "tripped" summary, and the cooldown's end emits ONE "released" summary.""" + events = [] # summary events, queued outside the lock with self._global_watch_lock: # Handle cooldown expiry first so we can emit the release summary. if self._global_watch_tripped_until and now >= self._global_watch_tripped_until: suppressed = self._global_watch_suppressed_during_trip self._global_watch_tripped_until = 0.0 self._global_watch_suppressed_during_trip = 0 - self._global_watch_window_start = now - self._global_watch_window_hits = 0 + self._global_watch_window_start, self._global_watch_window_hits = now, 0 if suppressed > 0: - # Queued outside the lock (below). - release_msg = self._global_watch_event( + events.append(self._global_watch_event( "watch_overflow_released", f"Watch-pattern notifications resumed. " f"{suppressed} match event(s) were suppressed during the flood.", - suppressed=suppressed, - ) - - # Still in cooldown — drop and count. + suppressed=suppressed)) if self._global_watch_tripped_until and now < self._global_watch_tripped_until: + # Still in cooldown — drop and count. self._global_watch_suppressed_during_trip += 1 admit = False - trip_now = None else: - # Slide the window. if now - self._global_watch_window_start >= WATCH_GLOBAL_WINDOW_SECONDS: - self._global_watch_window_start = now - self._global_watch_window_hits = 0 - - if self._global_watch_window_hits >= WATCH_GLOBAL_MAX_PER_WINDOW: - # Trip the breaker. + self._global_watch_window_start, self._global_watch_window_hits = now, 0 + admit = self._global_watch_window_hits < WATCH_GLOBAL_MAX_PER_WINDOW + if admit: + self._global_watch_window_hits += 1 + else: self._global_watch_tripped_until = now + WATCH_GLOBAL_COOLDOWN_SECONDS self._global_watch_suppressed_during_trip += 1 - trip_now = now - admit = False - else: - self._global_watch_window_hits += 1 - trip_now = None - admit = True - - # Queue summary events outside the lock. - if release_msg is not None: - self.completion_queue.put(release_msg) - if trip_now is not None: - self.completion_queue.put(self._global_watch_event( - "watch_overflow_tripped", - f"Watch-pattern overflow: >{WATCH_GLOBAL_MAX_PER_WINDOW} " - f"notifications in {WATCH_GLOBAL_WINDOW_SECONDS}s across all processes. " - f"Suppressing further watch_match events for " - f"{WATCH_GLOBAL_COOLDOWN_SECONDS}s.", - )) + events.append(self._global_watch_event( + "watch_overflow_tripped", + f"Watch-pattern overflow: >{WATCH_GLOBAL_MAX_PER_WINDOW} " + f"notifications in {WATCH_GLOBAL_WINDOW_SECONDS}s across all processes. " + f"Suppressing further watch_match events for " + f"{WATCH_GLOBAL_COOLDOWN_SECONDS}s.")) + for msg in events: + self.completion_queue.put(msg) return admit @staticmethod @@ -687,141 +541,108 @@ class ProcessRegistry: """Best-effort liveness check for host-visible PIDs.""" if not pid: return False - # ``os.kill(pid, 0)`` is NOT a no-op on Windows (bpo-14484) — use - # the cross-platform existence check. + # ``os.kill(pid, 0)`` is NOT a no-op on Windows (bpo-14484) — use the + # cross-platform existence check. from gateway.status import _pid_exists return _pid_exists(pid) @staticmethod def _safe_host_start_time(pid: Optional[int]) -> Optional[int]: """Kernel start ticks for a host PID, or None when unavailable.""" - if not pid: - return None try: from gateway.status import get_process_start_time - return get_process_start_time(pid) + return get_process_start_time(pid) if pid else None except Exception: return None @classmethod def _host_pid_is_ours(cls, pid: Optional[int], expected_start: Optional[int]) -> bool: """True only if ``pid`` is alive AND still the process we spawned. - The kernel recycles PIDs, so a stored number can later name an unrelated process (seen in the wild: a browser's session leader tree-killed). The kernel start time captured at spawn must match the live one; with no baseline - (legacy checkpoints, no ``/proc``) degrade to a bare liveness check. - """ - if not cls._is_host_pid_alive(pid): - return False - if expected_start is None: - return True - return cls._safe_host_start_time(pid) == expected_start + (legacy checkpoints, no ``/proc``) degrade to a bare liveness check.""" + return cls._is_host_pid_alive(pid) and ( + expected_start is None or cls._safe_host_start_time(pid) == expected_start) def _refresh_detached_session(self, session: Optional[ProcessSession]) -> Optional[ProcessSession]: """Update recovered host-PID sessions when the underlying process has exited.""" if session is None or session.exited or not session.detached or session.pid_scope != "host": return session - # A recycled PID (alive but not ours) counts as "our process exited" so a # later kill() can never tree-kill the stranger. if self._host_pid_is_ours(session.pid, session.host_start_time): return session - with session._lock: if session.exited: return session - session.exited = True - # Recovered sessions no longer have a waitable handle, so the real - # exit code is unavailable once the original process object is gone. - session.exit_code = None - + # No waitable handle survives recovery, so the real exit code is unknown. + session.exited, session.exit_code = True, None self._move_to_finished(session) return session @staticmethod def _proc_alive(proc) -> bool: - """True if a psutil.Process is running and not a zombie. - - A zombie is already dead (just unreaped), so there's nothing to SIGKILL. - """ + """True if a psutil.Process is running and not a zombie (already dead, just unreaped).""" try: import psutil - if not proc.is_running(): - return False - return proc.status() != psutil.STATUS_ZOMBIE + return proc.is_running() and proc.status() != psutil.STATUS_ZOMBIE except Exception: return False @staticmethod def _config_value(section: str, key: str, fallback): """``config.yaml`` value for ``section.key``, else the DEFAULT_CONFIG value. - Raises if config is unreadable; callers wrap with their own hard fallback so - registry code paths never crash on a broken config file. - """ + registry code paths never crash on a broken config file.""" from hermes_cli.config import DEFAULT_CONFIG, cfg_get, read_raw_config val = cfg_get(read_raw_config(), section, key) return DEFAULT_CONFIG[section][key] if val is None else val @staticmethod - def _daemon_term_grace_seconds() -> float: - """Grace window (s) between SIGTERM and escalated SIGKILL, floored at 0 - (0 disables escalation). ``terminal.daemon_term_grace_seconds``; 2.0 if - config is unreadable.""" + def _config_seconds(key: str, fallback: float) -> float: + """``terminal.`` as a non-negative float (0 disables); *fallback* if unreadable.""" try: - return max(float(ProcessRegistry._config_value("terminal", "daemon_term_grace_seconds", 2.0)), 0.0) + return max(float(ProcessRegistry._config_value("terminal", key, fallback)), 0.0) except Exception: - return 2.0 + return fallback + + @staticmethod + def _daemon_term_grace_seconds() -> float: + """Grace (s) between SIGTERM and escalated SIGKILL; 0 disables escalation.""" + return ProcessRegistry._config_seconds("daemon_term_grace_seconds", 2.0) @classmethod def _terminate_host_pid(cls, pid: int, expected_start: Optional[int] = None) -> None: """Terminate a host-visible PID and its descendants. - - ``expected_start`` (kernel start time at spawn) is re-validated first; a - mismatch or dead PID means the number was recycled onto a stranger and we - refuse to touch it — a leaked orphan beats tree-killing someone's browser. - - POSIX: psutil walks the tree and SIGTERMs children before the parent so - subprocess trees (Chromium renderers under an agent-browser daemon) aren't - reparented to init and survive. After ``terminal.daemon_term_grace_seconds`` - any survivor is SIGKILLed (0 disables escalation). - - Windows: ``taskkill /PID /T /F`` (same primitive as - ``gateway.status.terminate_pid``; ``/F`` is already a hard kill). The psutil - path is unusable there: PPID links go stale so ``children(recursive=True)`` - misses orphans, and ``terminate()`` is ``TerminateProcess()`` on one handle — - nothing cascades like a SIGTERM to a process group. The bare ``os.kill`` - fallback covers OSError/PermissionError and a missing ``taskkill.exe``. - """ + ``expected_start`` (kernel start time at spawn) is re-validated first: a mismatch + or dead PID means the number was recycled onto a stranger and we refuse to touch + it — a leaked orphan beats tree-killing someone's browser. POSIX: psutil SIGTERMs + children before the parent (so trees aren't reparented to init and survive), then + SIGKILLs survivors after ``terminal.daemon_term_grace_seconds``. Windows: + ``taskkill /T /F`` (psutil's stale PPID links miss orphans there); ``os.kill`` + is the fallback.""" if expected_start is not None and not cls._host_pid_is_ours(pid, expected_start): logger.warning( "Refusing to terminate host pid %d: start-time mismatch — " - "PID was recycled onto an unrelated process.", pid, - ) + "PID was recycled onto an unrelated process.", pid) return - def _sigterm_quietly(): - try: - os.kill(pid, signal.SIGTERM) - except (OSError, ProcessLookupError, PermissionError): - pass + def _sigterm_quietly(): + with suppress(OSError, ProcessLookupError, PermissionError): + os.kill(pid, signal.SIGTERM) if _IS_WINDOWS: try: subprocess.run( - ["taskkill", "/PID", str(pid), "/T", "/F"], - capture_output=True, - text=True, encoding='utf-8', errors='replace', - timeout=10, - creationflags=windows_hide_flags(), - stdin=subprocess.DEVNULL, - ) + ["taskkill", "/PID", str(pid), "/T", "/F"], capture_output=True, text=True, + encoding='utf-8', errors='replace', timeout=10, creationflags=windows_hide_flags(), + stdin=subprocess.DEVNULL) except (FileNotFoundError, subprocess.TimeoutExpired, OSError): _sigterm_quietly() return - import psutil + gone = (psutil.NoSuchProcess, psutil.AccessDenied, OSError) try: parent = psutil.Process(pid) except psutil.NoSuchProcess: @@ -829,47 +650,41 @@ class ProcessRegistry: except (OSError, PermissionError): _sigterm_quietly() return - # Snapshot the whole tree (children before parent) and SIGTERM each. try: targets = parent.children(recursive=True) - except (psutil.NoSuchProcess, psutil.AccessDenied, OSError): + except gone: targets = [] targets.append(parent) - for proc in targets: - try: + with suppress(gone): proc.terminate() - except (psutil.NoSuchProcess, psutil.AccessDenied, OSError): - pass - - # Escalate to SIGKILL for anything that ignored SIGTERM within the grace - # window. We deliberately do NOT trust ``psutil.wait_procs``' gone/alive - # partition: it reaps via ``Process.wait()`` and mis-partitions across - # zombie transitions in a parent/child tree, leaving survivors un-killed. - # A direct liveness re-probe of every target is deterministic. + # Escalate to SIGKILL for anything that ignored SIGTERM within the grace window. + # ``psutil.wait_procs``' gone/alive partition is deliberately NOT trusted: it + # reaps via ``Process.wait()`` and mis-partitions across zombie transitions in a + # parent/child tree, leaving survivors un-killed. Re-probing every target is + # deterministic. grace = cls._daemon_term_grace_seconds() if grace <= 0: return deadline = time.monotonic() + grace - while time.monotonic() < deadline: - if not any(cls._proc_alive(_p) for _p in targets): - break + while time.monotonic() < deadline and any(cls._proc_alive(_p) for _p in targets): time.sleep(0.05) for proc in targets: - try: - if not cls._proc_alive(proc): - continue - proc.kill() # SIGKILL on POSIX - logger.info( - "Escalated to SIGKILL for pid %d (ignored SIGTERM within " - "%.1fs grace)", proc.pid, grace, - ) - except (psutil.NoSuchProcess, psutil.AccessDenied, OSError): - pass + with suppress(gone): + if cls._proc_alive(proc): + proc.kill() # SIGKILL on POSIX + logger.info("Escalated to SIGKILL for pid %d (ignored SIGTERM within %.1fs grace)", proc.pid, grace) # ----- Spawn ----- + @staticmethod + def _new_session(command, task_id, owner_task_id, session_key, cwd, **extra) -> ProcessSession: + return ProcessSession( + id=f"proc_{uuid.uuid4().hex[:12]}", command=command, task_id=task_id, + owner_task_id=owner_task_id or task_id, session_key=session_key, cwd=cwd, + started_at=time.time(), **extra) + @staticmethod def _env_temp_dir(env: Any) -> str: """Return the writable sandbox temp dir for env-backed background tasks.""" @@ -883,31 +698,35 @@ class ProcessRegistry: logger.debug("Could not resolve environment temp dir: %s", exc) return "/tmp" - def _scope_argv(self, session: ProcessSession, argv: List[str], unit_suffix: str, label: str): - """Wrap *argv* in a transient systemd scope when we are the supervised gateway. - - Returns ``(argv, scoped)``. A scoped worker gets its own cgroup so an OOM kills - only the worker, not the gateway (and its messaging control plane). - """ + def _scope_argv(self, session: ProcessSession, safe_command: str, unit_suffix: str, label: str) -> List[str]: + """Login-shell argv for *safe_command* (parity with LocalEnvironment: rc files + sourced, user tools on PATH), wrapped in a transient systemd scope when we are + the supervised gateway (own cgroup: an OOM kills only the worker, not the + gateway and its messaging control plane).""" + argv = [_find_shell(), "-lic", f"set +m; {safe_command}"] in_supervised_gateway = _IS_LINUX and _is_supervised_gateway_process() if in_supervised_gateway and _systemd_run_user_scope_available(): session.systemd_unit = f"hermes-worker-{unit_suffix}.scope" - return _build_systemd_scope_argv(argv, unit_suffix=unit_suffix), True + return _build_systemd_scope_argv(argv, unit_suffix=unit_suffix) if in_supervised_gateway: - # Under a supervisor but no private cgroup: an OOM in the worker can - # still take the whole gateway down. + # Under a supervisor but no private cgroup: a worker OOM can still take + # the whole gateway down. logger.debug( "%s background executor not isolated in a systemd scope " - "(systemd-run --user unavailable); worker shares the gateway cgroup.", - label, - ) - return argv, False + "(systemd-run --user unavailable); worker shares the gateway cgroup.", label) + return argv - def _track_started(self, session: ProcessSession, reader_target, reader_name: str) -> None: + @staticmethod + def _spawn_env(env_vars: dict) -> dict: + """Sanitized child env; PYTHONUNBUFFERED so tqdm/datasets-style buffering + doesn't hide progress from process(action="poll").""" + env = _sanitize_subprocess_env(os.environ, env_vars) + env["PYTHONUNBUFFERED"] = "1" + return env + + def _track_started(self, session: ProcessSession, reader_target, reader_name: str, extra_args=()) -> None: """Start the output reader thread, register the session and checkpoint it.""" - reader = threading.Thread( - target=reader_target, args=(session,), daemon=True, name=reader_name, - ) + reader = threading.Thread(target=reader_target, args=(session, *extra_args), daemon=True, name=reader_name) session._reader_thread = reader reader.start() with self._lock: @@ -917,30 +736,19 @@ class ProcessRegistry: def _spawn_local_pty(self, session: ProcessSession, safe_command: str, env_vars: dict) -> ProcessSession: """PTY spawn for interactive CLI tools (Codex, Claude Code, REPLs). - Raises ImportError when no PTY backend is installed and re-raises any spawn - failure; ``spawn_local`` falls back to pipe mode in both cases. - """ + failure; ``spawn_local`` falls back to pipe mode in both cases.""" if _IS_WINDOWS: from winpty import PtyProcess as _PtyProcessCls else: from ptyprocess import PtyProcess as _PtyProcessCls - user_shell = _find_shell() - pty_env = _sanitize_subprocess_env(os.environ, env_vars) - pty_env["PYTHONUNBUFFERED"] = "1" + pty_env = self._spawn_env(env_vars) # A PTY is a real TTY, so pager-happy tools (git log/diff, man) WILL page and # hang waiting for `q` — default them to cat, honoring any pager the user set. pty_env.setdefault("GIT_PAGER", "cat") pty_env.setdefault("PAGER", "cat") - pty_argv, _ = self._scope_argv( - session, [user_shell, "-lic", f"set +m; {safe_command}"], session.id, "PTY", - ) - pty_proc = _PtyProcessCls.spawn( - pty_argv, - cwd=session.cwd, - env=pty_env, - dimensions=(30, 120), - ) + pty_argv = self._scope_argv(session, safe_command, session.id, "PTY") + pty_proc = _PtyProcessCls.spawn(pty_argv, cwd=session.cwd, env=pty_env, dimensions=(30, 120)) session.pid = pty_proc.pid session.host_start_time = self._safe_host_start_time(session.pid) session._pty = pty_proc @@ -948,38 +756,18 @@ class ProcessRegistry: return session def spawn_local( - self, - command: str, - cwd: str = None, - task_id: str = "", - session_key: str = "", - env_vars: dict = None, - use_pty: bool = False, - owner_task_id: str = "", - ) -> ProcessSession: - """Spawn a background process locally (TERMINAL_ENV=local only; other - backends use spawn_via_env()). - - ``use_pty`` requests a pseudo-terminal via ptyprocess/pywinpty for interactive - CLI tools; it falls back to a plain pipe when that is unavailable or fails. - """ + self, command: str, cwd: str = None, task_id: str = "", session_key: str = "", + env_vars: dict = None, use_pty: bool = False, owner_task_id: str = "") -> ProcessSession: + """Spawn a background process locally (TERMINAL_ENV=local; other backends use + spawn_via_env()). ``use_pty`` requests a pseudo-terminal via ptyprocess/pywinpty + for interactive CLIs, falling back to a plain pipe when unavailable or failing.""" # Bash parses ``A && B &`` as ``(A && B) &`` — a subshell that holds our stdout # pipe open forever when B is a long-running server. The rewriter turns it into # ``A && { B & }``. Lazy import: terminal_tool imports this module. from tools.terminal_tool import _rewrite_compound_background as _rewrite_bg safe_command = _rewrite_bg(command) - - session = ProcessSession( - id=f"proc_{uuid.uuid4().hex[:12]}", - command=command, - task_id=task_id, - owner_task_id=owner_task_id or task_id, - session_key=session_key, - cwd=_resolve_safe_cwd(cwd or os.getcwd()), - started_at=time.time(), - ) - + session = self._new_session(command, task_id, owner_task_id, session_key, _resolve_safe_cwd(cwd or os.getcwd())) pty_scope_attempted = False if use_pty: try: @@ -996,189 +784,99 @@ class ProcessRegistry: "to avoid duplicate command execution" ) from e session.systemd_unit = "" - - # Pipe path (non-PTY or PTY fallback). The user's login shell keeps parity with - # LocalEnvironment (rc files sourced, user tools on PATH). PYTHONUNBUFFERED so - # tqdm/datasets-style buffering doesn't hide progress from process(action="poll"). - user_shell = _find_shell() - bg_env = _sanitize_subprocess_env(os.environ, env_vars) - bg_env["PYTHONUNBUFFERED"] = "1" + # Pipe path (non-PTY or PTY fallback). _popen_kwargs = {"creationflags": windows_hide_flags()} if _IS_WINDOWS else {} - unit_suffix = f"{session.id}-pipe-fallback" if pty_scope_attempted else session.id - spawn_argv, _ = self._scope_argv( - session, [user_shell, "-lic", f"set +m; {safe_command}"], unit_suffix, "Local", - ) - + spawn_argv = self._scope_argv(session, safe_command, unit_suffix, "Local") # start_new_session is REQUIRED with systemd-run --scope too: the scope does not # give the worker a new session, so from an interactive TUI the worker would # share the foreground process group and background spawns would stop the whole # session (observed as dead TUIs in state T). Cgroup isolation is unaffected — # the scope attaches to the invoked process, not the spawning session. proc = subprocess.Popen( - spawn_argv, - text=True, - cwd=session.cwd, - env=bg_env, - encoding="utf-8", - errors="replace", - stdout=subprocess.PIPE, - stderr=subprocess.STDOUT, - stdin=subprocess.DEVNULL, - start_new_session=True, - **_popen_kwargs, - ) - + spawn_argv, text=True, cwd=session.cwd, env=self._spawn_env(env_vars), encoding="utf-8", + errors="replace", stdout=subprocess.PIPE, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, + start_new_session=True, **_popen_kwargs) session.process = proc session.pid = proc.pid session.host_start_time = self._safe_host_start_time(session.pid) - try: self._track_started(session, self._reader_loop, f"proc-reader-{session.id}") except Exception: - # Post-Popen setup failed — kill the orphaned subprocess (and any setsid - # descendants) before re-raising so nothing leaks untracked. - try: - if session.systemd_unit: - # Scope teardown is the authoritative cleanup for the worker cgroup - # (never killpg here); the wrapper PID is terminated as fallback. - _stop_systemd_unit(session.systemd_unit) - self._terminate_host_pid(proc.pid, session.host_start_time) - elif not _IS_WINDOWS: - try: - kill_signal = getattr(signal, "SIGKILL", signal.SIGTERM) - os.killpg(os.getpgid(proc.pid), kill_signal) # windows-footgun: ok - guarded by _IS_WINDOWS above - except (ProcessLookupError, PermissionError, OSError): - proc.kill() - else: - proc.kill() - except Exception: - pass - try: - proc.wait(timeout=5) - except Exception: - pass + self._reap_untracked(session, proc) raise - return session + def _reap_untracked(self, session: ProcessSession, proc: subprocess.Popen) -> None: + """Post-Popen setup failed: kill the orphaned subprocess (and any setsid + descendants) so nothing leaks untracked.""" + with suppress(Exception): + if session.systemd_unit: + # Scope teardown is the authoritative cleanup for the worker cgroup + # (never killpg here); the wrapper PID is terminated as fallback. + _stop_systemd_unit(session.systemd_unit) + self._terminate_host_pid(proc.pid, session.host_start_time) + elif not _IS_WINDOWS: + try: + kill_signal = getattr(signal, "SIGKILL", signal.SIGTERM) + os.killpg(os.getpgid(proc.pid), kill_signal) # windows-footgun: ok - guarded by _IS_WINDOWS above + except (ProcessLookupError, PermissionError, OSError): + proc.kill() + else: + proc.kill() + with suppress(Exception): + proc.wait(timeout=5) + def spawn_via_env( - self, - env: Any, - command: str, - cwd: str = None, - task_id: str = "", - session_key: str = "", - timeout: int = 10, - owner_task_id: str = "", - ) -> ProcessSession: + self, env: Any, command: str, cwd: str = None, task_id: str = "", session_key: str = "", + timeout: int = 10, owner_task_id: str = "") -> ProcessSession: """Spawn a background process inside a non-local backend's sandbox. - - The command is wrapped to capture its in-sandbox PID and redirect output to - a log file, which later execute() calls poll. Less capable than local spawn - (no live pipe, no stdin) but runs in the correct sandbox context. - """ - session = ProcessSession( - id=f"proc_{uuid.uuid4().hex[:12]}", - command=command, - task_id=task_id, - owner_task_id=owner_task_id or task_id, - session_key=session_key, - cwd=cwd, - started_at=time.time(), - env_ref=env, - pid_scope="sandbox", - ) - - # Run the command in the sandbox with output capture + The command is wrapped to capture its in-sandbox PID and redirect output to a + log file that later execute() calls poll. No live pipe or stdin, but it runs in + the correct sandbox context.""" + session = self._new_session(command, task_id, owner_task_id, session_key, cwd, env_ref=env, pid_scope="sandbox") temp_dir = self._env_temp_dir(env) - log_path = f"{temp_dir}/hermes_bg_{session.id}.log" - pid_path = f"{temp_dir}/hermes_bg_{session.id}.pid" - exit_path = f"{temp_dir}/hermes_bg_{session.id}.exit" - quoted_command = shlex.quote(command) - quoted_temp_dir = shlex.quote(temp_dir) - quoted_log_path = shlex.quote(log_path) - quoted_pid_path = shlex.quote(pid_path) - quoted_exit_path = shlex.quote(exit_path) + log_path, pid_path, exit_path = (f"{temp_dir}/hermes_bg_{session.id}.{ext}" for ext in ("log", "pid", "exit")) + q = shlex.quote bg_command = ( - f"mkdir -p {quoted_temp_dir} && " - f"( nohup bash -lc {quoted_command} > {quoted_log_path} 2>&1; " - f"rc=$?; printf '%s\\n' \"$rc\" > {quoted_exit_path} ) & " - f"echo $! > {quoted_pid_path} && cat {quoted_pid_path}" - ) - + f"mkdir -p {q(temp_dir)} && " + f"( nohup bash -lc {q(command)} > {q(log_path)} 2>&1; " + f"rc=$?; printf '%s\\n' \"$rc\" > {q(exit_path)} ) & " + f"echo $! > {q(pid_path)} && cat {q(pid_path)}") try: - result = env.execute( - bg_command, - timeout=timeout, - rewrite_compound_background=False, - ) + result = env.execute(bg_command, timeout=timeout, rewrite_compound_background=False) output = result.get("output", "").strip() - # Try to extract the PID from the output - for line in output.splitlines(): - line = line.strip() - if line.isdigit(): - session.pid = int(line) - break + session.pid = next((int(ln) for ln in map(str.strip, output.splitlines()) if ln.isdigit()), None) # No PID from the wrapper (syntax error, broken redirect): a failed launch, # not a fake running session. if session.pid is None: - session.exited = True - session.exit_code = int(result.get("returncode", -1)) - if session.exit_code == 0: - session.exit_code = -1 - session.completion_reason = "failed_start" - session.termination_source = "failed_start" - session.output_buffer = result.get("output", "").strip() + session.mark_exited(int(result.get("returncode", -1)) or -1, "failed_start", "failed_start") + session.output_buffer = output except Exception as e: - session.exited = True - session.exit_code = -1 - session.completion_reason = "failed_start" - session.termination_source = "failed_start" + session.mark_exited(-1, "failed_start", "failed_start") session.output_buffer = f"Failed to start: {e}" - - if not session.exited: - # Start a poller thread that periodically reads the log file - reader = threading.Thread( - target=self._env_poller_loop, - args=(session, env, log_path, pid_path, exit_path), - daemon=True, - name=f"proc-poller-{session.id}", - ) - session._reader_thread = reader - reader.start() - - with self._lock: - self._prune_if_needed() - if not session.exited: - self._running[session.id] = session - - if not session.exited: - self._write_checkpoint() - + if session.exited: + with self._lock: + self._prune_if_needed() + else: + self._track_started( + session, self._env_poller_loop, f"proc-poller-{session.id}", (env, log_path, pid_path, exit_path)) return session # ----- Reader / Poller Threads ----- def _reader_loop(self, session: ProcessSession): """Background thread: read stdout from a local Popen process. - - Uses ``buffer.read1(4096)`` not ``TextIOWrapper.read(4096)``: on pipes the - latter blocks until EOF, landing "live" output in one burst at exit. - - Orphaned-pipe guard: a backgrounded grandchild (``node server.js &``) - inherits our pipe's write end, so EOF never arrives while it lives — a - blocking read would park this thread, ``session.exited`` would never flip - and ``notify_on_complete`` never fire (``_reconcile_local_exit`` only runs - lazily from poll()/wait()). On POSIX we ``select()`` with a short interval - and stop draining shortly after the direct child exits, mirroring - ``tools/environments/base.py::_wait_for_process``. Windows pipes lack - select(); the blocking path stays and the lazy reconcile is the safety net. - """ + ``buffer.read1(4096)`` not ``TextIOWrapper.read(4096)``: on pipes the latter + blocks until EOF, landing "live" output in one burst at exit. Orphaned-pipe + guard: a backgrounded grandchild (``node server.js &``) inherits our pipe's write + end so EOF never arrives while it lives, which would park this thread and never + fire ``notify_on_complete``; on POSIX we ``select()`` and stop draining shortly + after the direct child exits (mirrors ``environments/base.py::_wait_for_process``). + Windows pipes lack select(), so the lazy ``_reconcile_local_exit`` is the net.""" first_chunk = True - # A multibyte UTF-8 char split across read1() chunks would become U+FFFD - # mojibake with stateless decoding; the incremental decoder holds the partial - # sequence until the continuation bytes arrive. + # A split multibyte UTF-8 char would become U+FFFD with stateless decoding; the + # incremental decoder holds the partial sequence until the rest arrives. decoder = codecs.getincrementaldecoder("utf-8")(errors="replace") def _append_chunk(chunk: str): @@ -1187,101 +885,80 @@ class ProcessRegistry: chunk = self._clean_shell_noise(chunk) first_chunk = False self._ingest_output(session, chunk) - try: proc = session.process if proc is None or proc.stdout is None: return stdout = proc.stdout - raw_read = getattr(getattr(stdout, "buffer", None), "read1", None) + def _read_once(): + """One 4 KiB read: decoded text ('' for a partial multibyte tail), None at EOF.""" + if raw_read is None: # mocked/alternate streams without a raw buffer: less "live" + return stdout.read(4096) or None + raw = raw_read(4096) + return decoder.decode(raw) if raw else None # select() needs a real OS fd; mocked streams (tests, adapters) may lack - # fileno() and use the blocking loop instead. - fd = None - if raw_read is not None and not _IS_WINDOWS: - fileno = getattr(stdout, "fileno", None) - try: - candidate = fileno() if callable(fileno) else None - except Exception: - candidate = None - if isinstance(candidate, int) and candidate >= 0: - fd = candidate - + # fileno() and use the blocking read instead. + try: + fd = stdout.fileno() if raw_read is not None and not _IS_WINDOWS else None + except Exception: + fd = None + if not (isinstance(fd, int) and fd >= 0): + fd = None if fd is not None: import select as _select - - idle_after_exit = 0 - while True: + idle_after_exit = 0 + while True: + if fd is not None: try: ready, _, _ = _select.select([fd], [], [], 0.2) except (ValueError, OSError): break # fd already closed - if ready: - raw = raw_read(4096) - if not raw: - break # true EOF — all writers closed - chunk = decoder.decode(raw) - if chunk: - _append_chunk(chunk) - idle_after_exit = 0 - elif proc.poll() is not None: - # Direct child gone and pipe idle ~200ms: a few more cycles - # for a buffered tail, then stop rather than wait forever on - # an orphaned grandchild's pipe. - idle_after_exit += 1 + if not ready: + # Direct child gone and pipe idle ~200ms: a few more cycles for a + # buffered tail, then stop rather than wait forever on an orphaned + # grandchild's pipe. + if proc.poll() is not None: + idle_after_exit += 1 if idle_after_exit >= 3: break - else: - while True: - if raw_read is not None: - raw = raw_read(4096) - if not raw: - break - chunk = decoder.decode(raw) - if not chunk: - continue # partial multibyte sequence — wait for more bytes - else: - # Mocked/alternate streams without a raw buffer: less "live". - chunk = stdout.read(4096) - if not chunk: - break - + continue + chunk = _read_once() + if chunk is None: + break # true EOF — all writers closed + if chunk: _append_chunk(chunk) + idle_after_exit = 0 except Exception as e: logger.debug("Process stdout reader ended: %s", e) finally: - # Flush the decoder: a truncated multibyte sequence at EOF becomes one - # U+FFFD instead of vanishing. - try: - tail = decoder.decode(b"", final=True) - if tail: - _append_chunk(tail) - except Exception: - pass - # Always reap the child to prevent zombie processes. - try: - session.process.wait(timeout=5) - except Exception as e: - logger.debug("Process wait timed out or failed: %s", e) - self._finish_exited(session, session.process.returncode) + self._finish_reader( + session, decoder, _append_chunk, "Process", + lambda: session.process.wait(timeout=5), lambda: session.process.returncode) - def _env_poller_loop( - self, session: ProcessSession, env: Any, log_path: str, pid_path: str, exit_path: str - ): + def _finish_reader(self, session, decoder, append, label, wait, exit_code) -> None: + """Reader-thread teardown: flush the decoder (a truncated multibyte tail becomes + one U+FFFD instead of vanishing), reap the child (no zombies), record the exit.""" + with suppress(Exception): + tail = decoder.decode(b"", final=True) + if tail: + append(tail) + try: + wait() + except Exception as e: + logger.debug("%s wait timed out or failed: %s", label, e) + self._finish_exited(session, exit_code()) + + def _env_poller_loop(self, session: ProcessSession, env: Any, log_path: str, pid_path: str, exit_path: str): """Background thread: poll a sandbox log file for non-local backends.""" - quoted_log_path = shlex.quote(log_path) - quoted_pid_path = shlex.quote(pid_path) - quoted_exit_path = shlex.quote(exit_path) - prev_output_len = 0 # track delta for watch pattern scanning + q = shlex.quote + prev_output_len = 0 # delta tracking for watch-pattern scanning while not session.exited: - time.sleep(2) # Poll every 2 seconds + time.sleep(2) try: - # Read new output from the log file - result = env.execute(f"cat {quoted_log_path} 2>/dev/null", timeout=10) - new_output = result.get("output", "") + new_output = env.execute(f"cat {q(log_path)} 2>/dev/null", timeout=10).get("output", "") if new_output: - # Delta since the previous read feeds watch-pattern scanning. delta = new_output[prev_output_len:] if len(new_output) > prev_output_len else "" prev_output_len = len(new_output) with session._lock: @@ -1289,36 +966,23 @@ class ProcessRegistry: if delta: self._check_watch_patterns(session, delta) self._emit_output(session, delta) - - # Check if process is still running check = env.execute( - f"kill -0 \"$(cat {quoted_pid_path} 2>/dev/null)\" 2>/dev/null; echo $?", - timeout=5, - ) + f"kill -0 \"$(cat {q(pid_path)} 2>/dev/null)\" 2>/dev/null; echo $?", timeout=5) check_output = check.get("output", "").strip() if check_output and check_output.splitlines()[-1].strip() != "0": - # Process has exited -- get exit code captured by the wrapper shell. - exit_result = env.execute( - f"cat {quoted_exit_path} 2>/dev/null", - timeout=5, - ) - exit_str = exit_result.get("output", "").strip() + # Exited -- read the exit code captured by the wrapper shell. + exit_str = env.execute(f"cat {q(exit_path)} 2>/dev/null", timeout=5).get("output", "").strip() try: - session.exit_code = int(exit_str.splitlines()[-1].strip()) + exit_code = int(exit_str.splitlines()[-1].strip()) except (ValueError, IndexError): - session.exit_code = -1 - session.exited = True - if session.completion_reason != "killed": - session.completion_reason = "exited" - self._move_to_finished(session) + exit_code = -1 + session.exit_code = exit_code # unlike mark_exited, a raced kill still takes this code + self._finish_exited(session, exit_code) return - except Exception: # Environment might be gone (sandbox reaped, etc.) - session.exited = True - session.exit_code = -1 - session.completion_reason = "lost" - session.termination_source = "backend_lost" + session.exited, session.exit_code = True, -1 + session.completion_reason, session.termination_source = "lost", "backend_lost" self._move_to_finished(session) return @@ -1327,9 +991,6 @@ class ProcessRegistry: pty = session._pty # Same split-multibyte handling as _reader_loop. decoder = codecs.getincrementaldecoder("utf-8")(errors="replace") - - _append_text = lambda text: self._ingest_output(session, text) # noqa: E731 - try: while pty.isalive(): try: @@ -1338,25 +999,14 @@ class ProcessRegistry: # ptyprocess returns bytes; pywinpty returns str text = chunk if isinstance(chunk, str) else decoder.decode(chunk) if text: - _append_text(text) + self._ingest_output(session, text) except Exception: # EOFError included break except Exception as e: logger.debug("PTY stdout reader ended: %s", e) - - # Flush any partial multibyte sequence held by the decoder. - try: - tail = decoder.decode(b"", final=True) - if tail: - _append_text(tail) - except Exception: - pass - - try: - pty.wait() - except Exception as e: - logger.debug("PTY wait timed out or failed: %s", e) - self._finish_exited(session, pty.exitstatus if hasattr(pty, 'exitstatus') else -1) + self._finish_reader( + session, decoder, lambda t: self._ingest_output(session, t), "PTY", + pty.wait, lambda: pty.exitstatus if hasattr(pty, 'exitstatus') else -1) def _ingest_output(self, session: ProcessSession, text: str) -> None: """Buffer a freshly-read chunk, then scan watch patterns and stream it live.""" @@ -1365,32 +1015,20 @@ class ProcessRegistry: self._emit_output(session, text) def _finish_exited(self, session: ProcessSession, exit_code) -> None: - """Mark a reader-observed exit and move the session to finished. - - A kill that raced the reader already recorded its own exit_code/reason; - don't overwrite it. - """ - session.exited = True - if session.completion_reason != "killed": - session.exit_code = exit_code - session.completion_reason = "exited" + """Mark a reader-observed exit (a raced kill keeps its own code/reason) and finish.""" + session.mark_exited(exit_code) self._move_to_finished(session) def _move_to_finished(self, session: ProcessSession): """Move a session from running to finished. - Idempotent: kill_process() and the reader thread can both call this; only - the FIRST move enqueues the completion notification, so no duplicates. - """ + the FIRST move enqueues the completion notification, so no duplicates.""" with self._lock: was_running = self._running.pop(session.id, None) is not None self._finished[session.id] = session session._completion_event.set() self._write_checkpoint() - if was_running and session.notify_on_complete: - from tools.ansi_strip import strip_ansi - output_tail = strip_ansi(session.output_buffer[-2000:]) if session.output_buffer else "" notification = { "type": "completion", "session_id": session.id, @@ -1398,10 +1036,8 @@ class ProcessRegistry: "task_id": session.task_id, "owner_task_id": session.owner_task_id or session.task_id, "command": session.command, - "exit_code": session.exit_code, - "completion_reason": session.completion_reason, - "termination_source": session.termination_source, - "output": output_tail, + **self._exit_fields(session), + "output": _output_tail(session, 2000), # Stable producer identity across checkpoint recovery (unlike a # consumer-observed completion timestamp). "started_at": session.started_at, @@ -1409,6 +1045,14 @@ class ProcessRegistry: _redact_process_result(notification) self.completion_queue.put(notification) + @staticmethod + def _exit_fields(session: ProcessSession) -> dict: + return { + "exit_code": session.exit_code, + "completion_reason": session.completion_reason, + "termination_source": session.termination_source, + } + # ----- Query Methods ----- def is_completion_consumed(self, session_id: str) -> bool: @@ -1416,57 +1060,38 @@ class ProcessRegistry: return session_id in self._completion_consumed def is_session_waiting(self, session_id: str) -> bool: - """Whether a goal loop (``hermes_cli.goals`` wait barrier) should stay parked - on this session: still running AND, if it has ``watch_patterns``, none has - matched yet (a long-lived watcher unblocks on its trigger, not on exit). - Unknown/exited/already-fired sessions return False so a stale barrier can - never wedge the loop.""" - if not session_id: - return False + """Whether a goal loop (``hermes_cli.goals`` wait barrier) should stay parked on + this session: still running AND, with ``watch_patterns``, none matched yet (a + long-lived watcher unblocks on its trigger, not on exit). Unknown/exited/ + already-fired sessions return False so a stale barrier can never wedge the loop.""" with self._lock: - session = self._running.get(session_id) or self._finished.get(session_id) + session = (self._running.get(session_id) or self._finished.get(session_id)) if session_id else None if session is None: return False - try: + with suppress(Exception): self._refresh_detached_session(session) - except Exception: - pass - if session.exited: - return False - return not (session.watch_patterns and not session._watch_disabled and session._watch_hits > 0) + return not session.exited and not ( + session.watch_patterns and not session._watch_disabled and session._watch_hits > 0) def wait_for_pending_completions( - self, - task_id: Optional[str] = None, - *, - timeout: float | None = None, - poll_interval: float = 1.0, + self, task_id: Optional[str] = None, *, timeout: float | None = None, poll_interval: float = 1.0, ) -> dict: """Bounded linger for ``notify_on_complete`` background processes at one-shot exit. - - A one-shot CLI run (``hermes -q/-Q/-z``) exits when its turn ends; any - background process it spawned still holds a stdout pipe owned by the dying - parent and dies of SIGPIPE seconds later (Bot Mode handoff replies sent via - message_agent/bot_relay were the visible casualty). Only ``notify_on_complete`` - processes carry a completion contract — servers/daemons/watchers aren't the - parent's to wait for. - - ``task_id=None`` waits on every tracked process (a one-shot process hosts one - agent). ``timeout=None`` reads ``terminal.oneshot_completion_wait_seconds``; - ``<= 0`` disables. Each ``poll_interval`` pass re-reconciles child state so an - orphaned-pipe exit can't wedge the linger. Returns - ``{"waited": [...], "completed": [...], "timed_out": [...]}`` of session ids. - """ + A one-shot CLI run (``hermes -q/-Q/-z``) exits when its turn ends; a background + process it spawned still holds a stdout pipe owned by the dying parent and dies of + SIGPIPE seconds later (Bot Mode handoff replies were the visible casualty). Only + ``notify_on_complete`` processes carry a completion contract — servers/daemons/ + watchers aren't the parent's to wait for. ``task_id=None`` waits on every tracked + process; ``timeout=None`` reads ``terminal.oneshot_completion_wait_seconds`` (``<= 0`` + disables). Each pass re-reconciles child state so an orphaned-pipe exit can't wedge + the linger. Returns ``{"waited", "completed", "timed_out"}`` id lists.""" if timeout is None: timeout = self._oneshot_completion_wait_seconds() result: dict = {"waited": [], "completed": [], "timed_out": []} with self._lock: pending = [ - s - for s in self._running.values() - if s.notify_on_complete - and not s.exited - and (task_id is None or s.task_id == task_id) + s for s in self._running.values() + if s.notify_on_complete and not s.exited and (task_id is None or s.task_id == task_id) ] if not pending or timeout <= 0: return result @@ -1474,17 +1099,13 @@ class ProcessRegistry: logger.info( "One-shot exit lingering (bounded %ss) for %d notify_on_complete " "background process(es): %s", - timeout, - len(pending), - ", ".join(s.id for s in pending), - ) + timeout, len(pending), ", ".join(s.id for s in pending)) deadline = time.monotonic() + max(float(timeout), 0.0) interval = max(float(poll_interval), 0.05) try: from tools.interrupt import is_interrupted as _is_interrupted except Exception: - def _is_interrupted() -> bool: - return False + _is_interrupted = lambda: False # noqa: E731 interrupted = False for session in pending: try: @@ -1496,11 +1117,9 @@ class ProcessRegistry: if remaining <= 0: break # Reconcile first so orphaned-pipe and detached exits fire the event. - try: + with suppress(Exception): self._reconcile_local_exit(session) self._refresh_detached_session(session) - except Exception: - pass if session.exited: break session._completion_event.wait(min(remaining, interval)) @@ -1508,73 +1127,62 @@ class ProcessRegistry: # Stop waiting, but never let the interrupt skip the caller's durable # teardown (session flush, end_session) that follows. interrupted = True - if session.exited: - result["completed"].append(session.id) - else: - result["timed_out"].append(session.id) + result["completed" if session.exited else "timed_out"].append(session.id) if result["timed_out"]: logger.warning( "One-shot exit linger timed out after %ss with %d background " "process(es) still running: %s — they may be killed when this " "process exits.", - timeout, - len(result["timed_out"]), - ", ".join(result["timed_out"]), - ) + timeout, len(result["timed_out"]), ", ".join(result["timed_out"])) return result @staticmethod def _oneshot_completion_wait_seconds() -> float: - """Bounded linger (s) for one-shot exits with pending notify_on_complete - processes: ``terminal.oneshot_completion_wait_seconds`` (0 disables), 600 - if config is unreadable.""" - try: - return max(float(ProcessRegistry._config_value("terminal", "oneshot_completion_wait_seconds", 600.0)), 0.0) - except Exception: - return 600.0 + """Linger (s) for one-shot exits with pending notify_on_complete processes; 0 disables.""" + return ProcessRegistry._config_seconds("oneshot_completion_wait_seconds", 600.0) - def _drain_should_skip( - self, session_id: str, *, skip_poll_observed: bool = True - ) -> bool: - """Skip a completion the CLI agent already has this turn — consumed via - wait/log or observed inline via poll(). Gateway/tui watchers check only - ``is_completion_consumed`` so a read-only poll never suppresses their - autonomous delivery turn.""" - return session_id in self._completion_consumed or ( - skip_poll_observed and session_id in self._poll_observed - ) + def _drain_should_skip(self, session_id: str, *, skip_poll_observed: bool = True) -> bool: + """Skip a completion the CLI agent already has this turn — consumed via wait/log + or observed inline via poll(). Gateway/tui watchers check only + ``is_completion_consumed`` so a read-only poll never suppresses their turn.""" + return session_id in self._completion_consumed or (skip_poll_observed and session_id in self._poll_observed) @staticmethod def _surface_child_process_notifications() -> bool: - """Whether subagent-owned process notifications surface in the parent - (``delegation.surface_child_process_notifications``; suppress on any config - error — never crash the drain loop).""" + """``delegation.surface_child_process_notifications``; False on any config + error — never crash the drain loop.""" try: return bool(ProcessRegistry._config_value("delegation", "surface_child_process_notifications", False)) except Exception: return False + @staticmethod + def _owns_event(evt: dict, session_key: str, owns_event, is_async_delegation: bool) -> bool: + """Routing verdict for one drained event (see drain_notifications); False = requeue.""" + evt_session_key = str(evt.get("session_key") or "") + requires_positive_proof = is_async_delegation or bool(evt_session_key or evt.get("origin_ui_session_id")) + if owns_event is not None and requires_positive_proof: + try: + return bool(owns_event(evt)) + except Exception: + return False # fail closed — never leak on a broken check + if session_key and requires_positive_proof: + return evt_session_key == session_key + # Restored payloads from a previous process: an unfiltered drain cannot prove + # ownership, so leave them for the owner. + return not (is_async_delegation and evt.get("restored")) + def drain_notifications( - self, - session_key: str = "", - owns_event=None, - *, - skip_poll_observed: bool = True, + self, session_key: str = "", owns_event=None, *, skip_poll_observed: bool = True, ) -> "list[tuple[dict, str]]": """Pop all pending events and return ``(raw_event, formatted_text)`` pairs. - - Skips completions per ``_drain_should_skip``; gateway/TUI callers pass - ``skip_poll_observed=False``. - - Routing: async-delegation events always need ownership proof; ordinary - events need it once they carry ``session_key`` or ``origin_ui_session_id``. - ``owns_event(evt)`` (strongest; the TUI passes a compression-chain-aware - check so a post-compression session still claims its pre-compression - dispatches) consumes ONLY on True; ``session_key`` uses plain equality. - Non-owned routed events are re-queued for their owner. With no filter every - event is consumed (legacy single-session), except restored delegation - payloads, which stay fail-closed. - """ + Skips completions per ``_drain_should_skip`` (gateway/TUI pass + ``skip_poll_observed=False``). Routing (``_owns_event``): async-delegation events + always need ownership proof, ordinary events once they carry ``session_key`` or + ``origin_ui_session_id``; ``owns_event(evt)`` (strongest; the TUI passes a + compression-chain-aware check) consumes ONLY on True, ``session_key`` uses plain + equality; non-owned events are re-queued for their owner. No filter consumes + everything (legacy single-session) except restored delegation payloads (fail-closed).""" results: "list[tuple[dict, str]]" = [] requeue: "list[dict]" = [] # delegation.surface_child_process_notifications, read at most once per drain @@ -1586,45 +1194,22 @@ class ProcessRegistry: except Exception: break is_async_delegation = evt.get("type") == "async_delegation" - evt_session_key = str(evt.get("session_key") or "") - evt_origin_sid = str(evt.get("origin_ui_session_id") or "") - requires_positive_proof = is_async_delegation or bool( - evt_session_key or evt_origin_sid - ) - if owns_event is not None and requires_positive_proof: - try: - owned = bool(owns_event(evt)) - except Exception: - owned = False # fail closed — never leak on a broken check - if not owned: - requeue.append(evt) - continue - elif session_key and requires_positive_proof: - if evt_session_key != session_key: - requeue.append(evt) - continue - elif is_async_delegation and evt.get("restored"): - # Restored payloads from a previous process: an unfiltered drain - # cannot prove ownership, so leave them for the owner. + if not self._owns_event(evt, session_key, owns_event, is_async_delegation): requeue.append(evt) continue # Routing happened first so a foreign session cannot drop the owner's # event via its own consumed/observed state. _evt_sid = evt.get("session_id", "") if evt.get("type") == "completion" and self._drain_should_skip( - _evt_sid, skip_poll_observed=skip_poll_observed - ): + _evt_sid, skip_poll_observed=skip_poll_observed): continue - # Subagent-owned process notifications are suppressed by default — the # child's delegation result is the deliverable. Judge ownership on # owner_task_id (RAW spawning id; task_id is the container key, collapsed # by _resolve_container_task_id). Dropped, NOT requeued: children never # drain, so a requeue would pin the event forever. 'async_delegation' # is the result itself and is NEVER suppressed. - _evt_task_id = str( - evt.get("owner_task_id") or evt.get("task_id") or "" - ) + _evt_task_id = str(evt.get("owner_task_id") or evt.get("task_id") or "") if not is_async_delegation and _evt_task_id.startswith("sa-"): if surface_child is None: surface_child = self._surface_child_process_notifications() @@ -1633,67 +1218,48 @@ class ProcessRegistry: "Suppressed subagent-owned process notification " "(delegation.surface_child_process_notifications=false): " "type=%s session_id=%s task_id=%s", - evt.get("type", "completion"), - _evt_sid, - _evt_task_id, - ) + evt.get("type", "completion"), _evt_sid, _evt_task_id) continue - - text = format_process_notification(evt) - if text: + if text := format_process_notification(evt): results.append((evt, text)) for evt in requeue: self.completion_queue.put(evt) return results - # Minimum characters of the random suffix required for prefix resolution. - # Short prefixes ("p", "pr", "proc_1") are too collision-prone to act on. + # Minimum suffix chars for prefix resolution; "p"/"proc_1" are too collision-prone. _MIN_PREFIX_CHARS = 4 def get(self, session_id: str) -> Optional[ProcessSession]: - """Get a session by full ID or unique prefix (``proc_4dae`` / bare ``4dae``, - like git short hashes). Ambiguous or too-short prefixes resolve to None, - never to an arbitrary pick.""" + """Session by full ID or unique prefix (``proc_4dae`` / bare ``4dae``, like git + short hashes); ambiguous or too-short prefixes resolve to None, never a guess.""" with self._lock: session = self._running.get(session_id) or self._finished.get(session_id) - if session is None: - session = self._resolve_prefix(session_id) - return self._refresh_detached_session(session) + return self._refresh_detached_session(session if session is not None else self._resolve_prefix(session_id)) def _resolve_prefix(self, session_id: str) -> Optional[ProcessSession]: - """Resolve a unique session-ID prefix (prefix-only, unique hit; a bare hex - tail is normalized to ``proc_``). :meth:`get` tries exact first.""" - if not session_id or not isinstance(session_id, str): - return None - query = session_id.strip() + """Resolve a unique session-ID prefix (a bare hex tail is normalized to + ``proc_``); :meth:`get` tries exact first.""" + query = session_id.strip() if isinstance(session_id, str) else "" if not query: return None - # Allow the bare suffix form: "4dae56" -> "proc_4dae56". if not query.startswith("proc_"): query = f"proc_{query}" - suffix = query[len("proc_"):] - if len(suffix) < self._MIN_PREFIX_CHARS: + if len(query) - len("proc_") < self._MIN_PREFIX_CHARS: return None with self._lock: matches = [ - s - for store in (self._running, self._finished) - for sid, s in store.items() - if sid.startswith(query) + s for store in (self._running, self._finished) + for sid, s in store.items() if sid.startswith(query) ] - if len(matches) == 1: - return matches[0] - return None + return matches[0] if len(matches) == 1 else None def _reconcile_local_exit(self, session: "ProcessSession") -> None: """Reconcile ``session.exited`` against the real child state. - - The reader flips ``exited`` only at EOF; when the direct child has exited but - a descendant (e.g. a daemon from ``hermes update``) holds the pipe open, poll() + The reader flips ``exited`` only at EOF; when the direct child has exited but a + descendant (e.g. a daemon from ``hermes update``) holds the pipe open, poll() would report "running" forever. If ``Popen.poll()`` has an exit code, drain - readable bytes non-blocking and flip ``exited``; the stuck daemon reader thread - is reaped with the process. No-op for env/PTY, exited and detached sessions. - """ + readable bytes non-blocking and flip ``exited``. No-op for env/PTY, exited and + detached sessions.""" if session is None or session.exited: return proc = getattr(session, "process", None) @@ -1705,9 +1271,7 @@ class ProcessRegistry: return if rc is None: return # Direct child still running — reader block is legitimate. - # Best-effort non-blocking drain of whatever the reader hasn't consumed. - drained = "" stdout = getattr(proc, "stdout", None) if stdout is not None and not _IS_WINDOWS: try: @@ -1716,67 +1280,46 @@ class ProcessRegistry: flags = fcntl.fcntl(fd, fcntl.F_GETFL) fcntl.fcntl(fd, fcntl.F_SETFL, flags | os.O_NONBLOCK) try: - chunk = stdout.read() - if chunk: - drained = chunk if isinstance(chunk, str) else chunk.decode("utf-8", errors="replace") - except (BlockingIOError, OSError, ValueError): - pass + with suppress(BlockingIOError, OSError, ValueError): + chunk = stdout.read() + if chunk: + session.append_output(chunk if isinstance(chunk, str) else chunk.decode("utf-8", errors="replace")) finally: - try: + with suppress(Exception): fcntl.fcntl(fd, fcntl.F_SETFL, flags) - except Exception: - pass except Exception as e: logger.debug("Non-blocking drain failed for %s: %s", session.id, e) - with session._lock: - if drained: - session.output_buffer += drained - if len(session.output_buffer) > session.max_output_chars: - session.output_buffer = session.output_buffer[-session.max_output_chars:] - session.exited = True - if session.completion_reason != "killed": - session.exit_code = rc - session.completion_reason = "exited" + session.mark_exited(rc) logger.info( "Reconciled session %s: direct child exited with code %s but reader " "was still blocked (orphaned pipe). Flipped to exited.", - session.id, rc, - ) + session.id, rc) self._move_to_finished(session) + @staticmethod + def _status_head(session: ProcessSession) -> dict: + return {"session_id": session.id, "command": session.command, "status": "exited" if session.exited else "running"} + def poll(self, session_id: str) -> dict: """Check status and get new output for a background process.""" - from tools.ansi_strip import strip_ansi - session = self.get(session_id) if session is None: - return {"status": "not_found", "error": f"No process with ID {session_id}"} - + return _not_found(session_id) self._reconcile_local_exit(session) # orphaned-pipe reader guard - with session._lock: - output_preview = strip_ansi(session.output_buffer[-1000:]) if session.output_buffer else "" - + output_preview = _output_tail(session, 1000) result = { - "session_id": session.id, - "command": session.command, - "status": "exited" if session.exited else "running", - "pid": session.pid, - "uptime_seconds": int(time.time() - session.started_at), - "output_preview": output_preview, - } + **self._status_head(session), "pid": session.pid, + "uptime_seconds": int(time.time() - session.started_at), "output_preview": output_preview} if session.exited: - result["exit_code"] = session.exit_code - result["completion_reason"] = session.completion_reason - result["termination_source"] = session.termination_source + result.update(self._exit_fields(session)) # Read-only: record in _poll_observed (CLI inline dedup) but NOT in # _completion_consumed, or a status check would suppress the watcher's # autonomous delivery turn. See __init__. self._poll_observed.add(session_id) if session.detached: - result["detached"] = True - result["note"] = "Process recovered after restart -- output history unavailable" + result.update(detached=True, note="Process recovered after restart -- output history unavailable") return result def read_log(self, session_id: str, offset: int | None = None, limit: int = 200) -> dict: @@ -1785,14 +1328,11 @@ class ProcessRegistry: session = self.get(session_id) if session is None: - return {"status": "not_found", "error": f"No process with ID {session_id}"} - + return _not_found(session_id) with session._lock: full_output = strip_ansi(session.output_buffer) - lines = full_output.splitlines() total_lines = len(lines) - # offset=None -> last N lines; an explicit offset=0 means the HEAD (don't # conflate the two via falsiness). if offset is None and limit > 0: @@ -1802,155 +1342,92 @@ class ProcessRegistry: offset = offset or 0 selected = lines[offset:offset + limit] stop = slice(offset, offset + limit).indices(total_lines)[1] - observed_completion_output = ( - total_lines == 0 or (bool(selected) and stop == total_lines) - ) - + observed_completion_output = total_lines == 0 or (bool(selected) and stop == total_lines) result = { - "session_id": session.id, - "command": session.command, - "status": "exited" if session.exited else "running", - "output": "\n".join(selected), - "total_lines": total_lines, - "showing": f"{len(selected)} lines", - } + **self._status_head(session), "output": "\n".join(selected), + "total_lines": total_lines, "showing": f"{len(selected)} lines"} if session.exited and observed_completion_output: self._completion_consumed.add(session_id) return result def wait(self, session_id: str, timeout: int = None) -> dict: """Block until the process exits, the timeout elapses, or the user interrupts. - ``timeout`` defaults to (and is clamped by) TERMINAL_TIMEOUT. Returns a dict - with status exited|timeout|interrupted|not_found|error and an output snapshot. - """ - from tools.ansi_strip import strip_ansi + with status exited|timeout|interrupted|not_found|error and an output snapshot.""" from tools.interrupt import is_interrupted as _is_interrupted try: - default_timeout = int(os.getenv("TERMINAL_TIMEOUT", "180")) + max_timeout = int(os.getenv("TERMINAL_TIMEOUT", "180")) except (ValueError, TypeError): - default_timeout = 180 - max_timeout = default_timeout - requested_timeout = timeout - timeout_note = None - + max_timeout = 180 # The schema says minimum=1 but not every caller enforces it; timeout=0 is # falsy and would silently fall through to the default wait. - if requested_timeout is not None and requested_timeout <= 0: - return { - "status": "error", - "error": f"timeout must be positive (got {requested_timeout})", - } - - if requested_timeout and requested_timeout > max_timeout: + if timeout is not None and timeout <= 0: + return {"status": "error", "error": f"timeout must be positive (got {timeout})"} + timeout_note = None + effective_timeout = timeout or max_timeout + if timeout and timeout > max_timeout: effective_timeout = max_timeout - timeout_note = ( - f"Requested wait of {requested_timeout}s was clamped " - f"to configured limit of {max_timeout}s" - ) - else: - effective_timeout = requested_timeout or max_timeout - + timeout_note = f"Requested wait of {timeout}s was clamped to configured limit of {max_timeout}s" session = self.get(session_id) if session is None: - return {"status": "not_found", "error": f"No process with ID {session_id}"} - + return _not_found(session_id) deadline = time.monotonic() + effective_timeout - while time.monotonic() < deadline: session = self._refresh_detached_session(session) if session is None: - return {"status": "not_found", "error": f"No process with ID {session_id}"} + return _not_found(session_id) self._reconcile_local_exit(session) # orphaned-pipe reader guard + result = None if session.exited: self._completion_consumed.add(session_id) result = self._exit_snapshot(session, "exited") elif _is_interrupted(): result = { - "status": "interrupted", - "command": session.command, - "output": strip_ansi(session.output_buffer[-1000:]), - "note": "User sent a new message -- wait interrupted", - } - else: - result = None + "status": "interrupted", "command": session.command, "output": _output_tail(session, 1000), + "note": "User sent a new message -- wait interrupted"} if result is not None: if timeout_note: result["timeout_note"] = timeout_note return result - remaining = deadline - time.monotonic() if remaining <= 0: break session._completion_event.wait(timeout=min(1.0, remaining)) - result = { - "status": "timeout", - "command": session.command, - "output": strip_ansi(session.output_buffer[-1000:]), - # Not a failure — models re-issued identical waits after misreading - # this result as an error. - "process_running": True, - } - uptime = time.time() - session.started_at if session.started_at else None + "status": "timeout", "command": session.command, "output": _output_tail(session, 1000), + # Not a failure — models re-issued identical waits after misreading this as an error. + "process_running": True} base_note = ( - f"Wait window of {effective_timeout}s elapsed — the process is " - "still running. This is not an error." - ) - if uptime is not None: - base_note += f" Uptime: {int(uptime)}s." - if session.notify_on_complete: - base_note += ( - " notify_on_complete is set: you will be notified on exit — " - "do more work instead of waiting again." - ) - else: - base_note += ( - " Poll again later or use terminal(background=true, " - "notify_on_complete=true) next time for automatic notification." - ) - if timeout_note: - result["timeout_note"] = f"{timeout_note}. {base_note}" - else: - result["timeout_note"] = base_note + f"Wait window of {effective_timeout}s elapsed — the process is still running. This is not an error.") + if session.started_at: + base_note += f" Uptime: {int(time.time() - session.started_at)}s." + base_note += ( + " notify_on_complete is set: you will be notified on exit — do more work instead of waiting again." + if session.notify_on_complete else + " Poll again later or use terminal(background=true, " + "notify_on_complete=true) next time for automatic notification.") + result["timeout_note"] = f"{timeout_note}. {base_note}" if timeout_note else base_note return result @staticmethod def _exit_snapshot(session: ProcessSession, status: str) -> dict: """Result dict for an exited session: exit metadata + last 2000 chars of output.""" - from tools.ansi_strip import strip_ansi - return { - "status": status, - "command": session.command, - "exit_code": session.exit_code, - "completion_reason": session.completion_reason, - "termination_source": session.termination_source, - "output": strip_ansi(session.output_buffer[-2000:]), - } + "status": status, "command": session.command, + **ProcessRegistry._exit_fields(session), "output": _output_tail(session, 2000)} def kill_process( - self, - session_id: str, - *, - source: str = "process.kill", - consume_output: bool = True, + self, session_id: str, *, source: str = "process.kill", consume_output: bool = True, ) -> dict: """Kill a background process and return its output snapshot. - ``consume_output`` is true for explicit tool/RPC kills (the caller sees the output). Bulk cleanup passes false so it doesn't suppress an autonomous - completion notification — except abandoned-turn reaping - (``kill_started_since``), which passes true so a killed abandoned process - can't enqueue a follow-up reviving work the timeout stopped. - """ - from tools.ansi_strip import strip_ansi - + completion notification — except abandoned-turn reaping (``kill_started_since``), + which passes true so a killed abandoned process can't revive stopped work.""" session = self.get(session_id) if session is None: - return {"status": "not_found", "error": f"No process with ID {session_id}"} - + return _not_found(session_id) if session.exited: # A double-forked descendant may still be alive in the systemd scope even # though the main process exited — stop the scope to reap survivors. @@ -1963,12 +1440,10 @@ class ProcessRegistry: if consume_output: self._completion_consumed.add(session_id) return result - try: early = self._signal_kill(session, session_id, consume_output) if early is not None: return early - # Additive to the PID kill: stopping the scope reaps double-forked # descendants reparented inside the cgroup. if session.systemd_unit: @@ -1976,7 +1451,7 @@ class ProcessRegistry: # Capture output, mark consumed, THEN expose ``exited`` to watcher tasks — # closes the delayed-notification race without losing the transcript. with session._lock: - output = strip_ansi(session.output_buffer[-2000:]) + output = _output_tail(session, 2000) if consume_output: self._completion_consumed.add(session_id) session.exited = True @@ -1986,12 +1461,8 @@ class ProcessRegistry: self._move_to_finished(session) self._write_checkpoint() return { - "status": "killed", - "session_id": session.id, - "completion_reason": session.completion_reason, - "termination_source": session.termination_source, - "output": output, - } + "status": "killed", "session_id": session.id, "completion_reason": session.completion_reason, + "termination_source": session.termination_source, "output": output} except Exception as e: return {"status": "error", "error": str(e)} @@ -1999,8 +1470,6 @@ class ProcessRegistry: """Deliver the kill via PTY, local Popen tree, sandbox exec or recovered host PID. Returns a final result dict when the kill cannot proceed (recycled/dead recovered PID, or no runtime handle), else None.""" - from tools.ansi_strip import strip_ansi - if session._pty: try: session._pty.terminate(force=True) @@ -2023,7 +1492,7 @@ class ProcessRegistry: with session._lock: session.exited = True session.exit_code = None - output = strip_ansi(session.output_buffer[-2000:]) + output = _output_tail(session, 2000) if consume_output: self._completion_consumed.add(session_id) self._move_to_finished(session) @@ -2032,140 +1501,97 @@ class ProcessRegistry: else: return { "status": "error", - "error": ( - "Recovered process cannot be killed after restart because " - "its original runtime handle is no longer available" - ), + "error": "Recovered process cannot be killed after restart because " + "its original runtime handle is no longer available", } return None - def _live_session(self, session_id: str): - """``(session, None)`` for a running session, else ``(None, error_result)``.""" + def _stdin_op(self, session_id: str, pty_op, pipe_op, ok: dict) -> dict: + """Run a stdin operation on a running session — ``pty_op(pty)`` under PTY mode, + else ``pipe_op(stdin)`` on the Popen pipe — and return *ok* on success.""" session = self.get(session_id) if session is None: - return None, {"status": "not_found", "error": f"No process with ID {session_id}"} + return _not_found(session_id) if session.exited: - return None, {"status": "already_exited", "error": "Process has already finished"} - return session, None + return {"status": "already_exited", "error": "Process has already finished"} + try: + if session._pty: + pty_op(session._pty) + elif not session.process or not session.process.stdin: + return {"status": "error", "error": "Process stdin not available (non-local backend or stdin closed)"} + else: + pipe_op(session.process.stdin) + return ok + except Exception as e: + return {"status": "error", "error": str(e)} def write_stdin(self, session_id: str, data: str) -> dict: """Send raw data to a running process's stdin (no newline appended).""" - session, err = self._live_session(session_id) - if err: - return err - # PTY mode -- write through pty handle. - if session._pty: - try: - # pywinpty expects str on Windows; ptyprocess expects bytes on POSIX. - if _IS_WINDOWS: - pty_data = data.decode("utf-8") if isinstance(data, bytes) else str(data) - else: - # surrogateescape: a PTY is a byte stream — round-trip the - # original bytes instead of crashing on surrogate content. - pty_data = data.encode("utf-8", "surrogateescape") if isinstance(data, str) else data - session._pty.write(pty_data) - return {"status": "ok", "bytes_written": len(data)} - except Exception as e: - return {"status": "error", "error": str(e)} + def via_pty(pty): + # pywinpty expects str on Windows; ptyprocess expects bytes on POSIX. + if _IS_WINDOWS: + pty.write(data.decode("utf-8") if isinstance(data, bytes) else str(data)) + else: + # surrogateescape: a PTY is a byte stream — round-trip the original + # bytes instead of crashing on surrogate content. + pty.write(data.encode("utf-8", "surrogateescape") if isinstance(data, str) else data) - # Popen mode -- write through stdin pipe - if not session.process or not session.process.stdin: - return {"status": "error", "error": "Process stdin not available (non-local backend or stdin closed)"} - try: - session.process.stdin.write(data) - session.process.stdin.flush() - return {"status": "ok", "bytes_written": len(data)} - except Exception as e: - return {"status": "error", "error": str(e)} + def via_pipe(stdin): + stdin.write(data) + stdin.flush() + return self._stdin_op(session_id, via_pty, via_pipe, {"status": "ok", "bytes_written": len(data)}) def submit_stdin(self, session_id: str, data: str = "") -> dict: """Send data + newline to stdin (like pressing Enter). - On a Windows PTY, Enter is a carriage return: ConPTY treats ``\\r`` as - end-of-line and a bare ``\\n`` through pywinpty is NOT a line terminator — - the child's blocking line read (``readline()``, Go ``bufio.Scanner`` in - ``gh auth login``) never returns and the process hangs looking healthy. - ``\\r\\n`` gives it both; POSIX PTYs and pipes keep ``\\n``. - """ + end-of-line and a bare ``\\n`` through pywinpty is NOT a line terminator — the + child's blocking line read (``readline()``, Go ``bufio.Scanner``) never returns + and the process hangs looking healthy. ``\\r\\n`` gives it both; POSIX keeps ``\\n``.""" session = self.get(session_id) - is_windows_pty = bool(_IS_WINDOWS and session is not None and session._pty) - return self.write_stdin(session_id, data + ("\r\n" if is_windows_pty else "\n")) + return self.write_stdin(session_id, data + ("\r\n" if _IS_WINDOWS and session and session._pty else "\n")) def request_close_terminal(self, session_id: str) -> dict: - """Ask the desktop GUI to close this process's read-only terminal tab. - - Does NOT kill the process — output keeps buffering and the tab can be - reopened from the status stack. Errors when no UI close sink is wired.""" - sink = self.on_close - if sink is None: - return { - "status": "error", - "error": "close_terminal is only available in the Hermes desktop app.", - } + """Ask the desktop GUI to close this process's read-only terminal tab. Does NOT + kill the process — output keeps buffering and the tab can be reopened from the + status stack. Errors when no UI close sink is wired.""" + if self.on_close is None: + return {"status": "error", "error": "close_terminal is only available in the Hermes desktop app."} # The session may already be finished (or pruned) — the tab can still # linger and be closed, so a missing session is not an error here. - session = self.get(session_id) try: - sink(session, session_id) + self.on_close(self.get(session_id), session_id) except Exception as e: return {"status": "error", "error": str(e)} return { - "status": "ok", - "closed": session_id, - "note": ( - "Closed the read-only terminal tab. The process was not killed; " - "its output remains available and the user can reopen the tab " - "from the status stack." - ), - } + "status": "ok", "closed": session_id, + "note": "Closed the read-only terminal tab. The process was not killed; " + "its output remains available and the user can reopen the tab " + "from the status stack."} def close_stdin(self, session_id: str) -> dict: """Close a running process's stdin / send EOF without killing the process.""" - session, err = self._live_session(session_id) - if err: - return err - - if session._pty: - try: - session._pty.sendeof() - return {"status": "ok", "message": "EOF sent"} - except Exception as e: - return {"status": "error", "error": str(e)} - - if not session.process or not session.process.stdin: - return {"status": "error", "error": "Process stdin not available (non-local backend or stdin closed)"} - try: - session.process.stdin.close() - return {"status": "ok", "message": "stdin closed"} - except Exception as e: - return {"status": "error", "error": str(e)} + session = self.get(session_id) + msg = "EOF sent" if session is not None and session._pty else "stdin closed" + return self._stdin_op( + session_id, lambda pty: pty.sendeof(), lambda stdin: stdin.close(), {"status": "ok", "message": msg}) def count_running(self) -> int: - """O(1) count of running processes for status-bar polling; CPython dict - ``len()`` is atomic so no lock is needed.""" - try: - return len(self._running) - except Exception: - return 0 + """O(1) running count for status-bar polling; dict ``len()`` is atomic, no lock.""" + return len(self._running) def list_sessions(self, task_id: str = None, session_key: str = None) -> list: - """List running and recently-finished processes for ``task_id`` and/or - ``session_key``. Cross-task entries that share the gateway session (a - forgotten preview server blocking session reset) are flagged - ``"session_scoped": true``.""" + """Running and recently-finished processes for ``task_id`` and/or ``session_key``; + cross-task entries sharing the gateway session (a forgotten preview server + blocking session reset) are flagged ``"session_scoped": true``.""" with self._lock: all_sessions = list(self._running.values()) + list(self._finished.values()) - all_sessions = [self._refresh_detached_session(s) for s in all_sessions] - if task_id or session_key: all_sessions = [ s for s in all_sessions - if (task_id and s.task_id == task_id) - or (session_key and s.session_key == session_key) + if (task_id and s.task_id == task_id) or (session_key and s.session_key == session_key) ] - result = [] for s in all_sessions: entry = { @@ -2182,8 +1608,7 @@ class ProcessRegistry: entry["session_scoped"] = True # Trigger metadata for goal-loop judges (a watcher may never exit). if s.watch_patterns and not s._watch_disabled: - entry["watch_patterns"] = list(s.watch_patterns) - entry["watch_hit"] = s._watch_hits > 0 + entry.update(watch_patterns=list(s.watch_patterns), watch_hit=s._watch_hits > 0) if s.notify_on_complete: entry["notify_on_complete"] = True if s.exited: @@ -2206,21 +1631,18 @@ class ProcessRegistry: return any(not s.exited and predicate(s) for s in self._running.values()) def has_active_processes(self, task_id: str) -> bool: - """Check if there are active (running) processes for a task_id.""" + """Whether any process for ``task_id`` is still running.""" return self._any_running(lambda s: s.task_id == task_id) - def has_active_for_session( - self, session_key: str, max_active_age: Optional[float] = None, - ) -> bool: + def has_active_for_session(self, session_key: str, max_active_age: Optional[float] = None) -> bool: """Active processes for a gateway session key. Processes older than - ``max_active_age`` seconds are ignored as stale so a forgotten - ``http.server`` can't freeze session idle/daily reset forever; ``None`` - keeps legacy behaviour (any running process blocks).""" + ``max_active_age`` seconds are ignored as stale so a forgotten ``http.server`` + can't freeze session idle/daily reset forever; ``None`` keeps legacy behaviour + (any running process blocks).""" now = time.time() return self._any_running( lambda s: s.session_key == session_key - and (max_active_age is None or (now - s.started_at) < max_active_age) - ) + and (max_active_age is None or (now - s.started_at) < max_active_age)) def has_any_active(self) -> bool: """Whether ANY background process is running — scale-to-zero must not @@ -2232,72 +1654,38 @@ class ProcessRegistry: only processes absent from the starting snapshot belong to the abandoned turn; older ones intentionally span turns and must survive.""" with self._lock: - return frozenset( - s.id - for s in self._running.values() - if s.task_id == task_id and not s.exited - ) + return frozenset(s.id for s in self._running.values() if s.task_id == task_id and not s.exited) - def kill_started_since( - self, - task_id: str, - baseline_ids, - *, - source: str, - ) -> int: + def kill_started_since(self, task_id: str, baseline_ids, *, source: str) -> int: """Kill ``task_id`` processes created after ``baseline_ids``. Output is consumed so an abandoned turn can't enqueue a follow-up reviving work the timeout deliberately stopped.""" - return self.kill_all( - task_id, - exclude_ids=frozenset(baseline_ids or ()), - source=source, - consume_output=True, - ) + return self.kill_all(task_id, exclude_ids=frozenset(baseline_ids or ()), source=source, consume_output=True) def kill_all( - self, - task_id: Optional[str] = None, - *, - exclude_ids: frozenset = frozenset(), - source: str = "kill_all", - consume_output: bool = False, - ) -> int: + self, task_id: Optional[str] = None, *, exclude_ids: frozenset = frozenset(), + source: str = "kill_all", consume_output: bool = False) -> int: """Kill all running processes, optionally filtered by task_id. Returns count killed.""" with self._lock: targets = [ s for s in self._running.values() - if (task_id is None or s.task_id == task_id) - and s.id not in exclude_ids - and not s.exited + if (task_id is None or s.task_id == task_id) and s.id not in exclude_ids and not s.exited ] - - killed = 0 - for session in targets: - result = self.kill_process( - session.id, - source=source, - consume_output=consume_output, - ) - if result.get("status") in {"killed", "already_exited"}: - killed += 1 - return killed + return sum( + self.kill_process(s.id, source=source, consume_output=consume_output).get("status") + in {"killed", "already_exited"} + for s in targets) # ----- Cleanup / Pruning ----- def _prune_if_needed(self): - """Remove oldest finished sessions if over MAX_PROCESSES. Must hold _lock.""" - # First prune expired finished sessions + """Drop expired finished sessions, then the oldest survivor while over + MAX_PROCESSES. Must hold _lock.""" now = time.time() - expired = [ - sid for sid, s in self._finished.items() - if (now - s.started_at) > FINISHED_TTL_SECONDS - ] - if len(self._running) + len(self._finished) - len(expired) >= MAX_PROCESSES: - # Still over the limit: also drop the oldest surviving finished session. - survivors = [sid for sid in self._finished if sid not in expired] - if survivors: - expired.append(min(survivors, key=lambda sid: self._finished[sid].started_at)) + expired = [sid for sid, s in self._finished.items() if (now - s.started_at) > FINISHED_TTL_SECONDS] + over_cap = len(self._running) + len(self._finished) - len(expired) >= MAX_PROCESSES + if over_cap and (survivors := [sid for sid in self._finished if sid not in expired]): + expired.append(min(survivors, key=lambda sid: self._finished[sid].started_at)) for sid in expired: del self._finished[sid] # Belt-and-suspenders against module-lifetime growth: forget consumed / @@ -2308,159 +1696,110 @@ class ProcessRegistry: # ----- Checkpoint (crash recovery) ----- - def _write_checkpoint( - self, - extra_entries: Optional[List[Dict[str, Any]]] = None, - ): - """Write running process metadata to checkpoint file atomically.""" + def _write_checkpoint(self, extra_entries: Optional[List[Dict[str, Any]]] = None): + """Write running process metadata to the checkpoint file atomically.""" try: with self._lock: entries = [] for s in self._running.values(): - if not s.exited: - # Backfill the start time so recovery can detect PID recycling - # even for sessions spawned before this field existed. - if s.host_start_time is None and s.pid_scope == "host" and s.pid: - s.host_start_time = self._safe_host_start_time(s.pid) - entry = {"session_id": s.id, **{f: getattr(s, f) for f in _CHECKPOINT_FIELDS}} - # Redact inline credentials before persisting: the file lives - # at ~/.hermes/processes.json. Recovery uses command only for - # display (adoption re-validates the PID, never re-runs it), - # so masking is lossless. - entry["command"] = redact_sensitive_text(s.command, code_file=True) - entry["owner_task_id"] = s.owner_task_id or s.task_id - entries.append(entry) + if s.exited: + continue + # Backfill the start time so recovery can detect PID recycling + # even for sessions spawned before this field existed. + if s.host_start_time is None and s.pid_scope == "host" and s.pid: + s.host_start_time = self._safe_host_start_time(s.pid) + entry = {"session_id": s.id, **{f: getattr(s, f) for f in _CHECKPOINT_FIELDS}} + # Redact inline credentials before persisting (~/.hermes/processes.json). + # Recovery uses command only for display (adoption re-validates the + # PID, never re-runs it), so masking is lossless. + entry["command"] = redact_sensitive_text(s.command, code_file=True) + entry["owner_task_id"] = s.owner_task_id or s.task_id + entries.append(entry) if extra_entries: tracked_ids = {item.get("session_id") for item in entries} - entries.extend( - item - for item in extra_entries - if item.get("session_id") not in tracked_ids - ) - - # Atomic write to avoid corruption on crash + entries.extend(item for item in extra_entries if item.get("session_id") not in tracked_ids) from utils import atomic_json_write atomic_json_write(CHECKPOINT_PATH, entries) except Exception as e: logger.debug("Failed to write checkpoint file: %s", e, exc_info=True) def recover_from_checkpoint(self) -> int: - """ - On gateway startup, probe PIDs from checkpoint file. - - Returns the number of processes recovered as detached. - """ + """On gateway startup, probe PIDs from the checkpoint file; returns how many + were recovered as detached sessions.""" if not CHECKPOINT_PATH.exists(): return 0 - try: entries = json.loads(CHECKPOINT_PATH.read_text(encoding="utf-8")) except Exception: return 0 - recovered = 0 unresolved_scope_entries: List[Dict[str, Any]] = [] for entry in entries: - pid = entry.get("pid") + pid, pid_scope = entry.get("pid"), entry.get("pid_scope", "host") if not pid: continue - - pid_scope = entry.get("pid_scope", "host") - if pid_scope != "host": - # In-sandbox PIDs mean nothing once the environment handle is gone. + if pid_scope != "host": # in-sandbox PIDs mean nothing once the env handle is gone logger.info( "Skipping recovery for non-host process: %s (pid=%s, scope=%s)", - entry.get("command", "unknown")[:60], - pid, - pid_scope, - ) + entry.get("command", "unknown")[:60], pid, pid_scope) continue - # Alive AND the same process: across a restart the kernel may have # recycled the PID onto a stranger, and adopting it would let a later # kill tree-kill e.g. a browser. - recorded_start = entry.get("host_start_time") - if not self._host_pid_is_ours(pid, recorded_start): + if not self._host_pid_is_ours(pid, entry.get("host_start_time")): if self._is_host_pid_alive(pid): logger.info( "Not recovering session %s: pid %d is alive but its " "start time no longer matches — PID was recycled onto " "an unrelated process; refusing to adopt it.", - entry.get("session_id", "?"), pid, - ) + entry.get("session_id", "?"), pid) systemd_unit = entry.get("systemd_unit", "") if systemd_unit and not _stop_systemd_unit(systemd_unit): logger.warning( "Could not reap persisted scope %s for dead wrapper pid %s; " "retaining checkpoint entry for the next startup", - systemd_unit, - pid, - ) + systemd_unit, pid) unresolved_scope_entries.append(entry) continue - fields = {f: entry.get(f, _CHECKPOINT_DEFAULTS[f]) for f in _CHECKPOINT_FIELDS} fields.update( command=entry.get("command", "unknown"), owner_task_id=entry.get("owner_task_id", "") or entry.get("task_id", ""), - pid=pid, - host_start_time=recorded_start, - pid_scope=pid_scope, - started_at=entry.get("started_at", time.time()), - ) - session = ProcessSession( - id=entry["session_id"], - detached=True, # Can't read output, but can report status + kill - **fields, - ) + started_at=entry.get("started_at", time.time())) + # detached: can't read output, but can report status + kill + session = ProcessSession(id=entry["session_id"], detached=True, **fields) with self._lock: self._running[session.id] = session recovered += 1 logger.info("Recovered detached process: %s (pid=%d)", session.command[:60], pid) - # Re-enqueue watcher so gateway can resume notifications if session.watcher_interval > 0: self.pending_watchers.append({ "session_id": session.id, "check_interval": session.watcher_interval, "session_key": session.session_key, - "platform": session.watcher_platform, - "chat_id": session.watcher_chat_id, - "user_id": session.watcher_user_id, - "user_name": session.watcher_user_name, - "thread_id": session.watcher_thread_id, - "message_id": session.watcher_message_id, + **{key: getattr(session, f"watcher_{key}") for key in _WATCHER_ROUTE_KEYS}, "notify_on_complete": session.notify_on_complete, "parent_session_id": session.parent_session_id, }) - self._write_checkpoint(extra_entries=unresolved_scope_entries) - return recovered -# Module-level singleton process_registry = ProcessRegistry() -# Notification rendering lives in tools.process_registry_notifications; the names are -# re-exported here so `from tools.process_registry import format_process_notification` -# and `patch("tools.process_registry._x")` keep resolving. +# Notification rendering lives in tools.process_registry_notifications; re-exported so +# `from tools.process_registry import format_process_notification` and +# `patch("tools.process_registry._x")` keep resolving. from tools.process_registry_notifications import ( # noqa: F401,E402 - _delegation_attribution_line, - _delegation_config, - _delegation_model_not_found, - _delegation_model_not_found_notice, - _format_age, - _format_async_delegation, - _model_not_found_patterns, - format_process_notification, + _delegation_attribution_line, _delegation_config, _delegation_model_not_found, + _delegation_model_not_found_notice, _format_age, _format_async_delegation, + _model_not_found_patterns, format_process_notification, ) -# --------------------------------------------------------------------------- -# Registry -- the "process" tool schema + handler -# --------------------------------------------------------------------------- +# --- the "process_manage" tool schema + handler ----------------------------------- from tools.registry import registry, tool_error PROCESS_SCHEMA = { @@ -2514,37 +1853,30 @@ def _redact_process_result(result: dict) -> dict: """Redact secrets from background-process output before it reaches the model, session.db and CLI, mirroring the foreground ``terminal`` redaction so the two surfaces can't diverge. Respects ``security.redact_secrets``; ``redact_terminal_output`` - picks ``code_file`` from the recorded command. The command itself is redacted too. - """ + picks ``code_file`` from the recorded command. The command itself is redacted too.""" if not isinstance(result, dict): return result from agent.redact import redact_sensitive_text, redact_terminal_output command = result.get("command") or "" - for field in ("output", "output_preview"): - value = result.get(field) - if isinstance(value, str) and value: - result[field] = redact_terminal_output(value, command) - if isinstance(result.get("command"), str) and result["command"]: - result["command"] = redact_sensitive_text(result["command"], code_file=True) + for key in ("output", "output_preview"): + if isinstance(value := result.get(key), str) and value: + result[key] = redact_terminal_output(value, command) + if isinstance(command, str) and command: + result["command"] = redact_sensitive_text(command, code_file=True) return result def _list_processes(task_id) -> dict: - # Surface session-scoped background processes (e.g. a forgotten preview - # server) in addition to this task's own — they share the gateway - # session_key and can block session reset. - try: + # Also surface session-scoped background processes (e.g. a forgotten preview + # server): they share the gateway session_key and can block session reset. + session_key = "" + with suppress(Exception): from tools.approval import get_current_session_key session_key = get_current_session_key(default="") or "" - except Exception: - session_key = "" - return { - "processes": [ - _redact_process_result(p) - for p in process_registry.list_sessions(task_id=task_id, session_key=session_key or None) - ] - } + return {"processes": [ + _redact_process_result(p) + for p in process_registry.list_sessions(task_id=task_id, session_key=session_key or None)]} # action -> (handler(session_id, args) -> dict, redact output?). Output-bearing @@ -2564,7 +1896,6 @@ def _handle_process(args, **kw): action = args.get("action", "") # Coerce to string — some models send session_id as an integer session_id = str(args.get("session_id", "")) if args.get("session_id") is not None else "" - if action == "list": return json.dumps(_list_processes(kw.get("task_id")), ensure_ascii=False) if action in _SESSION_ACTIONS: diff --git a/tools/process_registry_notifications.py b/tools/process_registry_notifications.py index ccd8d2fd60..b470be43a3 100644 --- a/tools/process_registry_notifications.py +++ b/tools/process_registry_notifications.py @@ -7,6 +7,7 @@ conversation by the CLI drain loop, the gateway, and the TUI. """ import time +from contextlib import suppress def _format_age(seconds: float) -> str: @@ -19,19 +20,15 @@ def _format_age(seconds: float) -> str: return f"{s}s" m, s = divmod(s, 60) if m < 60: - return f"{m}m" if s == 0 else f"{m}m{s}s" + return f"{m}m" + (f"{s}s" if s else "") h, m = divmod(m, 60) - return f"{h}h" if m == 0 else f"{h}h{m}m" + return f"{h}h" + (f"{m}m" if m else "") def _model_not_found_patterns() -> "list[str]": - """Model-not-found phrases shared with the failover classifier. - - Imported from ``agent.error_classifier`` so the batch renderer applies the - same classification the failover path uses (no hand-copied list to drift). - Fails open to a minimal built-in set so an import problem never hides the - per-task blocks. - """ + """Model-not-found phrases from ``agent.error_classifier`` (same classification + the failover path uses, no hand-copied list to drift); a minimal built-in set + if the import fails so per-task blocks are never hidden.""" try: from agent.error_classifier import _MODEL_NOT_FOUND_PATTERNS @@ -42,11 +39,8 @@ def _model_not_found_patterns() -> "list[str]": def _delegation_config() -> dict: """Active delegation config (model/provider/fallbacks); ``{}`` on any error. - - Mirrors ``tools.delegate_tool._load_config`` lazily so the renderer sees the - same model/provider the dispatcher used without importing the heavy - delegation module at import time. - """ + Lazy ``tools.delegate_tool._load_config`` so the renderer sees the dispatcher's + model/provider without importing the heavy delegation module at import time.""" try: from tools.delegate_tool import _load_config as _cfg @@ -57,24 +51,15 @@ def _delegation_config() -> dict: def _delegation_model_not_found(results, config) -> bool: """True when a result reflects a config-level model_not_found rejection. - Requires both a model-not-found phrase AND the currently-configured model name in the same error/summary text, so a stale task failing on a - different (removed) model is not mis-attributed to the config. - """ - model = (config or {}).get("model") + different (removed) model is not mis-attributed to the config.""" + model = str((config or {}).get("model") or "").lower() if not model: return False - model = str(model).lower() - for r in results or []: - text = " ".join( - str(part) for part in (r.get("error"), r.get("summary")) if part - ).lower() - if not text or model not in text: - continue - if any(p in text for p in _model_not_found_patterns()): - return True - return False + patterns = _model_not_found_patterns() + texts = (" ".join(str(x) for x in (r.get("error"), r.get("summary")) if x).lower() for r in results or []) + return any(model in text and any(p in text for p in patterns) for text in texts) def _delegation_model_not_found_notice(results) -> "list[str] | None": @@ -89,37 +74,34 @@ def _delegation_model_not_found_notice(results) -> "list[str] | None": f'"{model}" was rejected by provider "{provider}" ' "(HTTP 400: not a valid model ID).", "Every task in this batch failed for this reason before doing any work.", - "Check Settings → Advanced → Subagent Model (or: " - "hermes config get delegation.model).", + "Check Settings → Advanced → Subagent Model (or: hermes config get delegation.model).", ] - try: + with suppress(Exception): from hermes_cli.fallback_config import get_fallback_chain if not get_fallback_chain(config): - lines.append( - "No fallback chain is configured, so no failover was attempted." - ) - except Exception: - pass + lines.append("No fallback chain is configured, so no failover was attempted.") return lines _TRUNCATED_SUMMARY_NOTE = ( "[TRUNCATED — subagent hit its iteration cap; the summary below " "may be incomplete. Verify before relying on it, or re-dispatch " - "the unfinished part.]" -) + "the unfinished part.]") def _is_truncated(entry: dict) -> bool: return bool(entry.get("truncated") or entry.get("exit_reason") == "max_iterations") -def _dispatched_line(dispatched_at, completed_at) -> "str | None": - if not isinstance(dispatched_at, (int, float)): - return None - ts = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(dispatched_at)) - return f"Dispatched: {ts} ({_format_age(completed_at - dispatched_at)} ago)" +def _header_lines(evt: dict, title: str, intro: str, completed_at: float) -> "list[str]": + """Shared preamble: title, intro, blank, dispatch time and task-source lines.""" + lines = [title, intro, ""] + dispatched_at = evt.get("dispatched_at") + if isinstance(dispatched_at, (int, float)): + ts = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(dispatched_at)) + lines.append(f"Dispatched: {ts} ({_format_age(completed_at - dispatched_at)} ago)") + return lines def _task_source_lines(evt: dict) -> "list[str]": @@ -131,6 +113,10 @@ def _task_source_lines(evt: dict) -> "list[str]": return lines +def _role_model(evt: dict) -> str: + return f"Role: {evt.get('role') or 'leaf'} Model: {evt.get('model') or '?'}" + + def _format_batch_delegation(evt: dict, deleg_id: str, completed_at: float) -> str: """Consolidated block for a delegate_task fan-out that finished as one unit.""" results = evt.get("results") or [] @@ -138,32 +124,24 @@ def _format_batch_delegation(evt: dict, deleg_id: str, completed_at: float) -> s n = len(results) if results else len(goals) total_dur = evt.get("total_duration_seconds", evt.get("duration_seconds", "?")) error = evt.get("error") - lines = [ + lines = _header_lines( + evt, f"[ASYNC DELEGATION BATCH COMPLETE — {deleg_id}]", f"A background fan-out of {n} subagent(s) you dispatched earlier " "has finished. All ran in parallel and waited on each other; their " "consolidated results are below. You may have moved on since " "dispatching — act on these or re-dispatch if things have changed.", - "", - ] - dispatched = _dispatched_line(evt.get("dispatched_at"), completed_at) - if dispatched: - lines.append(dispatched) + completed_at) lines.extend(_task_source_lines(evt)) - lines.append( - f"Role: {evt.get('role') or 'leaf'} Model: {evt.get('model') or '?'}" - f" Total duration: {total_dur}s" - ) + lines.append(f"{_role_model(evt)} Total duration: {total_dur}s") if error and not results: - lines.append("--- ERROR ---") - lines.append(f"The batch did not complete successfully: {error}") + lines += ["--- ERROR ---", f"The batch did not complete successfully: {error}"] return "\n".join(lines) # Config-level rejection notice BEFORE the per-task wall — a rejected # delegation model fails every task identically and must not stay buried. _notice = _delegation_model_not_found_notice(results) if _notice: - lines.append("") - lines.extend(_notice) + lines += ["", *_notice] for r in sorted(results, key=lambda x: x.get("task_index", 0)): idx = r.get("task_index", 0) r_status = r.get("status", "?") @@ -172,19 +150,14 @@ def _format_batch_delegation(evt: dict, deleg_id: str, completed_at: float) -> s r_goal = goals[idx] if idx < len(goals) else r.get("goal", "") r_truncated = _is_truncated(r) icon = "⚠" if r_truncated else ("✓" if r_status in ("completed", "success") else "✗") - lines.append("") - header = f"--- {icon} TASK {idx + 1}/{n}" - if r_goal: - header += f": {r_goal}" - header += f" (status={r_status}" + header = f"--- {icon} TASK {idx + 1}/{n}" + (f": {r_goal}" if r_goal else "") + f" (status={r_status}" if r.get("api_calls"): header += f", api_calls={r['api_calls']}" if r.get("duration_seconds") is not None: header += f", {r['duration_seconds']}s" if r_truncated: header += ", TRUNCATED: hit max_iterations — work may be incomplete" - header += ") ---" - lines.append(header) + lines += ["", header + ") ---"] if r_status in ("completed", "success") and r_summary: if r_truncated: lines.append(_TRUNCATED_SUMMARY_NOTE) @@ -192,102 +165,77 @@ def _format_batch_delegation(evt: dict, deleg_id: str, completed_at: float) -> s elif r_summary: if r_error: lines.append(f"({r_status}: {r_error})") - lines.append("Partial output:") - lines.append(r_summary) + lines += ["Partial output:", r_summary] else: - lines.append( - f"(no summary — status={r_status}" - + (f": {r_error}" if r_error else "") - + ")" - ) - r_live = r.get("live_transcript") - if r_live: - lines.append( - f"Full live transcript (complete tool/assistant trace): {r_live}" - ) + lines.append(f"(no summary — status={r_status}" + (f": {r_error}" if r_error else "") + ")") + if r.get("live_transcript"): + lines.append(f"Full live transcript (complete tool/assistant trace): {r['live_transcript']}") return "\n".join(lines) def _format_async_delegation(evt: dict) -> str: """Format an async-delegation completion into a self-contained re-injection. - - Carries the FULL original task source (goal, context, toolsets, role, - model) plus dispatch time, status, and the complete result summary: when - this re-enters the conversation the agent may be deep in unrelated context - and must be able to use the result OR re-dispatch without remembering why - the subagent existed. - """ + Carries the FULL original task source (goal, context, toolsets, role, model) plus + dispatch time, status, and the complete result summary: when this re-enters the + conversation the agent may be deep in unrelated context and must be able to use + the result OR re-dispatch without remembering why the subagent existed.""" deleg_id = evt.get("delegation_id", "unknown") completed_at = evt.get("completed_at") or time.time() if evt.get("is_batch") or isinstance(evt.get("results"), list): return _format_batch_delegation(evt, deleg_id, completed_at) - status = evt.get("status") or "completed" summary = evt.get("summary") error = evt.get("error") truncated = _is_truncated(evt) - lines = [ + lines = _header_lines( + evt, f"[ASYNC DELEGATION COMPLETE — {deleg_id}]", "A background subagent you dispatched earlier has finished. You may " "have moved on since dispatching it; the full task source is below so " "you can act on the result or re-dispatch if things have changed.", - "", - ] - dispatched = _dispatched_line(evt.get("dispatched_at"), completed_at) - if dispatched: - lines.append(dispatched) + completed_at) lines.append(f"Original goal: {evt.get('goal', '') or ''}") lines.extend(_task_source_lines(evt)) - lines.append(f"Role: {evt.get('role') or 'leaf'} Model: {evt.get('model') or '?'}") + lines.append(_role_model(evt)) _notice = _delegation_model_not_found_notice([evt]) if _notice: - lines.append("") - lines.extend(_notice) + lines += ["", *_notice] _trunc = " [TRUNCATED: hit max_iterations — work may be incomplete]" if truncated else "" - lines.append( + lines += [ f"Status: {status} API calls: {evt.get('api_calls', 0)} " - f"Duration: {evt.get('duration_seconds', '?')}s{_trunc}" - ) - lines.append("--- RESULT ---") + f"Duration: {evt.get('duration_seconds', '?')}s{_trunc}", + "--- RESULT ---", + ] if status in ("completed", "success") and summary: if truncated: lines.append(_TRUNCATED_SUMMARY_NOTE) lines.append(summary) else: if status == "interrupted": - lines.append( - "The subagent was interrupted before completing" - + (f": {error}" if error else ".") - ) + lines.append("The subagent was interrupted before completing" + (f": {error}" if error else ".")) else: # error / timeout / failed lines.append( - f"The subagent did not complete successfully (status={status})." - + (f"\n{error}" if error else "") + f"The subagent did not complete successfully (status={status})." + (f"\n{error}" if error else "") ) if summary: - lines.append("Partial output:") - lines.append(summary) + lines += ["Partial output:", summary] return "\n".join(lines) def _delegation_attribution_line(evt: dict) -> "str | None": """One-line provenance for a subagent-owned process event, else None. - - Subagents run terminal sessions under ``task_id == subagent_id``; a - background process they started outlives the child and is routed to the - PARENT conversation, which otherwise sees an anonymous raw output wall. + A background process a subagent started outlives the child and is routed to + the PARENT conversation, which otherwise sees an anonymous raw output wall. Judged on ``owner_task_id`` (the raw spawning id) — ``task_id`` is the - container key and may be collapsed to the session key. - """ + container key and may be collapsed to the session key.""" task_id = str(evt.get("owner_task_id") or evt.get("task_id") or "") if not task_id.startswith("sa-"): return None - try: + info = None + with suppress(Exception): from tools.delegate_tool import get_subagent_attribution info = get_subagent_attribution(task_id) - except Exception: - info = None if not info: # Registry entry aged out — still attribute generically, not anonymously. return f"Started by subagent {task_id} (delegate_task)." @@ -295,27 +243,20 @@ def _delegation_attribution_line(evt: dict) -> "str | None": if len(goal) > 120: goal = goal[:117] + "..." deleg = info.get("delegation_id") - parts = [f"Started by subagent {task_id}"] - if deleg: - parts.append(f"of delegation {deleg}") - line = " ".join(parts) + "." - if goal: - line += f' Task: "{goal}"' - return line + line = f"Started by subagent {task_id}" + (f" of delegation {deleg}" if deleg else "") + "." + return line + (f' Task: "{goal}"' if goal else "") + + +_REASON_STATUS = {"lost": "marked lost because the process backend disappeared", "failed_start": "failed to start"} def _completion_status(evt: dict) -> str: - _exit = evt.get("exit_code", "?") - _reason = evt.get("completion_reason") or "exited" - if _reason == "killed": + reason = evt.get("completion_reason") or "exited" + if reason == "killed": return f"terminated by {evt.get('termination_source') or 'Hermes'}" - if _reason == "lost": - return "marked lost because the process backend disappeared" - if _reason == "failed_start": - return "failed to start" - if _exit == 0: - return "completed normally" - return "exited" + if reason in _REASON_STATUS: + return _REASON_STATUS[reason] + return "completed normally" if evt.get("exit_code", "?") == 0 else "exited" def format_process_notification(evt: dict) -> "str | None": @@ -325,45 +266,33 @@ def format_process_notification(evt: dict) -> "str | None": _cmd = evt.get("command", "unknown") _attribution = _delegation_attribution_line(evt) - # watch_disabled and overflow events carry their own human-readable - # `message`; without this branch overflow events would fall through to the - # completion formatter as a phantom "process exited (exit code ?)". + # watch_disabled and overflow events carry their own human-readable `message`; + # otherwise overflow events would fall through to the completion formatter as a + # phantom "process exited (exit code ?)". if evt_type in ("watch_disabled", "watch_overflow_tripped", "watch_overflow_released"): return f"[IMPORTANT: {evt.get('message', '')}]" - + attribution = f"{_attribution}\n" if _attribution else "" if evt_type == "watch_match": _sup = evt.get("suppressed", 0) text = ( - f"[IMPORTANT: Background process {_sid} matched " - f"watch pattern \"{evt.get('pattern', '?')}\".\n" - ) - if _attribution: - text += f"{_attribution}\n" - text += f"Command: {_cmd}\nMatched output:\n{evt.get('output', '')}" + f"[IMPORTANT: Background process {_sid} matched watch pattern \"{evt.get('pattern', '?')}\".\n" + f"{attribution}Command: {_cmd}\nMatched output:\n{evt.get('output', '')}") if _sup: text += f"\n({_sup} earlier matches were suppressed by rate limit)" return text + "]" - if evt_type == "async_delegation": return _format_async_delegation(evt) _exit = evt.get("exit_code", "?") _out = evt.get("output", "") _signal = ", SIGTERM" if _exit in {-15, 143, "-15", "143"} else "" - text = ( - f"[IMPORTANT: Background process {_sid} {_completion_status(evt)} " - f"(exit code {_exit}{_signal}).\n" - ) - if _attribution: - text += f"{_attribution}\n" - # A subagent-owned process's full output belongs in the child's - # transcript, not as a raw wall in the parent — trim hard but keep - # enough tail to recognise failures. - if isinstance(_out, str) and len(_out) > 600: - _out = ( - "...(output trimmed — subagent-owned process; see the " - "delegation's live transcript for full output)\n" - + _out[-600:] - ) - text += f"Command: {_cmd}\nOutput:\n{_out}]" - return text + # A subagent-owned process's full output belongs in the child's transcript, not as + # a raw wall in the parent — trim hard but keep enough tail to recognise failures. + if _attribution and isinstance(_out, str) and len(_out) > 600: + _out = ( + "...(output trimmed — subagent-owned process; see the " + "delegation's live transcript for full output)\n" + + _out[-600:]) + return ( + f"[IMPORTANT: Background process {_sid} {_completion_status(evt)} (exit code {_exit}{_signal}).\n" + f"{attribution}Command: {_cmd}\nOutput:\n{_out}]")