Merge pull request #3707 from pipecat-ai/aleix/fix-openai-stream-close-compat
fix(openai): use compatible stream closing for non-OpenAI providers
This commit is contained in:
1
changelog/3707.fixed.md
Normal file
1
changelog/3707.fixed.md
Normal file
@@ -0,0 +1 @@
|
|||||||
|
- Fixed stream closing compatibility for OpenAI-compatible providers (e.g. OpenPipe) that return async generators instead of `AsyncStream`.
|
||||||
@@ -9,6 +9,7 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
import base64
|
import base64
|
||||||
import json
|
import json
|
||||||
|
from contextlib import asynccontextmanager
|
||||||
from typing import Any, Dict, List, Mapping, Optional
|
from typing import Any, Dict, List, Mapping, Optional
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
@@ -374,9 +375,19 @@ class BaseOpenAILLMService(LLMService):
|
|||||||
else self._stream_chat_completions_universal_context(context)
|
else self._stream_chat_completions_universal_context(context)
|
||||||
)
|
)
|
||||||
|
|
||||||
# Use context manager to ensure stream is closed on cancellation/exception.
|
# Ensure stream is closed on cancellation/exception to prevent socket
|
||||||
# Without this, CancelledError during iteration leaves the underlying socket open.
|
# leaks. OpenAI's AsyncStream uses close(), async generators use aclose().
|
||||||
async with chunk_stream:
|
@asynccontextmanager
|
||||||
|
async def _closing(stream):
|
||||||
|
try:
|
||||||
|
yield stream
|
||||||
|
finally:
|
||||||
|
if hasattr(stream, "aclose"):
|
||||||
|
await stream.aclose()
|
||||||
|
elif hasattr(stream, "close"):
|
||||||
|
await stream.close()
|
||||||
|
|
||||||
|
async with _closing(chunk_stream):
|
||||||
async for chunk in chunk_stream:
|
async for chunk in chunk_stream:
|
||||||
if chunk.usage:
|
if chunk.usage:
|
||||||
cached_tokens = (
|
cached_tokens = (
|
||||||
|
|||||||
@@ -149,13 +149,9 @@ async def test_openai_llm_stream_closed_on_cancellation():
|
|||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.iteration_count = 0
|
self.iteration_count = 0
|
||||||
|
|
||||||
async def __aenter__(self):
|
async def close(self):
|
||||||
return self
|
|
||||||
|
|
||||||
async def __aexit__(self, exc_type, exc_val, exc_tb):
|
|
||||||
nonlocal stream_closed
|
nonlocal stream_closed
|
||||||
stream_closed = True
|
stream_closed = True
|
||||||
return False
|
|
||||||
|
|
||||||
def __aiter__(self):
|
def __aiter__(self):
|
||||||
return self
|
return self
|
||||||
|
|||||||
Reference in New Issue
Block a user