diff --git a/cli.py b/cli.py index 9fe19f1296..14bde0d670 100644 --- a/cli.py +++ b/cli.py @@ -3469,6 +3469,7 @@ class HermesCLI(CLIAgentSetupMixin, CLICommandsMixin, CLIBillingMixin, CLITuiMix self._check_termios_drift, lambda: self._drain_process_notifications("cli-idle"), self._maybe_fire_loop_tick, + self._maybe_resume_parked_goal, ): with suppress(Exception): step() diff --git a/gateway/run_goals.py b/gateway/run_goals.py index ef2bd26d41..011fc12c11 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -244,7 +244,9 @@ class GatewayGoalsMixin: _bg_procs = None with suppress(Exception): from hermes_cli.goals import gather_background_processes as _gather_bg - _bg_procs = _gather_bg() + # Only THIS session's processes (gateway turns register under turn_ctx.session_id): + # subagents' pollers must not park the parent's goal. + _bg_procs = _gather_bg(owner_task_id=getattr(session_entry, "session_id", None) or None) # judge_goal() is a synchronous aux-LLM HTTP call (10-40 s; would block Discord heartbeats). # _run_in_executor_with_context carries the profile secret scope / aux runtime contextvars diff --git a/hermes_cli/cli_loops_mixin.py b/hermes_cli/cli_loops_mixin.py index fa00827bff..1b7c3ce1b7 100644 --- a/hermes_cli/cli_loops_mixin.py +++ b/hermes_cli/cli_loops_mixin.py @@ -386,6 +386,35 @@ class CLILoopsMixin: self._heartbeat_watchdog_started = False threading.Thread(target=_loop, daemon=True, name="heartbeat-watchdog").start() + def _maybe_resume_parked_goal(self) -> None: + """Idle hook run from process_loop: when a parked /goal's barrier has lifted (the process + exited, the timer elapsed, or the wait aged past its cap), queue the continuation so the + loop resumes WITHOUT waiting for an unrelated turn to re-evaluate it. The barrier used to be + checked only lazily, on the next turn; a session with nothing else arriving stayed parked + indefinitely (one run: 3 h 22 min on a grandchild's poller).""" + now = time.time() + if now - getattr(self, "_last_goal_barrier_check", 0.0) < 5.0: + return + self._last_goal_barrier_check = now + try: + if not self._pending_input.empty(): + return + mgr = self._get_goal_manager() + state = getattr(mgr, "state", None) if mgr is not None else None + if state is None or state.status != "active": + return + if not (state.waiting_on_pid is not None or state.waiting_on_session is not None or state.waiting_until): + return # not parked + if mgr.is_waiting(): + return # barrier still holds (is_waiting() also applies the age cap) + prompt = mgr.next_continuation_prompt() + if prompt: + from cli import _DIM, _RST, _cprint + _cprint(f" {_DIM}▶ Goal barrier lifted — resuming.{_RST}") + self._pending_input.put(prompt) + except Exception as exc: + logging.debug("parked-goal resume check failed: %s", exc) + def _maybe_fire_loop_tick(self) -> None: """Idle hook run from process_loop: fire a due /loop wakeup. @@ -533,7 +562,8 @@ class CLILoopsMixin: return try: from hermes_cli.goals import gather_background_processes as _gather_bg - _bg_procs = _gather_bg() + # Only THIS session's processes: subagents' pollers must not park the parent's goal. + _bg_procs = _gather_bg(owner_task_id=getattr(self, "session_id", None) or None) except Exception: _bg_procs = None decision = mgr.evaluate_after_turn( diff --git a/hermes_cli/goals.py b/hermes_cli/goals.py index 8d0b24f679..d730230712 100644 --- a/hermes_cli/goals.py +++ b/hermes_cli/goals.py @@ -49,6 +49,9 @@ DEFAULT_MAX_CONSECUTIVE_TRANSPORT_FAILURES = 5 # on concrete evidence instead of a vibe check. DEFAULT_GATE_TIMEOUT_SECONDS = 300 DEFAULT_GATE_MAX_RETRIES = 3 +# Longest a pid/session wait barrier may hold the loop before judging resumes. Timed barriers +# (``waiting_until``) carry their own deadline and are exempt. +_MAX_BARRIER_WAIT_S = 30 * 60 # Bounded tail of a failed gate's combined stdout/stderr fed back to the agent. _GATE_OUTPUT_TAIL_CHARS = 3000 @@ -909,9 +912,15 @@ def judge_goal( return verdict, reason, parse_failed, wait_directive, False -def gather_background_processes(task_id: Optional[str] = None) -> List[Dict[str, Any]]: +def gather_background_processes(task_id: Optional[str] = None, *, owner_task_id: Optional[str] = None) -> List[Dict[str, Any]]: """Fail-safe snapshot of RUNNING ``process_registry`` sessions for the judge; ``[]`` on any error - so the loop degrades to its pre-wait-barrier behavior.""" + so the loop degrades to its pre-wait-barrier behavior. + + ``owner_task_id`` restricts the snapshot to processes the goal's OWN session spawned. The registry's + ``task_id`` is the container key, which collapses to one value for every agent in the process, so + without this filter a fan-out parent's judge saw every subagent's pollers and parked the goal on a + grandchild's ``proc_*`` session (one run: 7 of 7 root verdicts were WAIT on child-owned processes; + parked 3 h 22 min at the end while nothing of its own was running).""" try: from tools.process_registry import process_registry @@ -919,7 +928,10 @@ def gather_background_processes(task_id: Optional[str] = None) -> List[Dict[str, except Exception as exc: logger.debug("gather_background_processes failed: %s", exc) return [] - return [s for s in sessions if isinstance(s, dict) and s.get("status") != "exited"] + running = [s for s in sessions if isinstance(s, dict) and s.get("status") != "exited"] + if owner_task_id: + running = [s for s in running if str(s.get("owner_task_id") or s.get("task_id") or "") == str(owner_task_id)] + return running def draft_contract(objective: str, *, timeout: Optional[float] = None) -> Optional[GoalContract]: @@ -1282,7 +1294,10 @@ class GoalManager: def is_waiting(self) -> bool: """True iff a barrier is set AND not yet satisfied. A satisfied barrier is cleared here - (lazy auto-clear) so the next evaluation resumes normal judging.""" + (lazy auto-clear) so the next evaluation resumes normal judging. A pid/session barrier + also expires after ``_MAX_BARRIER_WAIT_S``: a watcher or poller that never exits would + otherwise park the goal indefinitely (one run sat 3 h 22 min on a poller that outlived + the work it was polling).""" s = self._state if s is None: return False @@ -1294,6 +1309,10 @@ class GoalManager: still = time.time() < s.waiting_until else: return False + if still and s.waiting_since and s.waiting_until == 0.0 and time.time() - s.waiting_since > _MAX_BARRIER_WAIT_S: + logger.info("goal %s: wait barrier on %s exceeded %ds; resuming judging", + self.session_id, s.waiting_on_session or s.waiting_on_pid, _MAX_BARRIER_WAIT_S) + still = False if not still: self.stop_waiting() return still diff --git a/tests/cli/test_cli_goal_parked_resume.py b/tests/cli/test_cli_goal_parked_resume.py new file mode 100644 index 0000000000..a5a729f342 --- /dev/null +++ b/tests/cli/test_cli_goal_parked_resume.py @@ -0,0 +1,58 @@ +"""A parked /goal resumes from the idle hook once its barrier lifts, without waiting for another turn.""" +import queue +import time +from unittest.mock import patch + +import pytest + +from hermes_cli import goals +from hermes_cli.cli_loops_mixin import CLILoopsMixin + + +@pytest.fixture +def hermes_home(tmp_path, monkeypatch): + from pathlib import Path + home = tmp_path / ".hermes"; home.mkdir() + monkeypatch.setattr(Path, "home", lambda: tmp_path) + monkeypatch.setenv("HERMES_HOME", str(home)) + goals._DB_CACHE.clear() + yield home + goals._DB_CACHE.clear() + + +class _Cli(CLILoopsMixin): + def __init__(self, mgr): + self._pending_input = queue.Queue() + self._mgr = mgr + + def _get_goal_manager(self): + return self._mgr + + +def test_idle_hook_queues_the_continuation_when_a_timed_barrier_has_elapsed(hermes_home): + mgr = goals.GoalManager(session_id="resume-idle") + mgr.set("finish the thing") + mgr.wait_for_seconds(1, reason="cooldown") + cli = _Cli(mgr) + with patch("cli._cprint"), patch("cli._DIM", ""), patch("cli._RST", ""): + cli._maybe_resume_parked_goal() + assert cli._pending_input.empty() # still parked + mgr.state.waiting_until = time.time() - 1 + mgr._save() + cli._last_goal_barrier_check = 0.0 + cli._maybe_resume_parked_goal() + assert not cli._pending_input.empty() # continuation queued + assert "finish the thing" in cli._pending_input.get() + assert mgr.state.waiting_until == 0.0 # barrier cleared + + +def test_idle_hook_is_a_no_op_for_an_unparked_or_inactive_goal(hermes_home): + mgr = goals.GoalManager(session_id="resume-noop") + mgr.set("g") + cli = _Cli(mgr) + cli._maybe_resume_parked_goal() + assert cli._pending_input.empty() + mgr.clear() + cli._last_goal_barrier_check = 0.0 + cli._maybe_resume_parked_goal() + assert cli._pending_input.empty() diff --git a/tests/hermes_cli/test_goals.py b/tests/hermes_cli/test_goals.py index 413a9330ed..f3b2ea7439 100644 --- a/tests/hermes_cli/test_goals.py +++ b/tests/hermes_cli/test_goals.py @@ -436,6 +436,46 @@ class TestWaitBarrier: proc.terminate() proc.wait(timeout=10) + def test_barrier_on_a_process_that_never_exits_expires(self, hermes_home): + """A poller that outlives the work parked one run for 3h22m; a live barrier ages out.""" + from hermes_cli import goals + from hermes_cli.goals import GoalManager + + proc = self._spawn_sleeper() + try: + mgr = GoalManager(session_id="wb-expire") + mgr.set("g") + mgr.wait_on(proc.pid, reason="poller") + assert mgr.is_waiting() is True + mgr.state.waiting_since = time.time() - goals._MAX_BARRIER_WAIT_S - 1 + mgr._save() + assert mgr.is_waiting() is False + assert mgr.state.waiting_on_pid is None + finally: + proc.terminate() + proc.wait(timeout=10) + + +class TestGatherBackgroundProcessesOwnership: + def test_only_the_owning_sessions_processes_are_seen(self, monkeypatch): + """The registry task_id collapses to one container key for every agent in the process, so a + fan-out parent's judge must filter by owner; otherwise a grandchild's poller parks the goal.""" + from hermes_cli import goals + + class _Reg: + def list_sessions(self, task_id=None, session_key=None): + return [ + {"session_id": "proc_mine", "status": "running", "owner_task_id": "root-sid", "task_id": "default"}, + {"session_id": "proc_child", "status": "running", "owner_task_id": "sa-1-abc", "task_id": "default"}, + {"session_id": "proc_done", "status": "exited", "owner_task_id": "root-sid", "task_id": "default"}, + ] + + import tools.process_registry as pr + monkeypatch.setattr(pr, "process_registry", _Reg()) + assert [p["session_id"] for p in goals.gather_background_processes()] == ["proc_mine", "proc_child"] + assert [p["session_id"] for p in goals.gather_background_processes(owner_task_id="root-sid")] == ["proc_mine"] + + # ────────────────────────────────────────────────────────────────────── # Judge-driven auto-wait — the judge parks the loop on its own diff --git a/tools/process_registry.py b/tools/process_registry.py index 2df6fe1eba..7a0fe442e3 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -1795,6 +1795,7 @@ class ProcessRegistry: "command": s.command[:200], "cwd": s.cwd, "pid": s.pid, + "owner_task_id": s.owner_task_id or s.task_id, "started_at": time.strftime("%Y-%m-%dT%H:%M:%S", time.localtime(s.started_at)), "uptime_seconds": int(time.time() - s.started_at), "status": "exited" if s.exited else "running", diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index 0146955742..f074c19dbf 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -294,7 +294,9 @@ def _goal_followup_after_turn( if session.get("session_key") and (goal_mgr := _active_goal_manager(session)) is not None: try: from hermes_cli.goals import gather_background_processes as _gather_bg - _bg_procs = _gather_bg() + # Only THIS session's processes (TUI turns register under session_key): subagents' + # pollers must not park the parent's goal. Same rule as the CLI and gateway loops. + _bg_procs = _gather_bg(owner_task_id=session.get("session_key") or None) except Exception: _bg_procs = None decision = goal_mgr.evaluate_after_turn( diff --git a/website/docs/user-guide/features/goals.md b/website/docs/user-guide/features/goals.md index d187dd9bef..40949e6493 100644 --- a/website/docs/user-guide/features/goals.md +++ b/website/docs/user-guide/features/goals.md @@ -152,7 +152,7 @@ Gates and contracts compose: use a contract to shape *what the agent aims for*, Some goals are gated on something that takes minutes and runs on its own — CI on a pushed PR, a long build, a test matrix, a deploy, a rate-limit cooldown. Without help, the goal loop would re-poke the agent every turn into "is it done yet?" busy-work while it waits. -**This is handled automatically.** Every turn, the judge is shown the agent's live background processes (the `terminal(background=true)` registry — pid, session id, command, uptime, recent output, and any `watch_patterns` / `notify_on_complete` trigger) alongside the goal and the agent's response. When the agent's progress is genuinely gated on one of them, the judge returns a **`wait`** verdict instead of `continue`, and the loop **parks**: the next turns are skipped (no judge call, no continuation, no turn consumed) until the wait is satisfied — then it resumes normally with the result in hand. The judge can also park on a **time** basis (`wait_for_seconds`) for backoff/cooldown waits. `/goal status` shows `⏳ Goal (parked …)` while parked. +**This is handled automatically.** Every turn, the judge is shown the agent's own live background processes (the `terminal(background=true)` registry entries this session spawned — pid, session id, command, uptime, recent output, and any `watch_patterns` / `notify_on_complete` trigger; processes started by delegated subagents are not shown, so a fan-out parent is never parked on a worker's poller) alongside the goal and the agent's response. When the agent's progress is genuinely gated on one of them, the judge returns a **`wait`** verdict instead of `continue`, and the loop **parks**: the next turns are skipped (no judge call, no continuation, no turn consumed) until the wait is satisfied — then it resumes normally with the result in hand. A pid/session wait is capped at 30 minutes; a process that never exits (a watcher, a forgotten poller) cannot park the goal indefinitely. The judge can also park on a **time** basis (`wait_for_seconds`) for backoff/cooldown waits. `/goal status` shows `⏳ Goal (parked …)` while parked. The judge picks the right kind of wait from the process's own signal: