fix(agent): close the interrupted tool tail on the overflow terminal

The overflow-terminal path ends the turn without reaching finalize_turn, so
a transcript that overflowed right after a tool batch ended on a raw tool
result; strict providers reject the next user turn (tool -> user). Close it
with the same final text, mirroring the truncated-tool-call terminal above.

Also: classify once before either log so an overflow no longer emits a
"so the loop can continue" WARNING followed by the contradicting "NOT
seeding" one; reset the stale-streak breaker once for both branches; drop
the "compression could not recover it" wording (this path never reached
compression); trim the test file to the three tests that bind behaviour
(stream -> terminal stub; 413 stays non-terminal; the terminal ends the
turn, closes the tool tail, carries compression_exhausted).
This commit is contained in:
kshitijk4poor
2026-09-09 17:26:33 +05:30
committed by kshitij
parent 3b0e81459b
commit 9e6c4100cb
3 changed files with 40 additions and 62 deletions

View File

@@ -3267,42 +3267,37 @@ class _StreamingCall(StreamingWaitMonitor):
logger.warning(
"Partial stream dropped tool call(s) %s after %s chars of text; surfaced warning to user: %s",
_partial_names, len(_partial_text or ""), error)
else:
logger.warning(
"Partial stream delivered before error; returning length-truncated stub with %s chars of "
"recovered content so the loop can continue from where the stream died: %s",
len(_partial_text or ""), error)
# Classify content filtering (MiniMax 1027, Azure content_filter, Anthropic refusal)
# before the error is swallowed into the stub: the loop reads the tag and falls back.
# Classify the error before it is swallowed into the stub: the loop reads the
# content-filter tag and falls back; a context overflow must not be continued at all.
_cls = None
with contextlib.suppress(Exception):
from agent.error_classifier import classify_api_error
_cls = classify_api_error(
error, provider=str(getattr(self.agent, "provider", "") or ""), model=str(getattr(self.agent, "model", "") or ""))
# #106260: a context-overflow error after partial delivery must NOT seed a
# continuation stub with the recovered text — the transcript already cannot fit
# (compression failed or protect_last_n covers it), and appending tens of KB only
# makes every later request larger. Return an EMPTY stub marked terminal: the loop
# ends the turn via the recovery contract instead of continuing into the same
# overflow. Scope is deliberately context_overflow ONLY: payload_too_large (413)
# has its own byte-scored recovery owner (turn_overflow._recover_payload_too_large,
# #88960/#47339) that must not be bypassed (review P1, andrexibiza).
_reset_stale_streak(self.agent) # deltas fired => provider responsive: clear the breaker
# #106260: continuing after a context-overflow error re-sends a larger request into the
# same overflow. Return an EMPTY stub marked terminal so the loop ends the turn instead.
# Scope is context_overflow ONLY: payload_too_large (413) has its own byte-scored recovery
# owner (turn_overflow._recover_payload_too_large, #88960/#47339) that must not be bypassed.
if _cls is not None and _cls.reason == FailoverReason.context_overflow:
logger.warning(
"Partial stream ended on a context-overflow error after %s chars; "
"NOT seeding a continuation stub (transcript is already over budget): %s",
len(_partial_text or ""), error,
)
_reset_stale_streak(self.agent)
return _build_partial_stream_stub(
"assistant", None, None, getattr(self.agent, "model", "unknown"), None,
dropped_tool_names=_partial_names, overflow_terminal=True,
)
if not _partial_names:
logger.warning(
"Partial stream delivered before error; returning length-truncated stub with %s chars of "
"recovered content so the loop can continue from where the stream died: %s",
len(_partial_text or ""), error)
_stub = _build_partial_stream_stub("assistant", _partial_text, None,
getattr(self.agent, "model", "unknown"), None, dropped_tool_names=_partial_names)
if _cls is not None and _cls.reason == FailoverReason.content_policy_blocked:
_stub._content_filter_terminated = True
_reset_stale_streak(self.agent) # deltas fired => provider responsive: clear the breaker
return _stub
def run(self):

View File

