From d72c5bbdf4523c6a6e6ef57ced595e046168eb73 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Fri, 7 Aug 2026 11:21:55 -0700 Subject: [PATCH 1/7] fix(streaming): keep n>1 choices separate in stream_chunk_builder build_base_response emits a single hardcoded choice and every assembly step reads choices[0], while get_combined_content joins the delta content of every choice of every chunk into one string. A streamed n=2 call therefore came back as one choice holding both completions concatenated, which is text no model produced, plus whichever finish_reason arrived last. The reassembled response is what feeds cost and token accounting, what the caller gets with complete_response=True, and what goes into the response cache, so a cached streamed n=2 call later replays the concatenation. Assemble each choice index on its own by recursing over the chunks filtered to that index, then swap the results in after usage has been calculated. Usage is still computed over the full chunk list, so it is unchanged. A single-choice stream sees no behaviour change: the recursion only runs when the chunks carry more than one distinct index. --- litellm/main.py | 74 +++++++++++++++++++++++++++++++++ tests/test_litellm/test_main.py | 56 +++++++++++++++++++++++++ 2 files changed, 130 insertions(+) diff --git a/litellm/main.py b/litellm/main.py index c70a41c891a..2ddded81a23 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8439,6 +8439,67 @@ def stream_chunk_builder_text_completion(chunks: list, messages: list | None = N return TextCompletionResponse(**response) +def _streaming_choice_index(choice: Any) -> int | None: + index = choice.get("index") if isinstance(choice, dict) else getattr(choice, "index", None) + return index if isinstance(index, int) else None + + +def _streaming_chunk_choices(chunk: Any) -> list[Any]: + choices = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) + return choices or [] + + +def _distinct_streaming_choice_indices(chunks: list[Any]) -> list[int]: + indices: Final = { + _streaming_choice_index(choice) for chunk in chunks for choice in _streaming_chunk_choices(chunk) + } + return sorted(index for index in indices if index is not None) + + +def _replace_chunk_choices(chunk: Any, choices: list[Any]) -> Any: + if not choices: + return chunk + if isinstance(chunk, dict): + return {**chunk, "choices": choices} + return chunk.model_copy(update={"choices": choices}) + + +def _chunks_for_streaming_choice_index(chunks: list[Any], index: int) -> list[Any]: + matched: Final = ( + (chunk, [c for c in _streaming_chunk_choices(chunk) if _streaming_choice_index(c) == index]) + for chunk in chunks + ) + return [ + _replace_chunk_choices(chunk, choices) + for chunk, choices in matched + if choices or not _streaming_chunk_choices(chunk) + ] + + +def _build_streaming_choices_per_index( + chunks: list[Any], + indices: list[int], + messages: list | None, + start_time, + end_time, +) -> list[Choices] | None: + built: Final = [ + stream_chunk_builder( + chunks=_chunks_for_streaming_choice_index(chunks, index), + messages=messages, + start_time=start_time, + end_time=end_time, + ) + for index in indices + ] + if any(not isinstance(response, ModelResponse) or not response.choices for response in built): + return None + return [ + cast(Choices, cast(ModelResponse, response).choices[0]).model_copy(update={"index": index}) + for index, response in zip(indices, built) + ] + + def stream_chunk_builder( chunks: list, messages: list | None = None, @@ -8470,6 +8531,13 @@ def stream_chunk_builder( ): # route to the text completion logic return stream_chunk_builder_text_completion(chunks=chunks, messages=messages) + choice_indices: Final = _distinct_streaming_choice_indices(chunks) + per_choice: Final = ( + _build_streaming_choices_per_index(chunks, choice_indices, messages, start_time, end_time) + if len(choice_indices) > 1 + else None + ) + model: Final = chunks[0]["model"] # Initialize the response dictionary response: Final = processor.build_base_response(chunks) @@ -8521,6 +8589,9 @@ def stream_chunk_builder( ) setattr(response, "usage", usage) + if per_choice is not None: + response.choices = per_choice + # Propagate provider_specific_fields from chunk hidden params when present. for chunk in reversed(chunks): if isinstance(chunk, dict): @@ -8694,6 +8765,9 @@ def stream_chunk_builder( setattr(response, "usage", usage) + if per_choice is not None: + response.choices = per_choice + # Propagate provider_specific_fields from the last chunk (contains provider # metadata like traffic_type set during streaming) for chunk in reversed(chunks): diff --git a/tests/test_litellm/test_main.py b/tests/test_litellm/test_main.py index 9e160370048..897bde73bc7 100644 --- a/tests/test_litellm/test_main.py +++ b/tests/test_litellm/test_main.py @@ -1449,6 +1449,62 @@ async def test_async_mock_delay(): assert delay >= 0.01 +def _n_choice_stream_chunk(index, content, finish_reason=None, role=None): + from litellm.types.utils import Delta, ModelResponseStream, StreamingChoices + + delta = {"content": content} + if role is not None: + delta["role"] = role + return ModelResponseStream( + id="chatcmpl-n2", + created=1, + model="gpt-4o", + object="chat.completion.chunk", + choices=[StreamingChoices(index=index, delta=Delta(**delta), finish_reason=finish_reason)], + ) + + +def test_stream_chunk_builder_keeps_n_greater_than_one_choices_separate(): + from litellm import stream_chunk_builder + + chunks = [ + _n_choice_stream_chunk(0, "AAA", role="assistant"), + _n_choice_stream_chunk(1, "BBB", role="assistant"), + _n_choice_stream_chunk(0, " aaa"), + _n_choice_stream_chunk(1, " bbb"), + _n_choice_stream_chunk(0, None, finish_reason="stop"), + _n_choice_stream_chunk(1, None, finish_reason="length"), + ] + + response = stream_chunk_builder(chunks, messages=[{"role": "user", "content": "hi"}]) + + assert response is not None + assert len(response.choices) == 2 + assert [choice.index for choice in response.choices] == [0, 1] + assert response.choices[0].message.content == "AAA aaa" + assert response.choices[1].message.content == "BBB bbb" + assert response.choices[0].finish_reason == "stop" + assert response.choices[1].finish_reason == "length" + + +def test_stream_chunk_builder_single_choice_is_unchanged(): + from litellm import stream_chunk_builder + + chunks = [ + _n_choice_stream_chunk(0, "AAA", role="assistant"), + _n_choice_stream_chunk(0, " aaa"), + _n_choice_stream_chunk(0, None, finish_reason="stop"), + ] + + response = stream_chunk_builder(chunks, messages=[{"role": "user", "content": "hi"}]) + + assert response is not None + assert len(response.choices) == 1 + assert response.choices[0].index == 0 + assert response.choices[0].message.content == "AAA aaa" + assert response.choices[0].finish_reason == "stop" + + def test_stream_chunk_builder_thinking_blocks(): from litellm import stream_chunk_builder from litellm.types.utils import Delta, ModelResponseStream, StreamingChoices From c62e6f3fb92299d88ffc16e3b76954c43202ffc2 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Tue, 25 Aug 2026 12:43:52 -0700 Subject: [PATCH 2/7] chore(streaming): give the per-index helpers real types The strict gate flagged the new helpers for ANN401 (bare `Any` on `choice`, `chunk` and the `_replace_chunk_choices` return) and ANN001 (`start_time` and `end_time` unannotated). Each one has a concrete type: the chunk and choice helpers already branch on `isinstance(..., dict)`, so the union of the pydantic model and a dict is what they actually take, and the two times are the `datetime.datetime | None` pair threaded into stream_chunk_builder. That also puts the already-imported StreamingChoices to use, clearing an F401. Net ruff difference on the file is 7 findings removed and none added. --- litellm/main.py | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index 2ddded81a23..a333a2bafd5 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8439,12 +8439,12 @@ def stream_chunk_builder_text_completion(chunks: list, messages: list | None = N return TextCompletionResponse(**response) -def _streaming_choice_index(choice: Any) -> int | None: +def _streaming_choice_index(choice: StreamingChoices | dict[str, Any]) -> int | None: index = choice.get("index") if isinstance(choice, dict) else getattr(choice, "index", None) return index if isinstance(index, int) else None -def _streaming_chunk_choices(chunk: Any) -> list[Any]: +def _streaming_chunk_choices(chunk: ModelResponseStream | dict[str, Any]) -> list[Any]: choices = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) return choices or [] @@ -8456,7 +8456,9 @@ def _distinct_streaming_choice_indices(chunks: list[Any]) -> list[int]: return sorted(index for index in indices if index is not None) -def _replace_chunk_choices(chunk: Any, choices: list[Any]) -> Any: +def _replace_chunk_choices( + chunk: ModelResponseStream | dict[str, Any], choices: list[Any] +) -> ModelResponseStream | dict[str, Any]: if not choices: return chunk if isinstance(chunk, dict): @@ -8480,8 +8482,8 @@ def _build_streaming_choices_per_index( chunks: list[Any], indices: list[int], messages: list | None, - start_time, - end_time, + start_time: datetime.datetime | None, + end_time: datetime.datetime | None, ) -> list[Choices] | None: built: Final = [ stream_chunk_builder( From 8fe54628b40a652f852fe0ce1479cfb9b8daf57e Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Tue, 25 Aug 2026 13:06:16 -0700 Subject: [PATCH 3/7] chore(streaming): use read-only collection types in the per-index helpers Giving the helpers concrete types traded ANN401 for LIT001, since `dict[...]` and `list[...]` are mutable collections in an annotation. Eleven of the fourteen are genuinely read-only and become `Mapping` / `Sequence`. `_replace_chunk_choices` narrowed with `isinstance(chunk, dict)`, which does not narrow a `Mapping`, so it now narrows on `ModelResponseStream` and the mapping branch stays a mapping. Same two expected inputs, same result for each. The remaining three have to be real lists and say why: two are forwarded to `stream_chunk_builder`, whose `chunks` and `messages` parameters are lists, and the third is assigned to `ModelResponse.choices`. `tests/test_litellm/test_main.py` is 98 passed before and 98 passed after. --- litellm/main.py | 29 ++++++++++++++++------------- 1 file changed, 16 insertions(+), 13 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index d80f5e3fa26..a372fb5c828 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8543,17 +8543,17 @@ def stream_chunk_builder_text_completion(chunks: list, messages: list | None = N return TextCompletionResponse(**response) -def _streaming_choice_index(choice: StreamingChoices | dict[str, Any]) -> int | None: +def _streaming_choice_index(choice: StreamingChoices | Mapping[str, Any]) -> int | None: index = choice.get("index") if isinstance(choice, dict) else getattr(choice, "index", None) return index if isinstance(index, int) else None -def _streaming_chunk_choices(chunk: ModelResponseStream | dict[str, Any]) -> list[Any]: +def _streaming_chunk_choices(chunk: ModelResponseStream | Mapping[str, Any]) -> Sequence[Any]: choices = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) return choices or [] -def _distinct_streaming_choice_indices(chunks: list[Any]) -> list[int]: +def _distinct_streaming_choice_indices(chunks: Sequence[Any]) -> Sequence[int]: indices: Final = { _streaming_choice_index(choice) for chunk in chunks for choice in _streaming_chunk_choices(chunk) } @@ -8561,16 +8561,19 @@ def _distinct_streaming_choice_indices(chunks: list[Any]) -> list[int]: def _replace_chunk_choices( - chunk: ModelResponseStream | dict[str, Any], choices: list[Any] -) -> ModelResponseStream | dict[str, Any]: + chunk: ModelResponseStream | Mapping[str, Any], choices: Sequence[Any] +) -> ModelResponseStream | Mapping[str, Any]: if not choices: return chunk - if isinstance(chunk, dict): - return {**chunk, "choices": choices} - return chunk.model_copy(update={"choices": choices}) + # Narrow on the model rather than on dict, so the mapping branch stays a Mapping. + if isinstance(chunk, ModelResponseStream): + return chunk.model_copy(update={"choices": list(choices)}) + return {**chunk, "choices": choices} -def _chunks_for_streaming_choice_index(chunks: list[Any], index: int) -> list[Any]: +def _chunks_for_streaming_choice_index( + chunks: Sequence[Any], index: int +) -> list[Any]: # mutable-ok: fed straight to stream_chunk_builder, whose `chunks` parameter is a list matched: Final = ( (chunk, [c for c in _streaming_chunk_choices(chunk) if _streaming_choice_index(c) == index]) for chunk in chunks @@ -8583,12 +8586,12 @@ def _chunks_for_streaming_choice_index(chunks: list[Any], index: int) -> list[An def _build_streaming_choices_per_index( - chunks: list[Any], - indices: list[int], - messages: list | None, + chunks: Sequence[Any], + indices: Sequence[int], + messages: list | None, # mutable-ok: forwarded verbatim to stream_chunk_builder, whose `messages` is a list start_time: datetime.datetime | None, end_time: datetime.datetime | None, -) -> list[Choices] | None: +) -> list[Choices] | None: # mutable-ok: assigned to ModelResponse.choices, which is a list built: Final = [ stream_chunk_builder( chunks=_chunks_for_streaming_choice_index(chunks, index), From be1fe001a99fe6ac87a38f8e9164f28a4b4053ea Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Tue, 25 Aug 2026 13:13:34 -0700 Subject: [PATCH 4/7] chore(streaming): give the nine collection constructions their reasons LIT002 covers construction as well as annotation, so the lists, the local set and the two mapping literals in the per-index helpers each carry a `# mutable-ok:` reason now. Nothing moves at runtime: tests/test_litellm/test_main.py is 98 passed before and after. --- litellm/main.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index a372fb5c828..c0c04646db4 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8550,11 +8550,11 @@ def _streaming_choice_index(choice: StreamingChoices | Mapping[str, Any]) -> int def _streaming_chunk_choices(chunk: ModelResponseStream | Mapping[str, Any]) -> Sequence[Any]: choices = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) - return choices or [] + return choices or [] # mutable-ok: the caller iterates it; an empty list is the natural absent value def _distinct_streaming_choice_indices(chunks: Sequence[Any]) -> Sequence[int]: - indices: Final = { + indices: Final = { # mutable-ok: a local set, discarded after sorting _streaming_choice_index(choice) for chunk in chunks for choice in _streaming_chunk_choices(chunk) } return sorted(index for index in indices if index is not None) @@ -8567,18 +8567,18 @@ def _replace_chunk_choices( return chunk # Narrow on the model rather than on dict, so the mapping branch stays a Mapping. if isinstance(chunk, ModelResponseStream): - return chunk.model_copy(update={"choices": list(choices)}) - return {**chunk, "choices": choices} + return chunk.model_copy(update={"choices": list(choices)}) # mutable-ok: model_copy stores what it is given, and ModelResponse.choices is a list + return {**chunk, "choices": choices} # mutable-ok: a new mapping so the caller's chunk is not mutated def _chunks_for_streaming_choice_index( chunks: Sequence[Any], index: int ) -> list[Any]: # mutable-ok: fed straight to stream_chunk_builder, whose `chunks` parameter is a list matched: Final = ( - (chunk, [c for c in _streaming_chunk_choices(chunk) if _streaming_choice_index(c) == index]) + (chunk, [c for c in _streaming_chunk_choices(chunk) if _streaming_choice_index(c) == index]) # mutable-ok: per-chunk filtered choices, consumed by the comprehension below for chunk in chunks ) - return [ + return [ # mutable-ok: stream_chunk_builder takes chunks as a list _replace_chunk_choices(chunk, choices) for chunk, choices in matched if choices or not _streaming_chunk_choices(chunk) @@ -8592,7 +8592,7 @@ def _build_streaming_choices_per_index( start_time: datetime.datetime | None, end_time: datetime.datetime | None, ) -> list[Choices] | None: # mutable-ok: assigned to ModelResponse.choices, which is a list - built: Final = [ + built: Final = [ # mutable-ok: one built response per choice index, zipped below stream_chunk_builder( chunks=_chunks_for_streaming_choice_index(chunks, index), messages=messages, @@ -8603,8 +8603,8 @@ def _build_streaming_choices_per_index( ] if any(not isinstance(response, ModelResponse) or not response.choices for response in built): return None - return [ - cast(Choices, cast(ModelResponse, response).choices[0]).model_copy(update={"index": index}) + return [ # mutable-ok: assigned to ModelResponse.choices, which is a list + cast(Choices, cast(ModelResponse, response).choices[0]).model_copy(update={"index": index}) # mutable-ok: model_copy stores what it is given for index, response in zip(indices, built) ] From 8441296f703c501b5b33052ffae11865b3c7d952 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Thu, 27 Aug 2026 09:47:55 -0700 Subject: [PATCH 5/7] fix(streaming): narrow the built choices with isinstance instead of cast The per-index builder asserted its result with two nested cast() calls, which the repo's LIT006 budget counts as unchecked assertions. Move the narrowing into a small helper that isinstance-checks the response and its first choice and returns None when either is not what the merge needs, so the caller's fallback to the single-response path is reached by a real check rather than by faith. The two one-shot locals in the index/choices readers pick up the Final annotation LIT010 asks for. No behaviour change: the same inputs produce the same choices, and the same inputs fall back. Signed-off-by: Vineeth Sai --- litellm/main.py | 43 ++++++++++++++++++++++++++++++------------- 1 file changed, 30 insertions(+), 13 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index 4cfab4d6db2..4d8708752fe 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8588,12 +8588,12 @@ def stream_chunk_builder_text_completion(chunks: list, messages: list | None = N def _streaming_choice_index(choice: StreamingChoices | Mapping[str, Any]) -> int | None: - index = choice.get("index") if isinstance(choice, dict) else getattr(choice, "index", None) + index: Final = choice.get("index") if isinstance(choice, dict) else getattr(choice, "index", None) return index if isinstance(index, int) else None def _streaming_chunk_choices(chunk: ModelResponseStream | Mapping[str, Any]) -> Sequence[Any]: - choices = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) + choices: Final = chunk.get("choices") if isinstance(chunk, dict) else getattr(chunk, "choices", None) return choices or [] # mutable-ok: the caller iterates it; an empty list is the natural absent value @@ -8629,6 +8629,21 @@ def _chunks_for_streaming_choice_index( ] +def _renumbered_first_choice(response: Any, index: int) -> Choices | None: + """The single choice a per-index build produced, renumbered to that index. + + None when the build did not yield a usable non-streaming choice: `stream_chunk_builder` + is typed to admit a streaming response and a streaming choice, and either one leaves + nothing to merge. + """ + if not isinstance(response, ModelResponse) or not response.choices: + return None + first: Final = response.choices[0] + if not isinstance(first, Choices): + return None + return first.model_copy(update={"index": index}) # mutable-ok: model_copy stores what it is given + + def _build_streaming_choices_per_index( chunks: Sequence[Any], indices: Sequence[int], @@ -8636,21 +8651,23 @@ def _build_streaming_choices_per_index( start_time: datetime.datetime | None, end_time: datetime.datetime | None, ) -> list[Choices] | None: # mutable-ok: assigned to ModelResponse.choices, which is a list - built: Final = [ # mutable-ok: one built response per choice index, zipped below - stream_chunk_builder( - chunks=_chunks_for_streaming_choice_index(chunks, index), - messages=messages, - start_time=start_time, - end_time=end_time, + built: Final = [ # mutable-ok: one choice per index, returned as ModelResponse.choices + _renumbered_first_choice( + stream_chunk_builder( + chunks=_chunks_for_streaming_choice_index(chunks, index), + messages=messages, + start_time=start_time, + end_time=end_time, + ), + index, ) for index in indices ] - if any(not isinstance(response, ModelResponse) or not response.choices for response in built): + # One unusable build means there is nothing to merge, so the caller falls back to + # the single-response path rather than being handed a partial list. + if any(choice is None for choice in built): return None - return [ # mutable-ok: assigned to ModelResponse.choices, which is a list - cast(Choices, cast(ModelResponse, response).choices[0]).model_copy(update={"index": index}) # mutable-ok: model_copy stores what it is given - for index, response in zip(indices, built) - ] + return [choice for choice in built if choice is not None] # mutable-ok: assigned to ModelResponse.choices, which is a list def stream_chunk_builder( From e1b0a578e30c4cc4239511fef01669482c3bd2b1 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Thu, 27 Aug 2026 10:25:36 -0700 Subject: [PATCH 6/7] fix(streaming): annotate the narrowing helper with object rather than Any The helper isinstance-narrows its argument, so it does not need a dynamically typed annotation; object works and keeps the repo's ANN401 budget unchanged. Signed-off-by: Vineeth Sai --- litellm/main.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/main.py b/litellm/main.py index 4d8708752fe..8fab6dbf0d0 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8629,7 +8629,7 @@ def _chunks_for_streaming_choice_index( ] -def _renumbered_first_choice(response: Any, index: int) -> Choices | None: +def _renumbered_first_choice(response: object, index: int) -> Choices | None: """The single choice a per-index build produced, renumbered to that index. None when the build did not yield a usable non-streaming choice: `stream_chunk_builder` From 4dd061f29ed979333cda8b453fe0f44160ddd150 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Thu, 27 Aug 2026 11:04:08 -0700 Subject: [PATCH 7/7] fix(streaming): drop the redundant choice narrowing ModelResponse.choices is declared list[Choices], so the second isinstance was an unnecessary check basedpyright counts against the budget; only the response itself can be the streaming variant. Signed-off-by: Vineeth Sai --- litellm/main.py | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/litellm/main.py b/litellm/main.py index 8fab6dbf0d0..cc131f6947f 100644 --- a/litellm/main.py +++ b/litellm/main.py @@ -8632,16 +8632,14 @@ def _chunks_for_streaming_choice_index( def _renumbered_first_choice(response: object, index: int) -> Choices | None: """The single choice a per-index build produced, renumbered to that index. - None when the build did not yield a usable non-streaming choice: `stream_chunk_builder` - is typed to admit a streaming response and a streaming choice, and either one leaves - nothing to merge. + None when the build did not yield a usable non-streaming response: `stream_chunk_builder` + is typed to admit a streaming one, which leaves nothing to merge. """ if not isinstance(response, ModelResponse) or not response.choices: return None - first: Final = response.choices[0] - if not isinstance(first, Choices): - return None - return first.model_copy(update={"index": index}) # mutable-ok: model_copy stores what it is given + # ModelResponse.choices is declared list[Choices], so the first entry needs no + # further narrowing; only the response itself can be the streaming variant. + return response.choices[0].model_copy(update={"index": index}) # mutable-ok: model_copy stores what it is given def _build_streaming_choices_per_index(