diff --git a/cli.py b/cli.py index 0d904afe04..50f1fc0a90 100644 --- a/cli.py +++ b/cli.py @@ -959,10 +959,16 @@ def _wait_for_oneshot_background_completions(cli) -> None: Waits on the whole registry: a one-shot process hosts one agent, and task_id filtering would skip processes registered before the session id settled. + Skipped when the quiet -Q notify-resume loop already consumed the run's linger + budget: it calls wait_for_pending_completions with a shared deadline, so a + re-wait here would double-block on the same stuck notify_on_complete child. + See #90879. """ from tools.process_registry import process_registry + if getattr(cli, "_quiet_notify_linger_done", False): + return _agent, task_id = _oneshot_agent_and_session(cli) result = process_registry.wait_for_pending_completions(None) if result.get("waited"): @@ -4079,7 +4085,9 @@ def _run_quiet_single_query(cli, effective_query): stranded receipt.""" from agent.interrupt_compat import _accepts_keyword from agent.turn_author import take_turn_author_from_env - from hermes_cli.quiet_single_query import bind_quiet_session_key, continue_quiet_notify_completions + from hermes_cli.quiet_single_query import ( + bind_quiet_session_key, continue_quiet_notify_completions, quiet_notify_linger_seconds, + ) author = take_turn_author_from_env() author_kwargs = {"turn_author": author} if author is not None and _accepts_keyword(cli.agent.run_conversation, "turn_author") else {} @@ -4112,10 +4120,15 @@ def _run_quiet_single_query(cli, effective_query): history = follow["messages"] return follow + # One shared linger budget for the whole run: the loop below and the later + # _wait_for_oneshot_background_completions pass must not each wait the full + # oneshot_completion_wait_seconds on the same stuck notify_on_complete child. + cli._quiet_notify_linger_done = True continued = continue_quiet_notify_completions( getattr(cli, "session_id", "") or "", _follow_up, owns_event=getattr(cli, "_owns_process_notification", None), + linger_budget=quiet_notify_linger_seconds(), ) if isinstance(continued, dict): result = continued diff --git a/hermes_cli/quiet_single_query.py b/hermes_cli/quiet_single_query.py index 9efd9113f7..c1508a9938 100644 --- a/hermes_cli/quiet_single_query.py +++ b/hermes_cli/quiet_single_query.py @@ -10,6 +10,7 @@ was never injected as a follow-up. from __future__ import annotations +import time from typing import Any, Callable # Nested A→B→C is one extra turn; this caps a runaway message_agent chain. @@ -24,29 +25,52 @@ def bind_quiet_session_key(session_id: str): return token, reset_current_session_key +def quiet_notify_linger_seconds() -> float: + """Total linger budget for one quiet run: the shared ``terminal.oneshot_completion_wait_seconds``. + + One budget covers the drain loop here AND the later ``_wait_for_oneshot_background_completions`` + pass, so a stuck ``notify_on_complete`` child cannot stack round-after-round waits on top of the + finalize re-wait (pre-fix worst case: 8 rounds x 600s + 600s). + """ + from tools.process_registry import ProcessRegistry + + return ProcessRegistry._oneshot_completion_wait_seconds() + + def continue_quiet_notify_completions( session_id: str, run_turn: Callable[[str], Any], *, owns_event=None, max_rounds: int = _MAX_QUIET_NOTIFY_ROUNDS, + linger_budget: float | None = None, ) -> Any: """Linger for ``notify_on_complete`` work, then run owned completion texts as follow-up turns. - Returns the last ``run_turn`` result, or ``None`` when nothing owned completed. + Returns the last ``run_turn`` result, or ``None`` when nothing owned completed. The whole + loop shares ONE linger budget (default: ``terminal.oneshot_completion_wait_seconds``) — a + process that times out is waited on no further this run: after the current round's drained + texts run, the loop stops (the finalize linger still covers it once, bounded, via the + budget handshake below). """ from tools.process_registry import process_registry last: Any = None key = session_id or "" - for _ in range(max(int(max_rounds), 0)): - process_registry.wait_for_pending_completions(None) + if linger_budget is None: + linger_budget = quiet_notify_linger_seconds() + deadline = time.monotonic() + max(float(linger_budget), 0.0) + for _ in range(max_rounds): + wait = process_registry.wait_for_pending_completions(None, timeout=max(deadline - time.monotonic(), 0.0)) drained = process_registry.drain_notifications(session_key=key, owns_event=owns_event) texts = [ text for event, text in drained if event.get("type", "completion") == "completion" and text ] + if texts: + last = run_turn("\n\n".join(texts)) + if wait.get("timed_out"): + break if not texts: return last - last = run_turn("\n\n".join(texts)) return last diff --git a/tests/hermes_cli/test_quiet_turn_author.py b/tests/hermes_cli/test_quiet_turn_author.py index 23bd040f5f..d5608d6001 100644 --- a/tests/hermes_cli/test_quiet_turn_author.py +++ b/tests/hermes_cli/test_quiet_turn_author.py @@ -105,3 +105,33 @@ def test_quiet_one_shot_resumes_nested_notify_on_this_session_not_parent(monkeyp finally: while not process_registry.completion_queue.empty(): process_registry.completion_queue.get_nowait() + + +def test_quiet_notify_loop_shares_one_linger_budget(monkeypatch): + """One stuck notify_on_complete child is waited on once, not once per round plus finalize. + + The loop shares a single deadline across rounds and stops after draining once a + wait times out; the finalize pass is skipped entirely because the loop ran. + """ + from hermes_cli import quiet_single_query as qsq + from tools import process_registry as pr + + waits = [] + + def fake_wait(task_id=None, *, timeout=None, poll_interval=1.0): + waits.append(timeout) + # First wait: one stuck process times out; second wait (same run): nothing pending. + return {"waited": ["proc-stuck"] if len(waits) == 1 else [], "completed": [], "timed_out": ["proc-stuck"] if len(waits) == 1 else []} + + monkeypatch.setattr(pr.process_registry, "wait_for_pending_completions", fake_wait) + monkeypatch.setattr(pr.process_registry, "drain_notifications", lambda *a, **k: []) + monkeypatch.setattr(qsq.time, "monotonic", lambda: 100.0) + + calls = [] + result = qsq.continue_quiet_notify_completions( + "session-B", lambda text: calls.append(text) or {"final_response": text}, linger_budget=600.0, + ) + # Round 1: full budget; timed out -> drained (empty) -> break. Exactly one wait call. + assert waits == [600.0] + assert calls == [] + assert result is None