fix(delegate): nested orchestrators get their workers' results back — no 420 s deadline on delegate_task, summary budget uses the current prompt not the session sum
Two defects in the same path, both measured on the 1,393-agent refactor run. 1. A nested orchestrator (depth > 0) runs delegate_task synchronously by design: it needs its workers' results inside its own turn. But the sequential tool runner put every tool call under the generic 420 s deadline, and delegate_task was not exempt, so every batch longer than seven minutes returned "timed out after 420.0s" while the children kept running as orphans. 332 such timeouts in 234 orchestrator sessions; only 89 nested delegate_task calls in the whole run ever returned a real result. The orchestrators then spent 388 h of wall time polling: 1,526 reads of the live transcript files, 551 list actions, 242 h of explicit sleep, about $4k of API turns. delegate_task is now exempt from the sequential deadline (the batch owns its liveness: per-child heartbeats, the stale monitor, delegation.child_timeout_seconds). Live A/B, depth-1 orchestrator dispatching a 75 s leaf with the deadline set to 40 s (glm-5.3-flash via Nous): main -> "Error executing tool 'delegate_task': timed out after 40.0s", leaf result lost; branch -> orchestrator blocked 161 s and returned the leaf's LEAF_DONE_MARKER. 2. _parent_summary_char_budget computed the parent's context headroom from session_prompt_tokens, which is the running SUM of prompt tokens over every API call in the session. After a few hundred calls it exceeds any window, headroom goes negative, and every child summary collapses to the 2,000-char floor with the full text spilled to disk. All 1,393 child summaries in the run were truncated this way; the orchestrator planned from stubs. The budget now reads the last call's prompt_tokens from _last_turn_usage. Tests: delegate_task is in the exempt set and the set is narrow; budget for a long-lived parent equals the budget for a fresh parent with the same current prompt, and exceeds the floor.
This commit is contained in:
+10
-1
@@ -769,6 +769,15 @@ def _resolve_sequential_tool_timeout() -> float | None:
|
||||
return resolve_timeout("tools.sequential_call", default=_resolve_concurrent_tool_timeout())
|
||||
|
||||
|
||||
# Tools whose call blocks on a long-running operation that supervises its own liveness: no generic
|
||||
# sequential deadline. ``delegate_task`` in a nested orchestrator blocks for the whole batch by design
|
||||
# (children carry heartbeats, the stale monitor, and ``delegation.child_timeout_seconds``); under the
|
||||
# 420 s deadline every real batch "timed out" while its children ran on as orphans, and the orchestrator
|
||||
# spent the following hours polling transcripts (measured: 332 timeouts, ~$4k of orchestrator turns in
|
||||
# one run).
|
||||
_SEQUENTIAL_DEADLINE_EXEMPT_TOOLS = frozenset({"delegate_task"})
|
||||
|
||||
|
||||
def _abandoned_sequential_result(agent, ref: _ToolCallRef, message: str, result_cls, **outcome) -> _ManagedToolResult:
|
||||
"""Emit the terminal post_tool_call for a worker the sequential runner gave up on
|
||||
(timeout / interrupt) and wrap ``message`` in its marker ``result_cls``."""
|
||||
@@ -815,7 +824,7 @@ def _run_sequential_tool_execution_middleware(
|
||||
"""Run one sequential call on a worker thread under the concurrent executor's deadline.
|
||||
Interactive tools (``clarify``) own their wait via ``agent.clarify_timeout``; the
|
||||
generic deadline would report ``tool_timeout`` while the prompt is still live."""
|
||||
timeout_s = _resolve_sequential_tool_timeout()
|
||||
timeout_s = None if function_name in _SEQUENTIAL_DEADLINE_EXEMPT_TOOLS else _resolve_sequential_tool_timeout()
|
||||
ref = _ToolCallRef(function_name, function_args, effective_task_id, tool_call_id, middleware_trace)
|
||||
kwargs = dict(ref.middleware_kwargs(), execute=execute, scope_block=scope_block, display_index=display_index)
|
||||
if function_name in _NEVER_PARALLEL_TOOLS:
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
"""``delegate_task`` runs a nested orchestrator's whole batch inside one tool call by design, so it must not
|
||||
sit under the generic sequential-call deadline: with it, every batch longer than the deadline "timed out"
|
||||
while its children kept running as orphans and the orchestrator polled transcripts for hours."""
|
||||
from agent import tool_executor as te
|
||||
|
||||
|
||||
def test_delegate_task_is_exempt_from_the_sequential_deadline():
|
||||
assert "delegate_task" in te._SEQUENTIAL_DEADLINE_EXEMPT_TOOLS
|
||||
|
||||
|
||||
def test_exemption_is_narrow():
|
||||
assert "terminal" not in te._SEQUENTIAL_DEADLINE_EXEMPT_TOOLS
|
||||
assert "execute_code" not in te._SEQUENTIAL_DEADLINE_EXEMPT_TOOLS
|
||||
@@ -13,6 +13,7 @@ import tempfile
|
||||
import pytest
|
||||
|
||||
import tools.delegate_tool as dt
|
||||
from tools.delegate_tool_results import _MIN_SUMMARY_CHARS, _parent_summary_char_budget
|
||||
|
||||
|
||||
class _FakeCompressor:
|
||||
@@ -22,9 +23,11 @@ class _FakeCompressor:
|
||||
|
||||
|
||||
class _FakeParent:
|
||||
def __init__(self, context_length, used_tokens, max_tokens):
|
||||
def __init__(self, context_length, used_tokens, max_tokens, session_total=None):
|
||||
self.context_compressor = _FakeCompressor(context_length, max_tokens)
|
||||
self.session_prompt_tokens = used_tokens
|
||||
# Current prompt size (last call) drives the budget; the cumulative session counter must not.
|
||||
self._last_turn_usage = {"prompt_tokens": used_tokens}
|
||||
self.session_prompt_tokens = session_total if session_total is not None else used_tokens
|
||||
|
||||
|
||||
def test_small_summaries_pass_through_untouched():
|
||||
@@ -76,3 +79,12 @@ def test_empty_results_is_noop():
|
||||
[{"task_index": 0, "status": "failed", "summary": None}],
|
||||
_FakeParent(131_000, 1_000, 8_000),
|
||||
)
|
||||
|
||||
|
||||
def test_budget_uses_current_prompt_size_not_the_session_sum():
|
||||
"""A long-lived parent has a session sum far past any window while its current prompt is small; the
|
||||
budget must follow the current prompt, otherwise every summary collapses to the floor."""
|
||||
long_lived = _FakeParent(context_length=200_000, used_tokens=30_000, max_tokens=8_000, session_total=25_000_000)
|
||||
fresh = _FakeParent(context_length=200_000, used_tokens=30_000, max_tokens=8_000)
|
||||
assert _parent_summary_char_budget(long_lived, 1) == _parent_summary_char_budget(fresh, 1)
|
||||
assert _parent_summary_char_budget(long_lived, 1) > _MIN_SUMMARY_CHARS
|
||||
|
||||
@@ -219,15 +219,22 @@ def _trim_summary_with_footer(summary: str, cap: int, task_index: int) -> tuple[
|
||||
return head + "\n\n[... middle omitted — see footer ...]\n\n" + tail + "\n".join(footer_lines), spill_path
|
||||
|
||||
def _parent_summary_char_budget(parent_agent, n_summaries: int) -> Optional[int]:
|
||||
"""Per-summary char budget from the parent's *remaining* context headroom (context length − prompt tokens − the
|
||||
compressor's output reserve), a fraction of it split across the batch at ~4 chars/token. None when the parent's
|
||||
context state is unknown — caller then uses the static ceiling only."""
|
||||
"""Per-summary char budget from the parent's *remaining* context headroom (context length − the parent's
|
||||
current prompt size − the compressor's output reserve), a fraction of it split across the batch at ~4
|
||||
chars/token. None when the parent's context state is unknown — caller then uses the static ceiling only.
|
||||
|
||||
"Current prompt size" is the last API call's ``prompt_tokens`` (``_last_turn_usage``), never
|
||||
``session_prompt_tokens``: that field is the running SUM over every call in the session, so after a
|
||||
few hundred calls it exceeds any context window and every summary collapsed to the 2,000-char floor
|
||||
(measured: all 1,393 child summaries in one run, each spilled to disk, the orchestrator working from
|
||||
the stub)."""
|
||||
try:
|
||||
compressor = getattr(parent_agent, "context_compressor", None)
|
||||
context_length = getattr(compressor, "context_length", None)
|
||||
if not isinstance(context_length, int) or context_length <= 0:
|
||||
return None
|
||||
used_tokens = getattr(parent_agent, "session_prompt_tokens", 0)
|
||||
last_usage = getattr(parent_agent, "_last_turn_usage", None) or {}
|
||||
used_tokens = last_usage.get("prompt_tokens") if isinstance(last_usage, dict) else None
|
||||
if not isinstance(used_tokens, (int, float)) or used_tokens < 0:
|
||||
used_tokens = 0
|
||||
headroom_tokens = context_length - int(used_tokens) - int(getattr(compressor, "max_tokens", 0) or 0)
|
||||
|
||||
Reference in New Issue
Block a user