From 00e62487ca7b44fb23e43e27260c04d11cbbbde1 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 02:25:40 +0530 Subject: [PATCH 1/8] fix(proxy): log Anthropic pass-through in-band stream errors as failures --- .../streaming_handler.py | 20 +++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/litellm/proxy/pass_through_endpoints/streaming_handler.py b/litellm/proxy/pass_through_endpoints/streaming_handler.py index 19d8b063dd7..6cda69bd8fe 100644 --- a/litellm/proxy/pass_through_endpoints/streaming_handler.py +++ b/litellm/proxy/pass_through_endpoints/streaming_handler.py @@ -12,6 +12,7 @@ from litellm.litellm_core_utils.asyncify import asyncify from litellm.litellm_core_utils.core_helpers import bind_budget_reservation_to_callbacks from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER +from litellm.llms.anthropic.common_utils import AnthropicError from litellm.proxy._types import PassThroughEndpointLoggingResultValues from litellm.proxy.common_request_processing import ProxyBaseLLMRequestProcessing from litellm.proxy.common_utils.sse_keepalive import split_complete_sse_frames @@ -297,6 +298,7 @@ class PassThroughStreamingHandler: from litellm.llms.anthropic.pass_through.messages.streaming_iterator import ( _is_message_stop_chunk, # pyright: ignore[reportPrivateUsage] # both native stream paths share terminal-event detection _is_provider_error_chunk, # pyright: ignore[reportPrivateUsage] # provider errors must not become cache evidence + parse_anthropic_error_event, ) # Transport reads can split event names and JSON payloads. Recognize terminal @@ -312,6 +314,24 @@ class PassThroughStreamingHandler: and _is_message_stop_chunk(complete_frames) and not _is_provider_error_chunk(complete_frames) ) + # An in-band ``event: error`` on a cleanly closed 200 stream is an upstream + # failure, not a success: log it as one (partial usage included). + provider_error: Final = parse_anthropic_error_event(complete_frames) if complete_frames else None + if endpoint_type == EndpointType.ANTHROPIC and provider_error is not None: + _, error_message, error_status_code = provider_error + await PassThroughStreamingHandler.schedule_stream_failure_logging( + litellm_logging_obj=litellm_logging_obj, + endpoint_type=endpoint_type, + request_body=request_body, + raw_bytes=raw_bytes, + exception=AnthropicError(status_code=error_status_code, message=error_message), + stream_context=PassThroughStreamContext( + passthrough_success_handler_obj=passthrough_success_handler_obj, + url_route=url_route, + start_time=start_time, + ), + ) + return try: # TinyFish billing is owned by the detached poller; the $0 fallback below is only for streams with no run_id if endpoint_type == EndpointType.TINYFISH: From 6853826396b545f5b5e932b21cda28fa82849612 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 02:26:20 +0530 Subject: [PATCH 2/8] test(proxy): cover Anthropic pass-through in-band error frame logging --- ...streaming_handler_anthropic_error_frame.py | 103 ++++++++++++++++++ 1 file changed, 103 insertions(+) create mode 100644 tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py diff --git a/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py b/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py new file mode 100644 index 00000000000..f24edc41c8c --- /dev/null +++ b/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py @@ -0,0 +1,103 @@ +import json +from datetime import datetime +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj +from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER +from litellm.proxy.pass_through_endpoints.streaming_handler import ( + PassThroughStreamingHandler, +) +from litellm.proxy.pass_through_endpoints.success_handler import ( + PassThroughEndpointLogging, +) +from litellm.types.passthrough_endpoints.pass_through_endpoints import EndpointType + +MODEL = "claude-fable-5" + + +def _sse(event: str, data: dict) -> bytes: + return f"event: {event}\ndata: {json.dumps(data)}\n\n".encode() + + +def _partial_anthropic_stream() -> list[bytes]: + message_start = { + "type": "message_start", + "message": { + "id": "msg_partial", + "type": "message", + "role": "assistant", + "model": MODEL, + "content": [], + "stop_reason": None, + "stop_sequence": None, + "usage": {"input_tokens": 29, "output_tokens": 2}, + }, + } + block_start = {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}} + delta = {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "partial answer"}} + return [ + _sse("message_start", message_start), + _sse("content_block_start", block_start), + _sse("content_block_delta", delta), + ] + + +def _error_frame(error_type: str, message: str) -> bytes: + return _sse("error", {"type": "error", "error": {"type": error_type, "message": message}}) + + +def _logging_obj() -> MagicMock: + logging_obj = MagicMock(spec=LiteLLMLoggingObj) + logging_obj.model_call_details = {"model": MODEL, "stream": True} + logging_obj.optional_params = {} + logging_obj.litellm_params = {} + logging_obj.litellm_call_id = "test-call-id" + logging_obj.get_router_model_id.return_value = None + logging_obj.dispatch_success_handlers = AsyncMock() + logging_obj.dispatch_failure_handlers = AsyncMock() + return logging_obj + + +async def _route(logging_obj: MagicMock, raw_bytes: list[bytes]) -> None: + await PassThroughStreamingHandler._route_streaming_logging_to_handler( + litellm_logging_obj=logging_obj, + passthrough_success_handler_obj=PassThroughEndpointLogging(), + url_route="/anthropic/v1/messages", + request_body={"model": MODEL, "stream": True}, + endpoint_type=EndpointType.ANTHROPIC, + start_time=datetime.now(), + raw_bytes=raw_bytes, + end_time=datetime.now(), + model=MODEL, + ) + await GLOBAL_LOGGING_WORKER.flush() + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "error_type, expected_status", + [("api_error", 500), ("overloaded_error", 503), ("rate_limit_error", 429)], +) +async def test_in_band_error_frame_is_logged_as_failure_not_success(error_type: str, expected_status: int): + logging_obj = _logging_obj() + + await _route(logging_obj, [*_partial_anthropic_stream(), _error_frame(error_type, "boom")]) + + logging_obj.dispatch_success_handlers.assert_not_awaited() + logging_obj.dispatch_failure_handlers.assert_awaited_once() + exception = logging_obj.dispatch_failure_handlers.await_args.args[0] + assert exception.status_code == expected_status + assert "boom" in str(exception) + logging_obj.record_partial_usage_for_failure.assert_called_once() + + +@pytest.mark.asyncio +async def test_stream_without_error_frame_still_logs_success(): + logging_obj = _logging_obj() + + await _route(logging_obj, _partial_anthropic_stream()) + + logging_obj.dispatch_success_handlers.assert_awaited_once() + logging_obj.dispatch_failure_handlers.assert_not_awaited() From 28197fe0cb6c705cf1f4394e7dbb8a8d62d03348 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 06:41:02 +0530 Subject: [PATCH 3/8] fix(proxy): drop narrating comment --- litellm/proxy/pass_through_endpoints/streaming_handler.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/litellm/proxy/pass_through_endpoints/streaming_handler.py b/litellm/proxy/pass_through_endpoints/streaming_handler.py index 6cda69bd8fe..e9bd3113693 100644 --- a/litellm/proxy/pass_through_endpoints/streaming_handler.py +++ b/litellm/proxy/pass_through_endpoints/streaming_handler.py @@ -314,8 +314,6 @@ class PassThroughStreamingHandler: and _is_message_stop_chunk(complete_frames) and not _is_provider_error_chunk(complete_frames) ) - # An in-band ``event: error`` on a cleanly closed 200 stream is an upstream - # failure, not a success: log it as one (partial usage included). provider_error: Final = parse_anthropic_error_event(complete_frames) if complete_frames else None if endpoint_type == EndpointType.ANTHROPIC and provider_error is not None: _, error_message, error_status_code = provider_error From 869a5df8cc7f43b2423e5a41940e95c2a86a1be5 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 06:43:04 +0530 Subject: [PATCH 4/8] test(proxy): move in-band error tests into the existing streaming handler test file --- .../test_streaming_handler.py | 70 ++++++++++++++++++- 1 file changed, 69 insertions(+), 1 deletion(-) diff --git a/tests/unit/proxy/pass_through_endpoints/test_streaming_handler.py b/tests/unit/proxy/pass_through_endpoints/test_streaming_handler.py index 9d4532df49a..897c951c482 100644 --- a/tests/unit/proxy/pass_through_endpoints/test_streaming_handler.py +++ b/tests/unit/proxy/pass_through_endpoints/test_streaming_handler.py @@ -1,7 +1,7 @@ import json from collections.abc import Iterator from datetime import datetime -from unittest.mock import MagicMock +from unittest.mock import AsyncMock, MagicMock import pytest @@ -230,3 +230,71 @@ async def test_failed_anthropic_stream_records_partial_usage_off_the_event_loop( partial_usage = logging_obj.record_partial_usage_for_failure.call_args.kwargs["usage"] assert partial_usage.completion_tokens > 100_000 assert_loop_stayed_free(took, lags) + + +CLAUDE_MODEL = "claude-fable-5" + + +def _anthropic_error_frame(error_type: str, message: str) -> bytes: + payload = {"type": "error", "error": {"type": error_type, "message": message}} + return f"event: error\ndata: {json.dumps(payload)}\n\n".encode() + + +def _anthropic_error_logging_obj() -> MagicMock: + logging_obj = MagicMock(spec=LiteLLMLoggingObj) + logging_obj.model_call_details = {"model": CLAUDE_MODEL, "stream": True} + logging_obj.optional_params = {} + logging_obj.litellm_params = {} + logging_obj.litellm_call_id = "test-call-id" + logging_obj.get_router_model_id.return_value = None + logging_obj.dispatch_success_handlers = AsyncMock() + logging_obj.dispatch_failure_handlers = AsyncMock() + return logging_obj + + +async def _route_anthropic_stream(logging_obj: MagicMock, raw_bytes: list[bytes]) -> None: + await PassThroughStreamingHandler._route_streaming_logging_to_handler( + litellm_logging_obj=logging_obj, + passthrough_success_handler_obj=PassThroughEndpointLogging(), + url_route="/anthropic/v1/messages", + request_body={"model": CLAUDE_MODEL, "stream": True}, + endpoint_type=EndpointType.ANTHROPIC, + start_time=datetime.now(), + raw_bytes=raw_bytes, + end_time=datetime.now(), + model=CLAUDE_MODEL, + ) + await GLOBAL_LOGGING_WORKER.flush() + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "error_type, expected_status", + [("api_error", 500), ("overloaded_error", 503), ("rate_limit_error", 429)], +) +async def test_in_band_error_frame_is_logged_as_failure_not_success(error_type: str, expected_status: int): + logging_obj = _anthropic_error_logging_obj() + + await _route_anthropic_stream( + logging_obj, + [*_interrupted_anthropic_stream(CLAUDE_MODEL, "partial answer"), _anthropic_error_frame(error_type, "boom")], + ) + + logging_obj.dispatch_success_handlers.assert_not_awaited() + logging_obj.dispatch_failure_handlers.assert_awaited_once() + exception = logging_obj.dispatch_failure_handlers.await_args.args[0] + assert exception.status_code == expected_status + assert "boom" in str(exception) + partial_usage = logging_obj.record_partial_usage_for_failure.call_args.kwargs["usage"] + assert partial_usage.prompt_tokens == 29 + assert partial_usage.completion_tokens > 0 + + +@pytest.mark.asyncio +async def test_stream_without_error_frame_still_logs_success(): + logging_obj = _anthropic_error_logging_obj() + + await _route_anthropic_stream(logging_obj, _interrupted_anthropic_stream(CLAUDE_MODEL, "partial answer")) + + logging_obj.dispatch_success_handlers.assert_awaited_once() + logging_obj.dispatch_failure_handlers.assert_not_awaited() From b855bf891fd014cedaa8de8f27a624eb4c8e501e Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 06:43:37 +0530 Subject: [PATCH 5/8] test(proxy): remove duplicate test file --- ...streaming_handler_anthropic_error_frame.py | 103 ------------------ 1 file changed, 103 deletions(-) delete mode 100644 tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py diff --git a/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py b/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py deleted file mode 100644 index f24edc41c8c..00000000000 --- a/tests/test_litellm/proxy/pass_through_endpoints/test_streaming_handler_anthropic_error_frame.py +++ /dev/null @@ -1,103 +0,0 @@ -import json -from datetime import datetime -from unittest.mock import AsyncMock, MagicMock - -import pytest - -from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj -from litellm.litellm_core_utils.logging_worker import GLOBAL_LOGGING_WORKER -from litellm.proxy.pass_through_endpoints.streaming_handler import ( - PassThroughStreamingHandler, -) -from litellm.proxy.pass_through_endpoints.success_handler import ( - PassThroughEndpointLogging, -) -from litellm.types.passthrough_endpoints.pass_through_endpoints import EndpointType - -MODEL = "claude-fable-5" - - -def _sse(event: str, data: dict) -> bytes: - return f"event: {event}\ndata: {json.dumps(data)}\n\n".encode() - - -def _partial_anthropic_stream() -> list[bytes]: - message_start = { - "type": "message_start", - "message": { - "id": "msg_partial", - "type": "message", - "role": "assistant", - "model": MODEL, - "content": [], - "stop_reason": None, - "stop_sequence": None, - "usage": {"input_tokens": 29, "output_tokens": 2}, - }, - } - block_start = {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}} - delta = {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "partial answer"}} - return [ - _sse("message_start", message_start), - _sse("content_block_start", block_start), - _sse("content_block_delta", delta), - ] - - -def _error_frame(error_type: str, message: str) -> bytes: - return _sse("error", {"type": "error", "error": {"type": error_type, "message": message}}) - - -def _logging_obj() -> MagicMock: - logging_obj = MagicMock(spec=LiteLLMLoggingObj) - logging_obj.model_call_details = {"model": MODEL, "stream": True} - logging_obj.optional_params = {} - logging_obj.litellm_params = {} - logging_obj.litellm_call_id = "test-call-id" - logging_obj.get_router_model_id.return_value = None - logging_obj.dispatch_success_handlers = AsyncMock() - logging_obj.dispatch_failure_handlers = AsyncMock() - return logging_obj - - -async def _route(logging_obj: MagicMock, raw_bytes: list[bytes]) -> None: - await PassThroughStreamingHandler._route_streaming_logging_to_handler( - litellm_logging_obj=logging_obj, - passthrough_success_handler_obj=PassThroughEndpointLogging(), - url_route="/anthropic/v1/messages", - request_body={"model": MODEL, "stream": True}, - endpoint_type=EndpointType.ANTHROPIC, - start_time=datetime.now(), - raw_bytes=raw_bytes, - end_time=datetime.now(), - model=MODEL, - ) - await GLOBAL_LOGGING_WORKER.flush() - - -@pytest.mark.asyncio -@pytest.mark.parametrize( - "error_type, expected_status", - [("api_error", 500), ("overloaded_error", 503), ("rate_limit_error", 429)], -) -async def test_in_band_error_frame_is_logged_as_failure_not_success(error_type: str, expected_status: int): - logging_obj = _logging_obj() - - await _route(logging_obj, [*_partial_anthropic_stream(), _error_frame(error_type, "boom")]) - - logging_obj.dispatch_success_handlers.assert_not_awaited() - logging_obj.dispatch_failure_handlers.assert_awaited_once() - exception = logging_obj.dispatch_failure_handlers.await_args.args[0] - assert exception.status_code == expected_status - assert "boom" in str(exception) - logging_obj.record_partial_usage_for_failure.assert_called_once() - - -@pytest.mark.asyncio -async def test_stream_without_error_frame_still_logs_success(): - logging_obj = _logging_obj() - - await _route(logging_obj, _partial_anthropic_stream()) - - logging_obj.dispatch_success_handlers.assert_awaited_once() - logging_obj.dispatch_failure_handlers.assert_not_awaited() From 40774a76ddc4abc49d245b2f7df8af9648a3cccb Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 08:44:02 +0530 Subject: [PATCH 6/8] fix(proxy): type the route logging request body to satisfy pyright --- litellm/proxy/pass_through_endpoints/streaming_handler.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/proxy/pass_through_endpoints/streaming_handler.py b/litellm/proxy/pass_through_endpoints/streaming_handler.py index e9bd3113693..630b87d4cb9 100644 --- a/litellm/proxy/pass_through_endpoints/streaming_handler.py +++ b/litellm/proxy/pass_through_endpoints/streaming_handler.py @@ -280,7 +280,7 @@ class PassThroughStreamingHandler: litellm_logging_obj: LiteLLMLoggingObj, passthrough_success_handler_obj: PassThroughEndpointLogging, url_route: str, - request_body: dict, + request_body: dict[str, object], endpoint_type: EndpointType, start_time: datetime, raw_bytes: Sequence[bytes], From 20a4df16d8d22419b4429e572e02ea11ae25b02c Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 11:45:34 +0530 Subject: [PATCH 7/8] fix(proxy): keep route logging signature, suppress one unknown-arg with reason --- litellm/proxy/pass_through_endpoints/streaming_handler.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/litellm/proxy/pass_through_endpoints/streaming_handler.py b/litellm/proxy/pass_through_endpoints/streaming_handler.py index 630b87d4cb9..dc1129f961b 100644 --- a/litellm/proxy/pass_through_endpoints/streaming_handler.py +++ b/litellm/proxy/pass_through_endpoints/streaming_handler.py @@ -280,7 +280,7 @@ class PassThroughStreamingHandler: litellm_logging_obj: LiteLLMLoggingObj, passthrough_success_handler_obj: PassThroughEndpointLogging, url_route: str, - request_body: dict[str, object], + request_body: dict, endpoint_type: EndpointType, start_time: datetime, raw_bytes: Sequence[bytes], @@ -320,7 +320,7 @@ class PassThroughStreamingHandler: await PassThroughStreamingHandler.schedule_stream_failure_logging( litellm_logging_obj=litellm_logging_obj, endpoint_type=endpoint_type, - request_body=request_body, + request_body=request_body, # pyright: ignore[reportUnknownArgumentType] # request_body is an untyped dict in this signature raw_bytes=raw_bytes, exception=AnthropicError(status_code=error_status_code, message=error_message), stream_context=PassThroughStreamContext( From 670c1fb9ba912f4052bbabadd536ad87004ac5f8 Mon Sep 17 00:00:00 2001 From: Deepanshu Pal <40927968+DeepanshuPal@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:45:31 +0530 Subject: [PATCH 8/8] test(proxy): let the prompt cache observer finish on failure events --- tests/unit/proxy/hooks/test_prompt_cache_observer.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/tests/unit/proxy/hooks/test_prompt_cache_observer.py b/tests/unit/proxy/hooks/test_prompt_cache_observer.py index 82af3e9a6ef..b5fcad3add4 100644 --- a/tests/unit/proxy/hooks/test_prompt_cache_observer.py +++ b/tests/unit/proxy/hooks/test_prompt_cache_observer.py @@ -204,6 +204,9 @@ class RecordingObserver(PromptCacheObserver): await super().async_log_success_event(kwargs, response_obj, start_time, end_time) self.finished.set() + async def async_log_failure_event(self, kwargs, response_obj, start_time, end_time): + self.finished.set() + def native_response(): return {