@@ -32,9 +32,9 @@ _FIRST_TRUNCATED_FINAL = "First response truncated due to output length limit"
# continuation — the transcript already cannot fit, and appending the partial stub grows every
# later request into the same overflow. End the turn via the recovery contract instead.
_CONTEXT_OVERFLOW_PARTIAL_FINAL = (
"The conversation exceeded the model's context window and compression could "
"not recover it, so the partial response was not continued. Start a new "
"session (or /new) to continue with a clean transcript."
"The request no longer fits the model's context window, so the partial "
"response was not continued. Continue in a fresh session (/new; gateway "
"chats are reset automatically)."
)
_THINKING_EXHAUSTED = (
@@ -345,18 +345,23 @@ def recover_from_truncation(
# #106260: a context-overflow error after partial delivery must not seed a
# continuation. _partial_stream_stub marks such stubs _overflow_terminal and
# leaves content empty; continuing would only re-send a larger request into
# the same overflow (compression already failed / protect_last_n covers it).
# the same overflow. The stub path never raises, so this class never reached
# recover_from_overflow's compress-and-retry on main either — ending the turn
# replaces a growth loop, not a compression attempt.
if getattr(st.response, "_overflow_terminal", False):
agent._flush_status_buffer()
agent._vprint(
f"{agent.log_prefix}⚠️ Stream ended on a context-overflow error after "
"partial delivery — not continuing (the transcript is already over "
"budget).",
"partial delivery — not continuing (the request no longer fits the model's "
"context window).",
force=True,
)
# Prior tool batches can leave a tool-result tail; this path never reaches
# finalize_turn (same as the truncated-tool-call terminal above).
close_interrupted_tool_sequence(st.messages, _CONTEXT_OVERFLOW_PARTIAL_FINAL)
# Carry the #98722 typed exhaustion bit so the gateway resets/moves future
# input to a clean session instead of leaving this bloated one authoritative
# for the next turn (review P1, andrexibiza).
# for the next turn.
return st.end_turn(
_CONTEXT_OVERFLOW_PARTIAL_FINAL,
error=_CONTEXT_OVERFLOW_PARTIAL_FINAL,

View File

@@ -40,21 +40,6 @@ def _make_stream_chunk(content=None, finish_reason=None):
class TestOverflowTerminalStub:
def test_build_stub_carries_marker_and_empty_content(self):
from agent.chat_completion_helpers import _build_partial_stream_stub
stub = _build_partial_stream_stub(
"assistant", None, None, "test/model", None, overflow_terminal=True)
assert stub._overflow_terminal is True
assert stub.id == PARTIAL_STREAM_STUB_ID
assert stub.choices[0].finish_reason == FINISH_REASON_LENGTH
assert stub.choices[0].message.content is None
normal = _build_partial_stream_stub(
"assistant", "some text", None, "test/model", None)
assert getattr(normal, "_overflow_terminal", False) is False
assert normal.choices[0].message.content == "some text"
@patch("run_agent.AIAgent._create_request_openai_client")
@patch("run_agent.AIAgent._close_request_openai_client")
def test_partial_stream_overflow_error_returns_terminal_stub(
@@ -145,9 +130,18 @@ class TestRecoverFromTruncationOverflowTerminal:
)
agent = self._mock_agent()
# The overflow fired right after a tool batch: the transcript tail is a
# raw tool result, which strict providers reject as tool -> user.
messages = [
{"role": "user", "content": "go"},
{"role": "assistant", "content": None,
"tool_calls": [{"id": "c1", "type": "function",
"function": {"name": "read_file", "arguments": "{}"}}]},
{"role": "tool", "tool_call_id": "c1", "content": "big file"},
]
verdict = recover_from_truncation(
agent, self._response(), FINISH_REASON_LENGTH, MagicMock(),
messages=[], conversation_history=None, api_kwargs={},
messages=messages, conversation_history=None, api_kwargs={},
api_call_count=0, effective_task_id=None, current_turn_user_idx=None,
length_continue_retries=0, truncated_response_parts=[],
truncated_tool_call_retries=0, retry_count=0, compression_attempts=0,
@@ -160,25 +154,9 @@ class TestRecoverFromTruncationOverflowTerminal:
# #98722 typed bit: the gateway consumes this to reset/move future input
# to a clean session instead of leaving the bloated one authoritative.
assert result.get("compression_exhausted") is True
# No fragment or nudge was appended to the transcript.
assert result.get("messages") == []
def test_normal_stub_is_not_treated_as_terminal(self):
from agent.turn_truncation import (
_CONTEXT_OVERFLOW_PARTIAL_FINAL,
recover_from_truncation,
)
agent = self._mock_agent()
verdict = recover_from_truncation(
agent, self._response(overflow_terminal=False), FINISH_REASON_LENGTH,
MagicMock(),
messages=[], conversation_history=None, api_kwargs={},
api_call_count=0, effective_task_id=None, current_turn_user_idx=None,
length_continue_retries=0, truncated_response_parts=["recovered text"],
truncated_tool_call_retries=0, retry_count=0, compression_attempts=0,
)
# Not the overflow-terminal final; the normal continuation path runs
# (no early return with the overflow message).
result = (verdict.result or {}) if verdict.action == "return" else None
assert result is None or result.get("final_response") != _CONTEXT_OVERFLOW_PARTIAL_FINAL
# The interrupted tool tail is closed so the next user turn alternates;
# no fragment or nudge was appended.
assert messages[-1]["role"] == "assistant"
assert messages[-1]["content"] == _CONTEXT_OVERFLOW_PARTIAL_FINAL
assert len(messages) == 4
assert result.get("messages") is messages