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..4622a60561 100644 --- a/hermes_cli/cli_loops_mixin.py +++ b/hermes_cli/cli_loops_mixin.py @@ -533,7 +533,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/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 75ae003b84..bc4ed7a419 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -1776,6 +1776,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",