Handle response.failed, response.incomplete, and response.cancelled (#23492)

* Handle response.failed, response.incomplete, and response.cancelled terminal events in background streaming

Previously the background streaming task only handled response.completed and
hardcoded the final status to "completed". This missed three other terminal
event types from the OpenAI streaming spec, causing failed/incomplete/cancelled
responses to be incorrectly marked as completed.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

* Remove unused terminal_response_data variable

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

* Address code review: derive fallback status from event type, rewrite tests as integration tests

1. Replace hardcoded "completed" fallback in response_data.get("status")
   with _event_to_status lookup so that response.incomplete and
   response.cancelled events get the correct fallback if the response
   body ever omits the status field.

2. Replace duplicated-logic unit tests with integration tests that
   exercise background_streaming_task directly using mocked streaming
   responses and assert on the final update_state call arguments.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

* Remove dead mock_processor and unused mock_response parameter from test helper

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

* Remove FastAPI and UserAPIKeyAuth imports from test file

These types were only used as Mock(spec=...) arguments. Drop the spec
constraints and remove the top-level imports to avoid pulling FastAPI
into test files outside litellm/proxy/.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

* Log warning when streaming response has no body_iterator

If base_process_llm_request returns a non-streaming response (no
body_iterator), log a warning since this likely indicates a
misconfiguration or provider error rather than a successful completion.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Committed-By-Agent: claude

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
xianzongxie-stripe
2026-03-13 23:02:09 -07:00
committed by GitHub
co-authored by Claude Opus 4.6
parent a94b961c18
commit 81474c17fe
2 changed files with 299 additions and 6 deletions
@@ -113,6 +113,16 @@ async def background_streaming_task( # noqa: PLR0915
last_update_time = asyncio.get_event_loop().time()
UPDATE_INTERVAL = 0.150 # 150ms batching interval
# Track the terminal event from the stream (may not be "completed")
terminal_status = None # Will be set by response.completed/failed/incomplete/cancelled
terminal_error = None
_event_to_status = {
"response.completed": "completed",
"response.failed": "failed",
"response.incomplete": "incomplete",
"response.cancelled": "cancelled",
}
async def flush_state_if_needed(force: bool = False) -> None:
"""Flush accumulated state to Redis if interval elapsed or forced"""
nonlocal state_dirty, last_update_time
@@ -131,6 +141,12 @@ async def background_streaming_task( # noqa: PLR0915
last_update_time = current_time
# Handle StreamingResponse
if not hasattr(response, "body_iterator"):
verbose_proxy_logger.warning(
f"background_streaming_task: response for {polling_id} has no "
"body_iterator; this may indicate a misconfiguration or provider error"
)
if hasattr(response, "body_iterator"):
async for chunk in response.body_iterator:
# Parse chunk
@@ -224,10 +240,23 @@ async def background_streaming_task( # noqa: PLR0915
status="in_progress",
)
elif event_type == "response.completed":
# Response completed - extract all ResponsesAPIResponse fields
# https://platform.openai.com/docs/api-reference/responses-streaming/response-completed
elif event_type in (
"response.completed",
"response.failed",
"response.incomplete",
"response.cancelled",
):
# Terminal event - extract all ResponsesAPIResponse fields
# https://platform.openai.com/docs/api-reference/responses-streaming
response_data = event.get("response", {})
terminal_status = response_data.get(
"status",
_event_to_status.get(event_type, "completed"),
)
# Extract error for failed responses
if event_type == "response.failed":
terminal_error = response_data.get("error")
# Core response fields
usage_data = response_data.get("usage")
@@ -278,11 +307,14 @@ async def background_streaming_task( # noqa: PLR0915
# Final flush to ensure all accumulated state is saved
await flush_state_if_needed(force=True)
# Mark as completed with all ResponsesAPIResponse fields
# Use the terminal status from the stream, default to "completed"
final_status = terminal_status or "completed"
await polling_handler.update_state(
polling_id=polling_id,
status="completed",
status=final_status,
usage=usage_data,
error=terminal_error,
reasoning=reasoning_data,
tool_choice=tool_choice_data,
tools=tools_data,
@@ -301,7 +333,7 @@ async def background_streaming_task( # noqa: PLR0915
)
verbose_proxy_logger.info(
f"Completed background streaming for {polling_id}, output_items={len(output_items)}"
f"Finished background streaming for {polling_id}, status={final_status}, output_items={len(output_items)}"
)
except Exception as e:
@@ -1318,6 +1318,267 @@ class TestStreamingEventParsing:
assert output_items["item_123"]["content"][0]["type"] == "text"
def _make_sse_stream(events: list) -> Mock:
"""Create a mock StreamingResponse with body_iterator from a list of event dicts."""
async def _body_iterator():
for event in events:
yield f"data: {json.dumps(event)}"
yield "data: [DONE]"
mock_response = Mock()
mock_response.body_iterator = _body_iterator()
return mock_response
def _make_background_streaming_kwargs(
polling_id: str,
polling_handler: ResponsePollingHandler,
) -> dict:
"""Build kwargs for background_streaming_task with all required mocks."""
return dict(
polling_id=polling_id,
data={"model": "gpt-4o", "stream": False, "background": True},
polling_handler=polling_handler,
request=Mock(),
fastapi_response=Mock(),
user_api_key_dict=Mock(),
general_settings={},
llm_router=None,
proxy_config=Mock(),
proxy_logging_obj=Mock(),
select_data_generator=Mock(),
user_model=None,
user_temperature=None,
user_request_timeout=None,
user_max_tokens=None,
user_api_base=None,
version=None,
)
@pytest.mark.xdist_group("heavy_imports")
class TestBackgroundStreamingTerminalEvents:
"""
Integration tests that exercise background_streaming_task with mocked
streaming responses, verifying the final update_state call for each
terminal event type.
"""
@pytest.mark.asyncio
async def test_response_failed_sets_failed_status_and_error(self):
"""Test that a response.failed stream event results in failed status with error"""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
error_payload = {
"type": "server_error",
"message": "The model encountered an error",
"code": "model_error",
}
events = [
{"type": "response.in_progress"},
{
"type": "response.failed",
"response": {
"id": "resp_123",
"status": "failed",
"error": error_payload,
"model": "gpt-4o",
"output": [],
},
},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_1", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
# Find the final update_state call (last one)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "failed"
assert final_call.kwargs["error"] == error_payload
@pytest.mark.asyncio
async def test_response_incomplete_sets_incomplete_status_and_details(self):
"""Test that a response.incomplete stream event results in incomplete status"""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
events = [
{"type": "response.in_progress"},
{
"type": "response.incomplete",
"response": {
"id": "resp_123",
"status": "incomplete",
"incomplete_details": {"reason": "max_output_tokens"},
"usage": {"input_tokens": 10, "output_tokens": 4096},
"model": "gpt-4o",
"output": [{"id": "item_1", "type": "message"}],
},
},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_2", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "incomplete"
assert final_call.kwargs["incomplete_details"] == {"reason": "max_output_tokens"}
assert final_call.kwargs["usage"] == {"input_tokens": 10, "output_tokens": 4096}
@pytest.mark.asyncio
async def test_response_cancelled_sets_cancelled_status(self):
"""Test that a response.cancelled stream event results in cancelled status"""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
events = [
{"type": "response.in_progress"},
{
"type": "response.cancelled",
"response": {
"id": "resp_123",
"status": "cancelled",
"model": "gpt-4o",
"output": [],
},
},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_3", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "cancelled"
@pytest.mark.asyncio
async def test_response_completed_sets_completed_status(self):
"""Test that a response.completed stream event results in completed status"""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
events = [
{"type": "response.in_progress"},
{
"type": "response.completed",
"response": {
"id": "resp_123",
"status": "completed",
"usage": {"input_tokens": 10, "output_tokens": 50},
"model": "gpt-4o",
"output": [{"id": "item_1", "type": "message"}],
},
},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_4", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "completed"
assert final_call.kwargs["usage"] == {"input_tokens": 10, "output_tokens": 50}
@pytest.mark.asyncio
async def test_fallback_status_derived_from_event_type_when_status_field_missing(self):
"""Test that when the response body lacks a status field, the fallback
is derived from the event type, not hardcoded to 'completed'."""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
# response.incomplete event with NO status field in the response body
events = [
{"type": "response.in_progress"},
{
"type": "response.incomplete",
"response": {
"id": "resp_123",
# "status" deliberately omitted
"incomplete_details": {"reason": "max_output_tokens"},
"model": "gpt-4o",
"output": [],
},
},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_5", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "incomplete"
@pytest.mark.asyncio
async def test_no_terminal_event_defaults_to_completed(self):
"""Test that when no terminal event is received, status defaults to completed"""
from litellm.proxy.response_polling.background_streaming import (
background_streaming_task,
)
# Stream with only in_progress, no terminal event
events = [
{"type": "response.in_progress"},
]
mock_response = _make_sse_stream(events)
handler = AsyncMock(spec=ResponsePollingHandler)
kwargs = _make_background_streaming_kwargs("poll_6", handler)
with patch(
"litellm.proxy.response_polling.background_streaming.ProxyBaseLLMRequestProcessing"
) as MockProcessor:
MockProcessor.return_value.base_process_llm_request = AsyncMock(
return_value=mock_response
)
await background_streaming_task(**kwargs)
final_call = handler.update_state.call_args_list[-1]
assert final_call.kwargs["status"] == "completed"
class TestEdgeCases:
"""Test edge cases and error scenarios"""