diff --git a/litellm/llms/anthropic/pass_through/adapters/streaming_iterator.py b/litellm/llms/anthropic/pass_through/adapters/streaming_iterator.py index 12eee663ca5..fd20488d1f6 100644 --- a/litellm/llms/anthropic/pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/pass_through/adapters/streaming_iterator.py @@ -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]: """ diff --git a/tests/unit/llms/anthropic/pass_through/adapters/test_streaming_iterator_mid_stream_error.py b/tests/unit/llms/anthropic/pass_through/adapters/test_streaming_iterator_mid_stream_error.py index 4798d522182..228de4e3843 100644 --- a/tests/unit/llms/anthropic/pass_through/adapters/test_streaming_iterator_mid_stream_error.py +++ b/tests/unit/llms/anthropic/pass_through/adapters/test_streaming_iterator_mid_stream_error.py @@ -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"