fix(agent): don't seed continuation stub after a context-overflow stream death
When a stream delivered text and then died on a context-overflow / payload-too-large error, the partial content (often tens of KB) was seeded as a length-continuation stub, growing the transcript monotonically. In a session whose transcript cannot be compressed back under budget (protect_last_n covers everything -> no_progress, or the summary would itself be larger -> would_grow), every later request is larger than the one that just failed — an unrecoverable loop where the user sees a 30+ minute fake hang and the only remedy is killing the session (#106260). classify_api_error already labels these errors context_overflow / payload_too_large (should_compress=True). _partial_stream_stub now returns an EMPTY stub marked _overflow_terminal for that class instead of seeding the recovered text, and recover_from_truncation treats the marker as terminal: the turn ends via the recovery contract with a clear message (start /new) and the transcript is not polluted with the partial. Normal partials (network stall, output-cap truncation, tool-call drops) are unchanged — only the overflow error class changes behavior. Tests: stub marker + empty content; a real streamed partial hitting a 'maximum context length' error returns the terminal stub; recover_from_ truncation ends the turn (no fragment/nudge appended) on the marker while a normal stub still runs the continuation path. 67 streaming/continuation regressions pass.
This commit is contained in:
@@ -2174,11 +2174,17 @@ def cleanup_task_resources(agent, task_id: str) -> None:
|
||||
|
||||
|
||||
def _build_partial_stream_stub(role, full_content, full_reasoning, model_name, usage_obj, *,
|
||||
dropped_tool_names=None):
|
||||
dropped_tool_names=None, overflow_terminal=False):
|
||||
"""Stub for an SSE stream that ended without ``finish_reason`` after
|
||||
delivering content. Tagged ``PARTIAL_STREAM_STUB_ID`` + ``FINISH_REASON_LENGTH``
|
||||
so the loop enters its continuation/retry path instead of accepting
|
||||
truncated output as a complete turn (#32086)."""
|
||||
truncated output as a complete turn (#32086).
|
||||
|
||||
``overflow_terminal`` (``full_content=None``): the stream died on a
|
||||
context-overflow error. Seeding the recovered text as a continuation stub
|
||||
would grow every later request into the same overflow (#106260); the loop
|
||||
treats the marker as terminal and ends the turn via the recovery contract.
|
||||
"""
|
||||
return SimpleNamespace(
|
||||
id=PARTIAL_STREAM_STUB_ID,
|
||||
model=model_name,
|
||||
@@ -2190,6 +2196,7 @@ def _build_partial_stream_stub(role, full_content, full_reasoning, model_name, u
|
||||
)],
|
||||
usage=usage_obj,
|
||||
_dropped_tool_names=dropped_tool_names or None,
|
||||
_overflow_terminal=overflow_terminal,
|
||||
)
|
||||
|
||||
|
||||
@@ -3267,14 +3274,33 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
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.
|
||||
_stub = _build_partial_stream_stub("assistant", _partial_text, None,
|
||||
getattr(self.agent, "model", "unknown"), None, dropped_tool_names=_partial_names)
|
||||
_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 ""))
|
||||
if _cls.reason == FailoverReason.content_policy_blocked:
|
||||
_stub._content_filter_terminated = True
|
||||
# #106260: a context-overflow / payload-too-large 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.
|
||||
if _cls is not None and _cls.reason in (
|
||||
FailoverReason.context_overflow, FailoverReason.payload_too_large,
|
||||
):
|
||||
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,
|
||||
)
|
||||
_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
|
||||
|
||||
|
||||
@@ -28,6 +28,14 @@ _CONTINUABLE_MODES = {"chat_completions", "bedrock_converse", "anthropic_message
|
||||
_THINK_TAG_RE = re.compile(r'<(?:think|thinking|reasoning|REASONING_SCRATCHPAD)[^>]*>', re.IGNORECASE)
|
||||
_TRUNCATED_FINAL = "Response truncated due to output length limit"
|
||||
_FIRST_TRUNCATED_FINAL = "First response truncated due to output length limit"
|
||||
# #106260: a stream that died on a context-overflow error after partial delivery must not seed a
|
||||
# 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."
|
||||
)
|
||||
|
||||
_THINKING_EXHAUSTED = (
|
||||
"💭 Reasoning exhausted the output token budget — no visible response was produced.",
|
||||
@@ -326,6 +334,24 @@ def recover_from_truncation(
|
||||
force=True,
|
||||
)
|
||||
|
||||
# #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).
|
||||
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).",
|
||||
force=True,
|
||||
)
|
||||
return st.end_turn(
|
||||
_CONTEXT_OVERFLOW_PARTIAL_FINAL,
|
||||
error=_CONTEXT_OVERFLOW_PARTIAL_FINAL,
|
||||
failed=True,
|
||||
)
|
||||
|
||||
_trunc_msg = normalize_response_for_agent(agent, response)
|
||||
_trunc_content = getattr(_trunc_msg, "content", None) if _trunc_msg else None
|
||||
_trunc_has_tool_calls = bool(getattr(_trunc_msg, "tool_calls", None)) if _trunc_msg else False
|
||||
|
||||
@@ -0,0 +1,151 @@
|
||||
"""Regression tests for #106260: a context-overflow error after partial stream
|
||||
delivery must NOT seed a continuation stub with the recovered text.
|
||||
|
||||
Seeding tens of KB of partial content as a length-continuation stub makes every
|
||||
later request larger, so a session whose transcript cannot fit (compression
|
||||
failed / protect_last_n covers it) loops forever growing the context. The stub
|
||||
is instead marked terminal (content empty) and the loop ends the turn via the
|
||||
recovery contract.
|
||||
"""
|
||||
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
from hermes_constants import PARTIAL_STREAM_STUB_ID, FINISH_REASON_LENGTH
|
||||
|
||||
|
||||
def _make_agent():
|
||||
from run_agent import AIAgent
|
||||
|
||||
agent = AIAgent(
|
||||
api_key="test-key",
|
||||
base_url="https://example.com/v1",
|
||||
model="test/model",
|
||||
quiet_mode=True,
|
||||
skip_context_files=True,
|
||||
skip_memory=True,
|
||||
)
|
||||
agent.api_mode = "chat_completions"
|
||||
agent._interrupt_requested = False
|
||||
return agent
|
||||
|
||||
|
||||
def _make_stream_chunk(content=None, finish_reason=None):
|
||||
delta = SimpleNamespace(content=content, tool_calls=None,
|
||||
reasoning_content=None, reasoning=None)
|
||||
choice = SimpleNamespace(index=0, delta=delta, finish_reason=finish_reason)
|
||||
return SimpleNamespace(choices=[choice], model=None, usage=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(
|
||||
self, _mock_close, mock_create, monkeypatch,
|
||||
):
|
||||
"""A stream that delivered text then hit the provider's
|
||||
maximum-context-length error must return an EMPTY stub marked terminal,
|
||||
not a continuation stub carrying the recovered text (#106260)."""
|
||||
def _overflowing_stream():
|
||||
yield _make_stream_chunk(content="Here's my long partial answer ...")
|
||||
raise RuntimeError(
|
||||
"This model's maximum context length is 128000 tokens. "
|
||||
"However, you requested 140000 tokens."
|
||||
)
|
||||
|
||||
mock_client = MagicMock()
|
||||
mock_client.chat.completions.create.side_effect = lambda *a, **kw: _overflowing_stream()
|
||||
mock_create.return_value = mock_client
|
||||
|
||||
agent = _make_agent()
|
||||
agent._current_streamed_assistant_text = "Here's my long partial answer ..."
|
||||
|
||||
monkeypatch.setenv("HERMES_STREAM_RETRIES", "0")
|
||||
response = agent._interruptible_streaming_api_call({})
|
||||
|
||||
assert response.id == PARTIAL_STREAM_STUB_ID
|
||||
assert getattr(response, "_overflow_terminal", False) is True
|
||||
# The recovered text must not be seeded for continuation.
|
||||
assert response.choices[0].message.content is None
|
||||
|
||||
|
||||
class TestRecoverFromTruncationOverflowTerminal:
|
||||
def _mock_agent(self):
|
||||
agent = MagicMock()
|
||||
agent._vprint = MagicMock()
|
||||
agent._flush_status_buffer = MagicMock()
|
||||
agent._cleanup_task_resources = MagicMock()
|
||||
agent._persist_session = MagicMock()
|
||||
return agent
|
||||
|
||||
def _response(self, overflow_terminal=True, content="recovered text"):
|
||||
return SimpleNamespace(
|
||||
id=PARTIAL_STREAM_STUB_ID,
|
||||
_overflow_terminal=overflow_terminal,
|
||||
_dropped_tool_names=None,
|
||||
choices=[SimpleNamespace(
|
||||
index=0,
|
||||
message=SimpleNamespace(role="assistant", content=content,
|
||||
tool_calls=None, reasoning_content=None),
|
||||
finish_reason=FINISH_REASON_LENGTH,
|
||||
)],
|
||||
)
|
||||
|
||||
def test_overflow_terminal_stub_ends_turn_without_continuation(self):
|
||||
from agent.turn_truncation import (
|
||||
_CONTEXT_OVERFLOW_PARTIAL_FINAL,
|
||||
recover_from_truncation,
|
||||
)
|
||||
|
||||
agent = self._mock_agent()
|
||||
verdict = recover_from_truncation(
|
||||
agent, self._response(), 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=[],
|
||||
truncated_tool_call_retries=0, retry_count=0, compression_attempts=0,
|
||||
)
|
||||
|
||||
assert verdict.action == "return"
|
||||
result = verdict.result or {}
|
||||
assert result.get("failed") is True
|
||||
assert result.get("final_response") == _CONTEXT_OVERFLOW_PARTIAL_FINAL
|
||||
# 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
|
||||
Reference in New Issue
Block a user