From 50a13d53eb8792829adf20bb7457607a66c6dc69 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Sat, 5 Sep 2026 01:42:39 -0700 Subject: [PATCH] fix(goal): the judge can WAIT on delegated subagents (4 of 5 nudges in one run re-poked an agent that was waiting on workers) JUDGE_SYSTEM_PROMPT allowed WAIT only for "a background process listed below" or a stated backoff. A fan-out orchestrator's in-flight work is delegated subagents, which are not registry processes, so when the agent said "waiting on 4 worker batches, nothing to dispatch" the judge answered CONTINUE ("no background process is listed to gate on"). In the 1,393-agent run 4 of 5 nudges fired 6-153 s after exactly such a turn and each bought a status recap: 19 API calls at ~450K context, ~$4.9, and no work a batch-complete notification would not have produced. The judge now receives an "Active delegations: N batch(es) still running" line (count of live async_delegation records whose parent_session_id is this session, via the existing _session_records selector) and a WAIT branch for it: wait_for_seconds 600-1800 when the response says it is waiting on them with nothing dispatchable. Both goal-loop callers (CLI, gateway) pass the count. Live A/B with the real auxiliary judge on the run's response shape ("4 worker batches still running ... nothing else is dispatchable until these return"): main 3/3 CONTINUE, branch 3/3 WAIT 1200 s. Tests (2): the prompt carries the delegations line only when N > 0 and a wait verdict parses through; count_active_delegations is scoped to the spawning session and ignores completed records. --- gateway/run_goals.py | 6 ++- hermes_cli/cli_loops_mixin.py | 6 ++- hermes_cli/goals.py | 29 ++++++++++- .../hermes_cli/test_goal_judge_delegations.py | 48 +++++++++++++++++++ 4 files changed, 83 insertions(+), 6 deletions(-) create mode 100644 tests/hermes_cli/test_goal_judge_delegations.py diff --git a/gateway/run_goals.py b/gateway/run_goals.py index 011fc12c11..2fe88d27d9 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -241,12 +241,13 @@ class GatewayGoalsMixin: if mgr is None or not mgr.is_active(): return - _bg_procs = None + _bg_procs, _active_deleg = None, 0 with suppress(Exception): - from hermes_cli.goals import gather_background_processes as _gather_bg + from hermes_cli.goals import count_active_delegations, gather_background_processes as _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) + _active_deleg = count_active_delegations(getattr(session_entry, "session_id", 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 @@ -254,6 +255,7 @@ class GatewayGoalsMixin: decision = await self._run_in_executor_with_context( lambda: mgr.evaluate_after_turn( final_response or "", user_initiated=True, background_processes=_bg_procs, + active_delegations=_active_deleg, ), ) msg = decision.get("message") or "" diff --git a/hermes_cli/cli_loops_mixin.py b/hermes_cli/cli_loops_mixin.py index 1b7c3ce1b7..1af5ef9720 100644 --- a/hermes_cli/cli_loops_mixin.py +++ b/hermes_cli/cli_loops_mixin.py @@ -560,14 +560,16 @@ class CLILoopsMixin: last_response = self._last_assistant_response_text() if not last_response.strip(): return + _active_deleg = 0 try: - from hermes_cli.goals import gather_background_processes as _gather_bg + from hermes_cli.goals import count_active_delegations, gather_background_processes as _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) + _active_deleg = count_active_delegations(getattr(self.agent, "session_id", None)) except Exception: _bg_procs = None decision = mgr.evaluate_after_turn( - last_response, user_initiated=True, background_processes=_bg_procs) + last_response, user_initiated=True, background_processes=_bg_procs, active_delegations=_active_deleg) _print_decision_message(decision) if decision.get("should_continue"): prompt = decision.get("continuation_prompt") diff --git a/hermes_cli/goals.py b/hermes_cli/goals.py index d730230712..ec5d2c6e5e 100644 --- a/hermes_cli/goals.py +++ b/hermes_cli/goals.py @@ -141,6 +141,11 @@ JUDGE_SYSTEM_PROMPT = ( "``wait_on_pid`` (releases on exit only).\n" "- The agent says it is rate-limited / backing off / must wait a fixed " "period — return seconds in ``wait_for_seconds``.\n" + "- The agent has delegated subagents still running (stated below as " + "active delegations) and the response says it is waiting on them with " + "nothing else dispatchable — return ``wait_for_seconds`` between 600 and " + "1800. Their results wake the agent on their own; re-poking it now only " + "produces a status recap.\n" "Picking WAIT parks the loop without burning a turn; it resumes " "automatically when the pid exits or the time elapses. Do NOT pick WAIT " "just because work remains — only when re-poking now would be pure " @@ -159,6 +164,12 @@ JUDGE_SYSTEM_PROMPT = ( "accepted (true=done, false=continue)." ) +# Judge prompt line for live delegated subagents (WAIT-for-seconds vs CONTINUE). +JUDGE_DELEGATIONS_BLOCK_TEMPLATE = ( + "Active delegations: the agent has {count} delegated subagent batch(es) still running; " + "their results are delivered to it automatically when they finish.\n\n" +) + # Judge prompt block listing running background processes (WAIT vs CONTINUE, which pid). JUDGE_BACKGROUND_BLOCK_TEMPLATE = ( "Background processes the agent currently has running (it may be waiting " @@ -860,6 +871,7 @@ def judge_goal( subgoals: Optional[List[str]] = None, background_processes: Optional[List[Dict[str, Any]]] = None, contract: Optional[GoalContract] = None, + active_delegations: int = 0, ) -> Tuple[str, str, bool, Optional[Dict[str, Any]], bool]: """Ask the auxiliary model whether the goal is satisfied. @@ -886,7 +898,8 @@ def judge_goal( common = dict( goal=_truncate(goal, 2000), response=_truncate(last_response, _JUDGE_RESPONSE_SNIPPET_CHARS), - background_block=_render_background_block(background_processes), + background_block=_render_background_block(background_processes) + + (JUDGE_DELEGATIONS_BLOCK_TEMPLATE.format(count=active_delegations) if active_delegations > 0 else ""), current_time=datetime.now(tz=timezone.utc).astimezone().strftime("%Y-%m-%d %H:%M:%S %Z"), ) if contract is not None and not contract.is_empty(): @@ -912,6 +925,17 @@ def judge_goal( return verdict, reason, parse_failed, wait_directive, False +def count_active_delegations(session_id: Optional[str]) -> int: + """Live async delegation batches spawned by this session (fail-safe 0).""" + if not session_id: + return 0 + try: + from tools.async_delegation import _LIVE_STATES, _session_records + return len(_session_records(_LIVE_STATES, "", "", str(session_id))) + except Exception: + return 0 + + 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. @@ -1351,6 +1375,7 @@ class GoalManager: def evaluate_after_turn( self, last_response: str, *, user_initiated: bool = True, background_processes: Optional[List[Dict[str, Any]]] = None, + active_delegations: int = 0, ) -> Dict[str, Any]: """Run gates + judge and update state. Return a decision dict (``status``, ``should_continue``, ``continuation_prompt``, ``verdict``, ``reason``, ``message``). Both real user prompts and our @@ -1376,7 +1401,7 @@ class GoalManager: verdict, reason, parse_failed, wait_directive, transport_failed = judge_goal( state.goal, last_response, subgoals=state.subgoals or None, background_processes=background_processes, - contract=state.contract if state.has_contract() else None, + contract=state.contract if state.has_contract() else None, active_delegations=active_delegations, ) state.last_verdict = verdict state.last_reason = reason diff --git a/tests/hermes_cli/test_goal_judge_delegations.py b/tests/hermes_cli/test_goal_judge_delegations.py new file mode 100644 index 0000000000..96275141b2 --- /dev/null +++ b/tests/hermes_cli/test_goal_judge_delegations.py @@ -0,0 +1,48 @@ +"""The goal judge knows about delegated subagents. + +In a fan-out run 4 of 5 /goal nudges fired 6-153 s after a turn that had said "waiting on workers, +nothing to dispatch": the judge prompt had no WAIT branch for delegated subagents (only for registry +processes), so it returned CONTINUE and each nudge bought a status recap (19 API calls, ~$4.9). +""" +from types import SimpleNamespace +from unittest.mock import patch + +from hermes_cli import goals + + +def _resp(content): + return SimpleNamespace(choices=[SimpleNamespace(message=SimpleNamespace(content=content))]) + + +def test_judge_prompt_states_active_delegations_and_a_wait_branch_for_them(): + seen = {} + + def fake_call_llm(*a, **kw): + seen["prompt"] = kw.get("messages") or a + return _resp('{"verdict": "wait", "wait_for_seconds": 900, "reason": "waiting on workers"}') + + with patch("agent.auxiliary_client.call_llm", side_effect=fake_call_llm): + verdict, _reason, parse_failed, directive, _transport = goals.judge_goal( + "refactor everything", "Waiting on 4 workers; nothing to dispatch.", active_delegations=4) + text = str(seen["prompt"]) + assert "Active delegations: the agent has 4 delegated subagent batch(es) still running" in text + assert "delegated subagents still running" in goals.JUDGE_SYSTEM_PROMPT + assert (verdict, parse_failed) == ("wait", False) and directive.get("seconds") == 900 + + with patch("agent.auxiliary_client.call_llm", side_effect=fake_call_llm): + goals.judge_goal("g", "r", active_delegations=0) + assert "Active delegations" not in str(seen["prompt"]) + + +def test_count_active_delegations_is_scoped_to_the_spawning_session(): + from tools import async_delegation as ad + + fake = { + "a": {"status": "running", "parent_session_id": "root", "session_key": "", "origin_ui_session_id": ""}, + "b": {"status": "completed", "parent_session_id": "root", "session_key": "", "origin_ui_session_id": ""}, + "c": {"status": "running", "parent_session_id": "other", "session_key": "", "origin_ui_session_id": ""}, + } + with patch.object(ad, "_records", fake): + assert goals.count_active_delegations("root") == 1 + assert goals.count_active_delegations("other") == 1 + assert goals.count_active_delegations(None) == 0