mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-07 02:59:05 +00:00
refactor(responses): extract shared SSE output-item recovery helpers
Both ChatGPTResponsesAPIConfig and LiteLLMResponsesTransformationHandler duplicated the same OUTPUT_ITEM_DONE / OUTPUT_TEXT_DONE recovery algorithm. Move that logic into litellm.responses.sse_output_recovery and have both call sites use the shared helpers, so future fixes apply in one place. Co-authored-by: Yassin Kortam <yassin@berri.ai>
This commit is contained in:
parent
1bc0cf375a
commit
a2720f0e1a
3 changed files with 119 additions and 174 deletions
|
|
@ -31,6 +31,10 @@ from litellm.llms.base_llm.base_model_iterator import BaseModelResponseIterator
|
|||
from litellm.llms.base_llm.bridges.completion_transformation import (
|
||||
CompletionTransformationBridge,
|
||||
)
|
||||
from litellm.responses.sse_output_recovery import (
|
||||
record_output_item_chunk,
|
||||
record_output_text_chunk,
|
||||
)
|
||||
from litellm.types.llms.openai import (
|
||||
ChatCompletionAnnotation,
|
||||
ChatCompletionReasoningItem,
|
||||
|
|
@ -633,90 +637,6 @@ class LiteLLMResponsesTransformationHandler(CompletionTransformationBridge):
|
|||
return None
|
||||
return cast(List[Dict[str, Any]], response_output)
|
||||
|
||||
@classmethod
|
||||
def _update_recovered_output_items(
|
||||
cls,
|
||||
parsed_chunk: Dict[str, Any],
|
||||
recovered_output_items: Dict[int, Dict[str, Any]],
|
||||
) -> None:
|
||||
item = parsed_chunk.get("item")
|
||||
if not isinstance(item, dict):
|
||||
return
|
||||
try:
|
||||
output_index_raw = parsed_chunk.get("output_index")
|
||||
if output_index_raw is None:
|
||||
raise ValueError("missing output_index")
|
||||
output_index = int(output_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
output_index = len(recovered_output_items)
|
||||
recovered_output_items[output_index] = item
|
||||
|
||||
@classmethod
|
||||
def _update_recovered_text_only_items(
|
||||
cls,
|
||||
parsed_chunk: Dict[str, Any],
|
||||
recovered_output_items: Dict[int, Dict[str, Any]],
|
||||
recovered_text_only_items: Dict[int, Dict[str, Any]],
|
||||
) -> None:
|
||||
text = parsed_chunk.get("text")
|
||||
if not isinstance(text, str):
|
||||
return
|
||||
|
||||
try:
|
||||
output_index_raw = parsed_chunk.get("output_index")
|
||||
if output_index_raw is None:
|
||||
raise ValueError("missing output_index")
|
||||
output_index = int(output_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
output_index = len(recovered_text_only_items)
|
||||
|
||||
if output_index in recovered_output_items:
|
||||
return
|
||||
|
||||
item = recovered_text_only_items.get(output_index)
|
||||
if item is None:
|
||||
item = {
|
||||
"type": "message",
|
||||
"id": parsed_chunk.get("item_id") or f"msg_{output_index}",
|
||||
"role": "assistant",
|
||||
"status": "completed",
|
||||
"content": [],
|
||||
}
|
||||
recovered_text_only_items[output_index] = item
|
||||
|
||||
content = item.setdefault("content", [])
|
||||
if not isinstance(content, list):
|
||||
return
|
||||
|
||||
try:
|
||||
content_index_raw = parsed_chunk.get("content_index")
|
||||
if content_index_raw is None:
|
||||
raise ValueError("missing content_index")
|
||||
content_index = int(content_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
content_index = len(content)
|
||||
|
||||
while len(content) <= content_index:
|
||||
content.append(
|
||||
{
|
||||
"type": "output_text",
|
||||
"text": "",
|
||||
"annotations": [],
|
||||
}
|
||||
)
|
||||
|
||||
content_item = content[content_index]
|
||||
if not isinstance(content_item, dict):
|
||||
content_item = {}
|
||||
content[content_index] = content_item
|
||||
|
||||
content_item["type"] = "output_text"
|
||||
content_item["text"] = text
|
||||
if parsed_chunk.get("annotations") is not None:
|
||||
content_item["annotations"] = parsed_chunk["annotations"]
|
||||
else:
|
||||
content_item.setdefault("annotations", [])
|
||||
|
||||
@classmethod
|
||||
def _recover_output_items_from_raw_sse(
|
||||
cls, raw_sse: Optional[str]
|
||||
|
|
@ -743,14 +663,17 @@ class LiteLLMResponsesTransformationHandler(CompletionTransformationBridge):
|
|||
continue
|
||||
|
||||
if event_type == ResponsesAPIStreamEvents.OUTPUT_ITEM_DONE:
|
||||
cls._update_recovered_output_items(parsed_chunk, recovered_output_items)
|
||||
record_output_item_chunk(
|
||||
parsed_chunk=parsed_chunk,
|
||||
output_items=recovered_output_items,
|
||||
)
|
||||
continue
|
||||
|
||||
if event_type == ResponsesAPIStreamEvents.OUTPUT_TEXT_DONE:
|
||||
cls._update_recovered_text_only_items(
|
||||
record_output_text_chunk(
|
||||
parsed_chunk=parsed_chunk,
|
||||
recovered_output_items=recovered_output_items,
|
||||
recovered_text_only_items=recovered_text_only_items,
|
||||
output_items=recovered_output_items,
|
||||
text_only_items=recovered_text_only_items,
|
||||
)
|
||||
|
||||
# Merge text-only items into the recovered output items. Real
|
||||
|
|
|
|||
|
|
@ -9,6 +9,10 @@ from litellm.litellm_core_utils.llm_response_utils.convert_dict_to_response impo
|
|||
)
|
||||
from litellm.llms.openai.common_utils import OpenAIError
|
||||
from litellm.llms.openai.responses.transformation import OpenAIResponsesAPIConfig
|
||||
from litellm.responses.sse_output_recovery import (
|
||||
record_output_item_chunk,
|
||||
record_output_text_chunk,
|
||||
)
|
||||
from litellm.types.llms.openai import (
|
||||
ResponsesAPIResponse,
|
||||
ResponsesAPIStreamEvents,
|
||||
|
|
@ -168,17 +172,17 @@ class ChatGPTResponsesAPIConfig(OpenAIResponsesAPIConfig):
|
|||
|
||||
event_type = parsed_chunk.get("type")
|
||||
if event_type == ResponsesAPIStreamEvents.OUTPUT_ITEM_DONE:
|
||||
self._record_output_item_chunk(
|
||||
record_output_item_chunk(
|
||||
parsed_chunk=parsed_chunk,
|
||||
streamed_output_items=streamed_output_items,
|
||||
output_items=streamed_output_items,
|
||||
)
|
||||
continue
|
||||
|
||||
if event_type == ResponsesAPIStreamEvents.OUTPUT_TEXT_DONE:
|
||||
self._record_output_text_chunk(
|
||||
record_output_text_chunk(
|
||||
parsed_chunk=parsed_chunk,
|
||||
streamed_output_items=streamed_output_items,
|
||||
text_only_output_items=text_only_output_items,
|
||||
output_items=streamed_output_items,
|
||||
text_only_items=text_only_output_items,
|
||||
)
|
||||
continue
|
||||
|
||||
|
|
@ -224,87 +228,6 @@ class ChatGPTResponsesAPIConfig(OpenAIResponsesAPIConfig):
|
|||
return None
|
||||
return parsed_chunk
|
||||
|
||||
def _record_output_item_chunk(
|
||||
self, parsed_chunk: Dict[str, Any], streamed_output_items: Dict[int, dict]
|
||||
) -> None:
|
||||
item = parsed_chunk.get("item")
|
||||
output_index = parsed_chunk.get("output_index")
|
||||
if not isinstance(item, dict):
|
||||
return
|
||||
try:
|
||||
if output_index is None:
|
||||
raise ValueError("missing output_index")
|
||||
index = int(output_index)
|
||||
except (TypeError, ValueError):
|
||||
index = len(streamed_output_items)
|
||||
streamed_output_items[index] = item
|
||||
|
||||
def _record_output_text_chunk(
|
||||
self,
|
||||
parsed_chunk: Dict[str, Any],
|
||||
streamed_output_items: Dict[int, dict],
|
||||
text_only_output_items: Dict[int, dict],
|
||||
) -> None:
|
||||
text = parsed_chunk.get("text")
|
||||
if not isinstance(text, str):
|
||||
return
|
||||
|
||||
try:
|
||||
output_index_raw = parsed_chunk.get("output_index")
|
||||
if output_index_raw is None:
|
||||
raise ValueError("missing output_index")
|
||||
output_index = int(output_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
output_index = len(text_only_output_items)
|
||||
|
||||
# If a real OUTPUT_ITEM_DONE already covered this index, prefer it.
|
||||
if output_index in streamed_output_items:
|
||||
return
|
||||
|
||||
item = text_only_output_items.get(output_index)
|
||||
if item is None:
|
||||
item = {
|
||||
"type": "message",
|
||||
"id": parsed_chunk.get("item_id") or f"msg_{output_index}",
|
||||
"role": "assistant",
|
||||
"status": "completed",
|
||||
"content": [],
|
||||
}
|
||||
text_only_output_items[output_index] = item
|
||||
|
||||
content = item.setdefault("content", [])
|
||||
if not isinstance(content, list):
|
||||
return
|
||||
|
||||
try:
|
||||
content_index_raw = parsed_chunk.get("content_index")
|
||||
if content_index_raw is None:
|
||||
raise ValueError("missing content_index")
|
||||
content_index = int(content_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
content_index = len(content)
|
||||
|
||||
while len(content) <= content_index:
|
||||
content.append(
|
||||
{
|
||||
"type": "output_text",
|
||||
"text": "",
|
||||
"annotations": [],
|
||||
}
|
||||
)
|
||||
|
||||
content_item = content[content_index]
|
||||
if not isinstance(content_item, dict):
|
||||
content_item = {}
|
||||
content[content_index] = content_item
|
||||
|
||||
content_item["type"] = "output_text"
|
||||
content_item["text"] = text
|
||||
if parsed_chunk.get("annotations") is not None:
|
||||
content_item["annotations"] = parsed_chunk["annotations"]
|
||||
else:
|
||||
content_item.setdefault("annotations", [])
|
||||
|
||||
def _build_completed_response_from_chunk(
|
||||
self, parsed_chunk: Dict[str, Any], streamed_output_items: Dict[int, dict]
|
||||
) -> Optional[ResponsesAPIResponse]:
|
||||
|
|
|
|||
99
litellm/responses/sse_output_recovery.py
Normal file
99
litellm/responses/sse_output_recovery.py
Normal file
|
|
@ -0,0 +1,99 @@
|
|||
"""
|
||||
Shared helpers for recovering Responses API output items from raw SSE chunks.
|
||||
|
||||
The same recovery logic is needed in multiple places (e.g. the ChatGPT
|
||||
Responses transformation and the LiteLLM Responses-to-Chat-Completions
|
||||
bridge). Keep the implementation in a single module so a fix in one
|
||||
caller automatically applies to all of them.
|
||||
"""
|
||||
|
||||
from typing import Any, Dict
|
||||
|
||||
|
||||
def record_output_item_chunk(
|
||||
parsed_chunk: Dict[str, Any],
|
||||
output_items: Dict[int, Dict[str, Any]],
|
||||
) -> None:
|
||||
"""Record an OUTPUT_ITEM_DONE chunk into ``output_items`` keyed by
|
||||
``output_index`` (falling back to the next free slot when missing).
|
||||
"""
|
||||
item = parsed_chunk.get("item")
|
||||
if not isinstance(item, dict):
|
||||
return
|
||||
try:
|
||||
output_index_raw = parsed_chunk.get("output_index")
|
||||
if output_index_raw is None:
|
||||
raise ValueError("missing output_index")
|
||||
output_index = int(output_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
output_index = len(output_items)
|
||||
output_items[output_index] = item
|
||||
|
||||
|
||||
def record_output_text_chunk(
|
||||
parsed_chunk: Dict[str, Any],
|
||||
output_items: Dict[int, Dict[str, Any]],
|
||||
text_only_items: Dict[int, Dict[str, Any]],
|
||||
) -> None:
|
||||
"""Record an OUTPUT_TEXT_DONE chunk as a synthetic message item in
|
||||
``text_only_items``. Real OUTPUT_ITEM_DONE events already captured in
|
||||
``output_items`` take precedence at the same ``output_index``.
|
||||
"""
|
||||
text = parsed_chunk.get("text")
|
||||
if not isinstance(text, str):
|
||||
return
|
||||
|
||||
try:
|
||||
output_index_raw = parsed_chunk.get("output_index")
|
||||
if output_index_raw is None:
|
||||
raise ValueError("missing output_index")
|
||||
output_index = int(output_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
output_index = len(text_only_items)
|
||||
|
||||
if output_index in output_items:
|
||||
return
|
||||
|
||||
item = text_only_items.get(output_index)
|
||||
if item is None:
|
||||
item = {
|
||||
"type": "message",
|
||||
"id": parsed_chunk.get("item_id") or f"msg_{output_index}",
|
||||
"role": "assistant",
|
||||
"status": "completed",
|
||||
"content": [],
|
||||
}
|
||||
text_only_items[output_index] = item
|
||||
|
||||
content = item.setdefault("content", [])
|
||||
if not isinstance(content, list):
|
||||
return
|
||||
|
||||
try:
|
||||
content_index_raw = parsed_chunk.get("content_index")
|
||||
if content_index_raw is None:
|
||||
raise ValueError("missing content_index")
|
||||
content_index = int(content_index_raw)
|
||||
except (TypeError, ValueError):
|
||||
content_index = len(content)
|
||||
|
||||
while len(content) <= content_index:
|
||||
content.append(
|
||||
{
|
||||
"type": "output_text",
|
||||
"text": "",
|
||||
"annotations": [],
|
||||
}
|
||||
)
|
||||
|
||||
content_item = content[content_index]
|
||||
if not isinstance(content_item, dict):
|
||||
content_item = {}
|
||||
content[content_index] = content_item
|
||||
|
||||
content_item["type"] = "output_text"
|
||||
content_item["text"] = text
|
||||
if parsed_chunk.get("annotations") is not None:
|
||||
content_item["annotations"] = parsed_chunk["annotations"]
|
||||
else:
|
||||
content_item.setdefault("annotations", [])
|
||||
Loading…
Add table
Reference in a new issue