From a2720f0e1aaec0ae8d8d31c6d67edc1086628385 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 20 May 2026 21:07:24 +0000 Subject: [PATCH] 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 --- .../transformation.py | 99 +++---------------- .../llms/chatgpt/responses/transformation.py | 95 ++---------------- litellm/responses/sse_output_recovery.py | 99 +++++++++++++++++++ 3 files changed, 119 insertions(+), 174 deletions(-) create mode 100644 litellm/responses/sse_output_recovery.py diff --git a/litellm/completion_extras/litellm_responses_transformation/transformation.py b/litellm/completion_extras/litellm_responses_transformation/transformation.py index 8071451c9f6..caee7358abb 100644 --- a/litellm/completion_extras/litellm_responses_transformation/transformation.py +++ b/litellm/completion_extras/litellm_responses_transformation/transformation.py @@ -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 diff --git a/litellm/llms/chatgpt/responses/transformation.py b/litellm/llms/chatgpt/responses/transformation.py index e35af48ef54..5ef64fd6cd9 100644 --- a/litellm/llms/chatgpt/responses/transformation.py +++ b/litellm/llms/chatgpt/responses/transformation.py @@ -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]: diff --git a/litellm/responses/sse_output_recovery.py b/litellm/responses/sse_output_recovery.py new file mode 100644 index 00000000000..70f470660ca --- /dev/null +++ b/litellm/responses/sse_output_recovery.py @@ -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", [])