fix(anthropic): surface mid-stream errors on the sync /v1/messages adapter

#43601 taught the async adapter to convert a mid-stream provider failure
into a terminal Anthropic `error` event instead of dropping the socket. The
sync path, returned by `transformation.py` whenever `is_async` is False,
never got the same treatment and was worse off: `__next__` caught every
exception and raised `StopIteration`, so the SSE stream simply ended as
though the response had completed. The client got a silently truncated
answer with no error at all.

Let `__next__` propagate, as `__anext__` already does, and give
`anthropic_sse_wrapper` the same guard as `async_anthropic_sse_wrapper` so
the failure is emitted as an Anthropic `error` event.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
revanth-045 2026-09-29 15:54:25 +05:30
parent 85dc7cb62e
commit 3de13bc79b
2 changed files with 70 additions and 9 deletions

View file

@ -752,7 +752,11 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
raise StopIteration
except Exception as e:
verbose_logger.error("Anthropic Adapter - %s\n%s", e, traceback.format_exc())
raise StopIteration
# Propagate, as ``__anext__`` does: ``anthropic_sse_wrapper`` turns
# this into a terminal Anthropic ``error`` event. Converting it to
# StopIteration here ended the SSE stream as though the response had
# completed, so the client saw a truncated answer and no error.
raise
async def __anext__(self):
from .transformation import LiteLLMAnthropicMessagesAdapter
@ -991,14 +995,18 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
This wrapper ensures dict chunks are SSE formatted with both event and data lines.
"""
for chunk in self:
if isinstance(chunk, dict):
event_type: str = str(chunk.get("type", "message"))
payload = f"event: {event_type}\ndata: {json.dumps(chunk)}\n\n"
yield payload.encode()
else:
# For non-dict chunks, forward the original value unchanged
yield chunk
try:
for chunk in self:
if isinstance(chunk, dict):
event_type: str = str(chunk.get("type", "message"))
payload = f"event: {event_type}\ndata: {json.dumps(chunk)}\n\n"
yield payload.encode()
else:
# For non-dict chunks, forward the original value unchanged
yield chunk
except Exception as e: # noqa: BLE001 # boundary before the socket: any upstream failure becomes an Anthropic error event
verbose_logger.exception("Anthropic Adapter - mid-stream error, emitting Anthropic error event: %s", e)
yield _mid_stream_error_sse_event(e)
async def async_anthropic_sse_wrapper(self) -> AsyncIterator[bytes]:
"""

View file

@ -61,6 +61,23 @@ class _AsyncStreamThenRaise:
raise self._exc
class _SyncStreamThenRaise:
"""Sync counterpart of ``_AsyncStreamThenRaise``."""
def __init__(self, items: List[MagicMock], exc: BaseException):
self._it = iter(items)
self._exc = exc
def __iter__(self):
return self
def __next__(self):
try:
return next(self._it)
except StopIteration:
raise self._exc
def _parse_sse(raw: bytes) -> tuple[str, dict]:
text = raw.decode()
event_line, data_line = text.strip().split("\n", 1)
@ -140,3 +157,39 @@ def test_error_event_preserves_midstream_fallback_error():
assert name == "error"
assert payload["error"]["type"] == "api_error"
assert "internalServerException" in payload["error"]["message"]
def test_sync_mid_stream_bedrock_error_becomes_anthropic_error_event():
"""The sync wrapper is reachable too (``is_async=False`` in
``transformation.py``), and has to surface the same terminal ``error``
event rather than letting the exception tear down the connection."""
chunks = [_make_chunk(Delta(content="Creating a file"))]
bedrock_err = BedrockError(
status_code=500,
message="Bedrock ConverseStream ended without a terminal 'messageStop' event",
)
wrapper = AnthropicStreamWrapper(
completion_stream=_SyncStreamThenRaise(chunks, bedrock_err),
model="bedrock-converse-sonnet-4-6",
)
events = list(wrapper.anthropic_sse_wrapper())
parsed = [_parse_sse(e) for e in events]
event_types = [name for name, _ in parsed]
assert "message_start" in event_types
assert event_types[-1] == "error"
_, error_payload = parsed[-1]
assert error_payload["type"] == "error"
assert error_payload["error"]["type"] == "api_error"
assert "messageStop" in error_payload["error"]["message"]
def test_sync_mid_stream_error_does_not_raise_out_of_wrapper():
"""Draining the sync wrapper must not re-raise the upstream exception."""
wrapper = AnthropicStreamWrapper(
completion_stream=_SyncStreamThenRaise([], BedrockError(status_code=500, message="boom")),
model="claude-x",
)
events = list(wrapper.anthropic_sse_wrapper())
assert _parse_sse(events[-1])[0] == "error"