refactor(streaming): parameterize delegated chunks and messages types

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
kerry 2026-09-25 06:02:54 +00:00
parent d343d3f6ba
commit 818e9e4f8b
3 changed files with 34 additions and 22 deletions

View file

@ -11,6 +11,7 @@ from typing import (
Final,
Literal,
Protocol,
cast,
get_args,
)
@ -35,6 +36,7 @@ from litellm.types.utils import AdapterCompletionStreamWrapper, Delta
if TYPE_CHECKING:
from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObject
from litellm.types.llms.openai import AllMessageValues
from litellm.types.utils import ModelResponseStream
@ -120,12 +122,16 @@ class _CombinedChunkSplitter:
self._buffer: deque[ModelResponseStream] = deque()
@property
def chunks(self) -> list | None:
return getattr(self._stream, "chunks", None)
def chunks(self) -> "list[ModelResponseStream] | None":
return cast( # cast-ok: chunks is a list of ModelResponseStream on the inner stream
"list[ModelResponseStream] | None", getattr(self._stream, "chunks", None)
)
@property
def messages(self) -> list | None:
return getattr(self._stream, "messages", None)
def messages(self) -> "list[AllMessageValues] | None":
return cast( # cast-ok: messages is a list of AllMessageValues on the inner stream
"list[AllMessageValues] | None", getattr(self._stream, "messages", None)
)
@staticmethod
def _is_combined(chunk: "ModelResponseStream") -> bool:
@ -364,12 +370,16 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
)
@property
def chunks(self) -> list | None:
return getattr(self.completion_stream, "chunks", None)
def chunks(self) -> "list[ModelResponseStream] | None":
return cast( # cast-ok: chunks is a list of ModelResponseStream on the inner stream
"list[ModelResponseStream] | None", getattr(self.completion_stream, "chunks", None)
)
@property
def messages(self) -> list | None:
return getattr(self.completion_stream, "messages", None)
def messages(self) -> "list[AllMessageValues] | None":
return cast( # cast-ok: messages is a list of AllMessageValues on the inner stream
"list[AllMessageValues] | None", getattr(self.completion_stream, "messages", None)
)
def _merge_usage_into_held_stop_reason_chunk(self, chunk: Any) -> MessageBlockDelta:
"""Merge usage data from ``chunk`` into the held ``message_delta`` chunk.
@ -1211,11 +1221,11 @@ class AnthropicSSEStream(AsyncIterator[bytes]):
] = {} # mutable-ok: the proxy merges provider headers onto _hidden_params in place
@property
def chunks(self) -> list | None:
def chunks(self) -> "list[ModelResponseStream] | None":
return self._anthropic_wrapper.chunks
@property
def messages(self) -> list | None:
def messages(self) -> "list[AllMessageValues] | None":
return self._anthropic_wrapper.messages
@property

View file

@ -16,6 +16,8 @@ from litellm.llms.anthropic.experimental_pass_through.messages.streaming_iterato
if TYPE_CHECKING:
from litellm.caching.caching_handler import LLMCachingHandler
from litellm.litellm_core_utils.litellm_logging import Logging as LiteLLMLoggingObj
from litellm.types.llms.openai import AllMessageValues
from litellm.types.utils import ModelResponseStream
CACHED_STREAM_EVENTS_KEY: Final = "litellm_cached_anthropic_sse_events"
@ -51,15 +53,15 @@ class AnthropicMessagesStreamCacheWriter:
return getattr(self.stream, "has_buffered_provider_output", False) is True
@property
def chunks(self) -> list | None:
return cast( # cast-ok: the billing helper itself treats chunks as an opaque getattr
"list | None", getattr(self.stream, "chunks", None)
def chunks(self) -> "list[ModelResponseStream] | None":
return cast( # cast-ok: chunks is a list of ModelResponseStream on the inner stream
"list[ModelResponseStream] | None", getattr(self.stream, "chunks", None)
)
@property
def messages(self) -> list | None:
return cast( # cast-ok: messages is a plain list on the inner stream
"list | None", getattr(self.stream, "messages", None)
def messages(self) -> "list[AllMessageValues] | None":
return cast( # cast-ok: messages is a list of AllMessageValues on the inner stream
"list[AllMessageValues] | None", getattr(self.stream, "messages", None)
)
@property

View file

@ -636,15 +636,15 @@ class FallbackAwareAnthropicMessagesStream:
return getattr(self._source_iterator, "has_buffered_provider_output", False) is True
@property
def chunks(self) -> list | None:
return cast( # cast-ok: the billing helper itself treats chunks as an opaque getattr
"list | None", getattr(self._source_iterator, "chunks", None)
def chunks(self) -> list[ModelResponseStream] | None:
return cast( # cast-ok: chunks is a list of ModelResponseStream on the inner stream
"list[ModelResponseStream] | None", getattr(self._source_iterator, "chunks", None)
)
@property
def messages(self) -> list | None:
return cast( # cast-ok: messages is a plain list on the inner stream
"list | None", getattr(self._source_iterator, "messages", None)
def messages(self) -> list[AllMessageValues] | None:
return cast( # cast-ok: messages is a list of AllMessageValues on the inner stream
"list[AllMessageValues] | None", getattr(self._source_iterator, "messages", None)
)
@property