fix(anthropic-adapter): return a chunks-exposing stream so disconnects bill partial spend

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
kerry 2026-09-25 01:28:13 +00:00
parent 909bb4c87c
commit bf86914da4
3 changed files with 46 additions and 2 deletions

View file

@ -1197,3 +1197,35 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
return True
return False
class AnthropicSSEStream(AsyncIterator[bytes]):
"""
AsyncIterator[bytes] view of AnthropicStreamWrapper returned to callers of
translate_completion_output_params_streaming. Keeps the wrapper reachable so
the proxy's disconnect-time partial billing can read the inner chat stream's
collected chunks, messages, and model; a bare async generator would hide them.
"""
def __init__(self, anthropic_wrapper: AnthropicStreamWrapper) -> None:
self._anthropic_wrapper = anthropic_wrapper
self._byte_stream: Final[AsyncIterator[bytes]] = anthropic_wrapper.async_anthropic_sse_wrapper()
self._hidden_params: dict[str, object] = {} # mutable-ok: the proxy merges provider headers onto _hidden_params in place
@property
def chunks(self) -> list | None:
return self._anthropic_wrapper.chunks
@property
def messages(self) -> list | None:
return self._anthropic_wrapper.messages
@property
def model(self) -> str:
return self._anthropic_wrapper.model
async def __anext__(self) -> bytes:
return await self._byte_stream.__anext__()
async def aclose(self) -> None:
await self._byte_stream.aclose()

View file

@ -201,7 +201,7 @@ from litellm.types.llms.openai import (
from litellm.types.utils import Choices, ModelResponse, StreamingChoices, Usage
from litellm.utils import supports_mid_conversation_system
from .streaming_iterator import AnthropicStreamWrapper
from .streaming_iterator import AnthropicSSEStream, AnthropicStreamWrapper
if TYPE_CHECKING:
from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObject
@ -341,7 +341,7 @@ class AnthropicAdapter:
)
# Return the SSE-wrapped version for proper event formatting.
if is_async:
return anthropic_wrapper.async_anthropic_sse_wrapper()
return AnthropicSSEStream(anthropic_wrapper)
return anthropic_wrapper.anthropic_sse_wrapper()

View file

@ -630,6 +630,18 @@ class FallbackAwareAnthropicMessagesStream:
def has_buffered_provider_output(self) -> bool:
return getattr(self._source_iterator, "has_buffered_provider_output", False) is True
@property
def chunks(self) -> list | None:
return cast("list | None", getattr(self._source_iterator, "chunks", None)) # cast-ok: the billing helper itself treats chunks as an opaque getattr
@property
def messages(self) -> list | None:
return cast("list | None", getattr(self._source_iterator, "messages", None)) # cast-ok: messages is a plain list on the inner stream
@property
def model(self) -> str | None:
return cast("str | None", getattr(self._source_iterator, "model", None)) # cast-ok: model is a str on the inner stream
def adopt_fallback_source(self, fallback_response: object) -> None:
self._source_iterator = fallback_response
self.fallback_headers_adopted = True