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..0f1d870209 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 " @@ -428,6 +439,9 @@ class GoalState: waiting_on_pid: Optional[int] = None waiting_on_session: Optional[str] = None waiting_until: float = 0.0 + # Live delegation batches when a timed WAIT was set because of them; the barrier lifts as soon + # as that count drops (a batch returned), not only when the timer runs out. + waiting_on_delegations: int = 0 waiting_reason: Optional[str] = None waiting_since: float = 0.0 contract: GoalContract = field(default_factory=GoalContract) @@ -441,7 +455,7 @@ class GoalState: def from_json(cls, raw: str) -> "GoalState": data = json.loads(raw) raw_subgoals = data.get("subgoals") or [] - ints = {k: int(data.get(k) or 0) for k in ("turns_used", "consecutive_parse_failures", "consecutive_transport_failures")} + ints = {k: int(data.get(k) or 0) for k in ("turns_used", "consecutive_parse_failures", "consecutive_transport_failures", "waiting_on_delegations")} floats = {k: float(data.get(k) or 0.0) for k in ("created_at", "last_turn_at", "waiting_until", "waiting_since")} return cls( goal=data.get("goal", ""), @@ -473,6 +487,7 @@ class GoalState: self.waiting_on_pid = None self.waiting_on_session = None self.waiting_until = 0.0 + self.waiting_on_delegations = 0 self.waiting_reason = None self.waiting_since = 0.0 @@ -860,6 +875,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 +902,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 +929,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. @@ -1275,13 +1303,17 @@ class GoalManager: raise ValueError("session_id must be a non-empty string") return self._park(reason, waiting_on_session=session_id) - def wait_for_seconds(self, seconds: int, reason: str = "") -> GoalState: - """Park until ``seconds`` from now (backoff/cooldown waits with no process to track).""" + def wait_for_seconds(self, seconds: int, reason: str = "", *, on_delegations: int = 0) -> GoalState: + """Park until ``seconds`` from now (backoff/cooldown waits with no process to track). With + ``on_delegations`` the wait is FOR those live delegation batches: it also lifts as soon as + fewer are live (a batch result came back), so the loop re-judges with the result in hand + instead of sleeping out a 20-minute timer (independent review: results arrived with 1,199 s + left on the timer and nothing re-judged).""" self._require_active() seconds = int(seconds) if seconds <= 0: raise ValueError("seconds must be a positive integer") - return self._park(reason, waiting_until=time.time() + seconds) + return self._park(reason, waiting_until=time.time() + seconds, waiting_on_delegations=max(0, int(on_delegations))) def stop_waiting(self) -> bool: """Clear any active wait barrier (pid / session / time). Returns True if one was cleared.""" @@ -1307,6 +1339,11 @@ class GoalManager: still = _pid_alive(s.waiting_on_pid) elif s.waiting_until: still = time.time() < s.waiting_until + if still and s.waiting_on_delegations > 0: + # Set because of live delegations: lift the moment one of them returned. + live = count_active_delegations(self.session_id) + if live < s.waiting_on_delegations: + still = False 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: @@ -1329,7 +1366,7 @@ class GoalManager: reason = state.waiting_reason or tgt return _decision("active", False, None, "waiting", reason, f"⏳ Goal parked — waiting on {tgt}: {reason}") - def _apply_wait_directive(self, wait_directive: Dict[str, Any], reason: str) -> Dict[str, Any]: + def _apply_wait_directive(self, wait_directive: Dict[str, Any], reason: str, *, active_delegations: int = 0) -> Dict[str, Any]: """Judge said WAIT: set the barrier and park. The counted turn stands (the judge ran) but no continuation fires; the loop resumes once the barrier clears.""" if wait_directive.get("session_id"): @@ -1337,7 +1374,7 @@ class GoalManager: elif wait_directive.get("pid"): tgt = f"pid {self.wait_on(int(wait_directive['pid']), reason=reason).waiting_on_pid}" else: - self.wait_for_seconds(int(wait_directive["seconds"]), reason=reason) + self.wait_for_seconds(int(wait_directive["seconds"]), reason=reason, on_delegations=active_delegations) tgt = f"{wait_directive['seconds']}s" return _decision("active", False, None, "wait", reason, f"⏳ Goal parked (judge) — waiting on {tgt}: {reason}") @@ -1351,6 +1388,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 +1414,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 @@ -1387,7 +1425,7 @@ class GoalManager: state.consecutive_transport_failures = state.consecutive_transport_failures + 1 if transport_failed else 0 if verdict == "wait" and wait_directive: - return self._apply_wait_directive(wait_directive, reason) + return self._apply_wait_directive(wait_directive, reason, active_delegations=active_delegations) # BLOCKED is NOT done: pause so the user sees the judge's reason and can re-scope or override, # instead of burning turns on an unachievable goal or waving it through as complete. 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..28d38c3082 --- /dev/null +++ b/tests/hermes_cli/test_goal_judge_delegations.py @@ -0,0 +1,70 @@ +"""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 + + +def test_a_delegation_wait_lifts_when_a_batch_returns_not_only_when_the_timer_runs_out(tmp_path, monkeypatch): + """Independent-review witness: results arrived with 1,199 s left on a 1,200 s WAIT and nothing + re-judged; integration sat unfinished until the timer ran out.""" + from pathlib import Path + monkeypatch.setattr(Path, "home", lambda: tmp_path) + monkeypatch.setenv("HERMES_HOME", str(tmp_path / ".hermes")); (tmp_path / ".hermes").mkdir() + goals._DB_CACHE.clear() + mgr = goals.GoalManager(session_id="root-wait") + mgr.set("integrate the rounds") + with patch.object(goals, "count_active_delegations", return_value=4): + mgr.wait_for_seconds(1200, reason="4 batches running", on_delegations=4) + assert mgr.is_waiting() is True # all four still live: parked + with patch.object(goals, "count_active_delegations", return_value=3): + assert mgr.is_waiting() is False # one returned: barrier lifted early + assert mgr.state.waiting_until == 0.0 and mgr.state.waiting_on_delegations == 0 + # a plain timed wait (no delegations) is unaffected by the delegation count + mgr.wait_for_seconds(1200, reason="cooldown") + with patch.object(goals, "count_active_delegations", return_value=0): + assert mgr.is_waiting() is True + goals._DB_CACHE.clear() diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index f074c19dbf..f89bbe8855 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -292,15 +292,17 @@ def _goal_followup_after_turn( return goal_followup try: if session.get("session_key") and (goal_mgr := _active_goal_manager(session)) is not None: + _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 (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) + _active_deleg = count_active_delegations(getattr(session.get("agent"), "session_id", None)) except Exception: _bg_procs = None decision = goal_mgr.evaluate_after_turn( - raw, user_initiated=True, background_processes=_bg_procs) + raw, user_initiated=True, background_processes=_bg_procs, active_delegations=_active_deleg) if verdict_msg := decision.get("message") or "": _emit("status.update", sid, {"kind": "goal", "text": verdict_msg}) if decision.get("should_continue") and (