Merge pull request #103496 from NousResearch/fix/goal-judge-own-processes-only
fix(goal): the judge sees only its own session's background processes; pid/session waits expire after 30 min (parked 3h22m on a grandchild's poller)
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
+23
-4
@@ -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
|
||||
|
||||
@@ -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()
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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:
|
||||
|
||||
|
||||
Reference in New Issue
Block a user