refactor(tools): split wait_for_message into respond_to_user + wait_for_agents

One tool was doing three jobs (wait on the user, wait on other agents, and
- wrongly - wait for a long-running command), so the driver had to guess which
one an agent meant and used parent_id as the proxy: the root waits for a human,
everyone else waits for agents. That proxy is wrong, since the user can message
any agent from the TUI's agent tree.

Tool identity now carries the intent, and the coordinator records it as a
wait_kind that survives snapshot/restore:

  respond_to_user  -> wait_kind="user",   never auto-resumed (root or not)
  wait_for_agents  -> wait_kind="agents", auto-resumed on a 300s timer
  recovery exhaust -> wait_kind="stalled"

respond_to_user fuses the message and the yield into one call, so there is no
way to answer and then forget to stop - the two-step that gpt-4o-mini skipped
2/2 in live testing. Plain text still renders as before.

Auto-resume is also bounded now: an agent that re-parks after every timeout
burned a model turn every 300s for the rest of the scan (and, since parked
children notify their parent, spammed the parent's inbox on the same cycle).
After _MAX_IDLE_AUTO_RESUMES it stays parked until a real message arrives.
This commit is contained in:
Ahmed Allam
2026-08-02 02:15:51 +03:00
committed by Ahmed Allam
parent 742f382836
commit 1c1fa49961
20 changed files with 539 additions and 128 deletions
+39 -1
View File
@@ -24,6 +24,11 @@ logger = logging.getLogger(__name__)
Status = Literal["running", "waiting", "completed", "stopped", "crashed", "failed", "budget_paused"]
# Why an agent parked. The user can message any agent, so this - not the agent's
# position in the tree - decides whether waiting is bounded: only an agent waiting
# on other agents is re-checked on a timer.
WaitKind = Literal["user", "agents", "stalled"]
@dataclass(slots=True)
class AgentRuntime:
@@ -47,6 +52,8 @@ class AgentCoordinator:
self.pending_counts: dict[str, int] = {}
self.errors: dict[str, str] = {}
self.recovery_counts: dict[str, int] = {}
self.idle_resume_counts: dict[str, int] = {}
self.wait_kinds: dict[str, WaitKind] = {}
self.runtimes: dict[str, AgentRuntime] = {}
self._lock = asyncio.Lock()
self._snapshot_path: Path | None = None
@@ -182,12 +189,21 @@ class AgentCoordinator:
if agent_id in self.statuses:
self.statuses[agent_id] = "running"
self.errors.pop(agent_id, None)
self.wait_kinds.pop(agent_id, None)
self.runtimes.setdefault(agent_id, AgentRuntime()).user_wake_required = False
await self._maybe_snapshot()
async def park_waiting(self, agent_id: str) -> None:
async def park_waiting(self, agent_id: str, *, wait_kind: WaitKind) -> None:
"""Park an agent, recording what it is waiting on so the driver can time it."""
async with self._lock:
if agent_id in self.statuses:
self.wait_kinds[agent_id] = wait_kind
await self.set_status(agent_id, "waiting")
async def wait_kind_of(self, agent_id: str) -> WaitKind | None:
async with self._lock:
return self.wait_kinds.get(agent_id)
async def record_recovery(self, agent_id: str) -> int:
"""Count a turn that ended without a lifecycle tool call; return the new total.
@@ -207,6 +223,24 @@ class AgentCoordinator:
return
await self._maybe_snapshot()
async def record_idle_resume(self, agent_id: str) -> int:
"""Count an auto-resume that no message triggered; return the new total.
An agent that parks again after every auto-resume would otherwise burn a
model turn per timeout for the rest of the scan.
"""
async with self._lock:
count = self.idle_resume_counts.get(agent_id, 0) + 1
self.idle_resume_counts[agent_id] = count
await self._maybe_snapshot()
return count
async def reset_idle_resumes(self, agent_id: str) -> None:
async with self._lock:
if self.idle_resume_counts.pop(agent_id, None) is None:
return
await self._maybe_snapshot()
async def set_status(
self, agent_id: str, status: Status | str, *, error: str | None = None
) -> None:
@@ -418,6 +452,8 @@ class AgentCoordinator:
"metadata": {aid: dict(md) for aid, md in self.metadata.items()},
"pending_counts": dict(self.pending_counts),
"recovery_counts": dict(self.recovery_counts),
"idle_resume_counts": dict(self.idle_resume_counts),
"wait_kinds": dict(self.wait_kinds),
"mailboxes": {
aid: [dict(m) for m in runtime.mailbox]
for aid, runtime in self.runtimes.items()
@@ -438,6 +474,8 @@ class AgentCoordinator:
self.pending_counts = dict(snap.get("pending_counts", {}))
self.errors = dict(snap.get("errors", {}))
self.recovery_counts = dict(snap.get("recovery_counts", {}))
self.idle_resume_counts = dict(snap.get("idle_resume_counts", {}))
self.wait_kinds = dict(snap.get("wait_kinds", {}))
mailboxes = snap.get("mailboxes", {})
if isinstance(mailboxes, dict):
for aid, msgs in mailboxes.items():
+43 -18
View File
@@ -227,7 +227,7 @@ async def run_agent_loop(
return result
while True:
timeout = await _plain_waiting_timeout(coordinator, agent_id, context)
timeout = await _plain_waiting_timeout(coordinator, agent_id)
try:
woke = await coordinator.wait_for_message(agent_id, timeout=timeout)
except asyncio.CancelledError:
@@ -245,7 +245,19 @@ async def run_agent_loop(
# Real input is real progress, so the nudge budget starts over. A bare
# auto-resume is not: it must not hand a wedged agent a fresh budget.
await coordinator.reset_recovery(agent_id)
await coordinator.reset_idle_resumes(agent_id)
else:
idle_resumes = await coordinator.record_idle_resume(agent_id)
if idle_resumes >= _MAX_IDLE_AUTO_RESUMES:
logger.warning(
"agent %s auto-resumed %d times without hearing from anyone; "
"leaving it parked until a real message arrives",
agent_id,
idle_resumes,
)
await coordinator.park_waiting(agent_id, wait_kind="stalled")
await _notify_parent_on_stall(coordinator, agent_id)
continue
logger.info("agent %s reached its waiting timeout; auto-resuming", agent_id)
await coordinator.send(
agent_id,
@@ -438,10 +450,10 @@ async def _run_until_lifecycle(
) -> RunResultBase | None:
"""Drive an agent until an explicit lifecycle tool settles its status.
A turn that ends without ``finish_scan``, ``agent_finish``, or
``wait_for_message`` leaves the agent ``running``: plain text never
terminates a run and never yields to the user. Such a turn is nudged back
into a tool call, bounded by a recovery limit.
A turn that ends without ``finish_scan``, ``agent_finish``,
``respond_to_user``, or ``wait_for_agents`` leaves the agent ``running``:
plain text never terminates a run and never yields to the user. Such a turn
is nudged back into a tool call, bounded by a recovery limit.
"""
result: RunResultBase | None = None
input_data: Any = initial_input
@@ -521,8 +533,8 @@ async def _exhausted_recovery(
) -> RunResultBase | None:
"""Settle an agent that never recovered into a tool call.
Interactive runs park instead of dying: a human is present, so the scan
stays resumable by sending another message. Autonomous runs have nobody to
Interactive runs park instead of dying: a human is attached and can message
any agent, so the scan stays resumable. Autonomous runs have nobody to
resume them, so they fail loudly.
"""
if not interactive:
@@ -536,7 +548,7 @@ async def _exhausted_recovery(
"agent %s exhausted tool-call recovery attempts; parking until a message arrives",
agent_id,
)
await coordinator.set_status(agent_id, "waiting")
await coordinator.park_waiting(agent_id, wait_kind="stalled")
# A parked child owes its parent a completion report it can no longer send. The
# parent is an agent, not a watching human, so nothing else tells it to stop
# waiting and it burns its full timeout on a message that is never coming.
@@ -546,23 +558,35 @@ async def _exhausted_recovery(
_WAITING_AUTO_RESUME_TIMEOUT_S = 300.0
# An agent that parks again after every auto-resume makes no progress, so stop
# spending a model turn per timeout and leave it parked for a real message.
_MAX_IDLE_AUTO_RESUMES = 3
async def _plain_waiting_timeout(
coordinator: AgentCoordinator,
agent_id: str,
context: dict[str, Any],
) -> float | None:
"""Auto-resume timeout for a plainly-waiting subagent; None waits forever."""
if context.get("parent_id") is None:
return None
"""Auto-resume timeout for a parked agent; None waits until a message arrives.
Driven by what the agent is waiting on, not by where it sits in the graph:
the user can message any agent, so an agent awaiting a human parks
indefinitely whether or not it is the root. Only an agent awaiting other
agents is re-checked on a timer, and only until it has spent its idle
budget re-parking without hearing anything.
"""
async with coordinator._lock:
status = coordinator.statuses.get(agent_id)
has_error = agent_id in coordinator.errors
runtime = coordinator.runtimes.get(agent_id)
gated = runtime.user_wake_required if runtime is not None else False
if status == "waiting" and not has_error and not gated:
return _WAITING_AUTO_RESUME_TIMEOUT_S
return None
wait_kind = coordinator.wait_kinds.get(agent_id)
idle_resumes = coordinator.idle_resume_counts.get(agent_id, 0)
if status != "waiting" or has_error or gated:
return None
if wait_kind != "agents" or idle_resumes >= _MAX_IDLE_AUTO_RESUMES:
return None
return _WAITING_AUTO_RESUME_TIMEOUT_S
async def _run_cycle_parked(
@@ -801,8 +825,9 @@ async def _append_tool_required_message(
"Your previous message ended a turn without a tool call. Plain text never ends "
"execution and never hands control to the user: it is shown to the user, and the "
"run continues. Continue immediately and call exactly one tool. "
"If you have finished responding and want the user's next message, call "
"wait_for_message. "
"If you have something to tell the user and nothing to do until they reply, "
"call respond_to_user. "
"If you are blocked waiting for another agent, call wait_for_agents. "
f"If the whole engagement is complete, call {finish_tool}. "
"Otherwise use the appropriate execution or planning tool. "
f"This is recovery attempt {attempt}/{limit}."
@@ -813,7 +838,7 @@ async def _append_tool_required_message(
"call. That is invalid in non-interactive mode; plain text final answers are "
"ignored. Continue immediately and call exactly one tool. "
f"If your work is complete, call {finish_tool}. "
"If you are blocked waiting for another agent, call wait_for_message. "
"If you are blocked waiting for another agent, call wait_for_agents. "
"Otherwise use the appropriate execution or planning tool. "
f"This is recovery attempt {attempt}/{limit}."
)