fix(cli): one shared linger budget for the quiet -Q notify-resume loop
The salvaged loop called wait_for_pending_completions(None) with a fresh default 600s deadline on every round, and _finalize_single_query then re-waited the full timeout on the same stuck notify_on_complete child: a hung child blocked a quiet one-shot 2x-9x longer than before. One deadline now covers the whole run — the loop passes the remaining budget each round, stops after draining once a wait times out (a timed-out process never fires this run), and the finalize pass skips its re-wait when the loop already consumed the budget. Regression: test_quiet_notify_loop_shares_one_linger_budget (one wait call, full budget, on a stuck child).
This commit is contained in:
15
cli.py
15
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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user