fix(goal): the judge only sees the goal's own background processes, and a pid/session wait barrier expires after 30 min
In a fan-out run the /goal loop parked for 3 h 22 min at the end on a grandchild's poller (waiting_on_session=proc_21a6fe2369a1, 00:47 -> 04:09) while the root itself had nothing running. Every one of the root's 7 logged judge verdicts was WAIT on a child-owned process. Two causes. gather_background_processes() was called with no task_id from both goal-loop callers (CLI cli_loops_mixin, gateway run_goals), and the registry's task_id is the CONTAINER key, which collapses to one value for every agent in the process, so the judge's process list was every one of ~1,300 subagents' pollers. And a pid/session barrier had no ceiling: once the judge said WAIT on a session that never exits, nothing resumed judging; waiting_since was recorded and never read. Now list_sessions() reports owner_task_id (the RAW spawning id the registry already keeps for ownership checks), gather_background_processes takes owner_task_id and both callers pass their own session id (CLI turns register processes under self.session_id; gateway turns under turn_ctx.session_id), and is_waiting() ages out a pid/session barrier after _MAX_BARRIER_WAIT_S (30 min). Timed barriers keep their own deadline. Tests: only the owning session's running processes are returned when owner_task_id is given (unfiltered behaviour unchanged); a live barrier older than the ceiling clears and judging resumes.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
+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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user