fix(streaming): shape anthropic partial stubs correctly

Hand-applied from 5ab566fea1 (#45919): _build_partial_stream_stub now
takes api_mode and returns a Messages-shaped stub in anthropic_messages
mode, including the overflow_terminal and content-filter stub sites.

(cherry picked from commit 5ab566fea1, grafted onto the split helper)
This commit is contained in:
LeonSGP43
2026-09-24 16:29:57 +05:30
committed by kshitij
parent d26909c99e
commit fca694c3e7

View File

@@ -2365,7 +2365,7 @@ 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, overflow_terminal=False):
dropped_tool_names=None, overflow_terminal=False, api_mode=None):
"""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
@@ -2375,7 +2375,26 @@ def _build_partial_stream_stub(role, full_content, full_reasoning, model_name, u
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.
``api_mode="anthropic_messages"`` returns a Messages-shaped stub (``content``
block list + ``stop_reason="max_tokens"``) so AnthropicTransport validates it
and the loop continues instead of entering the invalid-response retry ladder
(#45908). Empty content keeps one empty text block: validate_response rejects
an empty list for ``max_tokens``.
"""
if api_mode == "anthropic_messages":
return SimpleNamespace(
id=PARTIAL_STREAM_STUB_ID,
type="message",
role=role,
model=model_name,
content=[SimpleNamespace(type="text", text=full_content or "")],
stop_reason="max_tokens",
stop_sequence=None,
usage=usage_obj,
_dropped_tool_names=dropped_tool_names or None,
_overflow_terminal=overflow_terminal,
)
return SimpleNamespace(
id=PARTIAL_STREAM_STUB_ID,
model=model_name,
@@ -3825,6 +3844,7 @@ class _StreamingCall(StreamingWaitMonitor):
return _build_partial_stream_stub(
"assistant", None, None, getattr(self.agent, "model", "unknown"), None,
dropped_tool_names=_partial_names, overflow_terminal=True,
api_mode=getattr(self.agent, "api_mode", None),
)
if not _partial_names:
logger.warning(
@@ -3832,7 +3852,8 @@ class _StreamingCall(StreamingWaitMonitor):
"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)
getattr(self.agent, "model", "unknown"), None, dropped_tool_names=_partial_names,
api_mode=getattr(self.agent, "api_mode", None))
if _cls is not None and _cls.reason == FailoverReason.content_policy_blocked:
_stub._content_filter_terminated = True
return _stub