fix(openai): use compatible stream closing for non-OpenAI providers
OpenAI's AsyncStream uses close() while async generators (e.g. from OpenPipe) use aclose(). Replace direct async-with on the stream with a helper that handles both protocols.
This commit is contained in:
@@ -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 = (
|
||||||
|
|||||||
Reference in New Issue
Block a user