Create streaming response, checking if the first chunk is an error. If the first chunk is an error, return a standard JSON error response. Otherwise, return StreamingResponse and stream all content.
(
generator: AsyncGenerator[str, None],
media_type: str,
headers: dict,
default_status_code: int = status.HTTP_200_OK,
request: Request | None = None,
)
| 608 | |
| 609 | |
| 610 | async def create_response( |
| 611 | generator: AsyncGenerator[str, None], |
| 612 | media_type: str, |
| 613 | headers: dict, |
| 614 | default_status_code: int = status.HTTP_200_OK, |
| 615 | request: Request | None = None, |
| 616 | ) -> StreamingResponse | JSONResponse: |
| 617 | """ |
| 618 | Create streaming response, checking if the first chunk is an error. |
| 619 | If the first chunk is an error, return a standard JSON error response. |
| 620 | Otherwise, return StreamingResponse and stream all content. |
| 621 | """ |
| 622 | # Tell buffering reverse proxies (nginx, ingress-nginx, Envoy) to flush SSE |
| 623 | # immediately instead of releasing the whole stream in one batch (issue #28384). |
| 624 | streaming_headers: Final = { |
| 625 | **headers, |
| 626 | "Cache-Control": "no-cache", |
| 627 | "X-Accel-Buffering": "no", |
| 628 | } |
| 629 | first_chunk_value: str | None = None |
| 630 | final_status_code = default_status_code |
| 631 | |
| 632 | try: |
| 633 | # Handle coroutine that returns a generator |
| 634 | if asyncio.iscoroutine(generator): |
| 635 | generator = await generator |
| 636 | |
| 637 | # Now get the first chunk from the actual generator |
| 638 | first_chunk_value = await _buffer_first_chunk_honoring_disconnect(generator, request) |
| 639 | |
| 640 | if first_chunk_value is not None: |
| 641 | try: |
| 642 | error_code_from_chunk: Final = await _parse_event_data_for_error(first_chunk_value) |
| 643 | if error_code_from_chunk is not None: |
| 644 | # First chunk is an error, stream hasn't really started yet |
| 645 | # Should return standard JSON error response instead of SSE format |
| 646 | final_status_code = error_code_from_chunk |
| 647 | verbose_proxy_logger.debug( |
| 648 | "Error detected in first stream chunk. Returning JSON error response with status code: %s", |
| 649 | final_status_code, |
| 650 | ) |
| 651 | |
| 652 | # Parse error content |
| 653 | error_dict: Final = _extract_error_from_sse_chunk(first_chunk_value) |
| 654 | |
| 655 | # Consume and close generator (avoid resource leak) |
| 656 | try: |
| 657 | await generator.aclose() |
| 658 | except Exception: |
| 659 | pass |
| 660 | |
| 661 | # Return JSON format error response |
| 662 | return JSONResponse( |
| 663 | status_code=final_status_code, |
| 664 | content={"error": error_dict}, |
| 665 | headers=headers, |
| 666 | ) |
| 667 | except Exception as e: |