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.
This commit is contained in:
Teknium
2026-09-05 01:42:39 -07:00
parent 1f3912d277
commit 50a13d53eb
4 changed files with 83 additions and 6 deletions
+4 -2
View File
@@ -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 ""
+4 -2
View File
@@ -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")
+27 -2
View File
@@ -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
@@ -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