Search / proxy_streams.py
VoiceOfML
Separate runtime responsibilities and close proxy streams reliably
9f9f582
Raw History Blame Contribute Delete
1.82 kB
"""Own upstream responses across the route-to-ASGI streaming handoff."""
from fastapi.responses import StreamingResponse
class OwnedStreamingResponse(StreamingResponse):
def __init__(self, content, *, cleanup, **kwargs):
super().__init__(content, **kwargs)
self.cleanup = cleanup
async def __call__(self, scope, receive, send):
try:
await super().__call__(scope, receive, send)
finally:
# A disconnect may happen before the generator's first iteration.
try:
self.cleanup()
finally:
await self.body_iterator.aclose()
class UpstreamLease:
"""Own an upstream response and an optional, already-acquired pool slot."""
def __init__(self, semaphore=None):
self.response = None
self.semaphore = semaphore
self.closed = False
self.transferred = False
def close(self):
if self.closed:
return
self.closed = True
try:
if self.response is not None:
self.response.release()
finally:
if self.semaphore is not None:
self.semaphore.release()
def __enter__(self):
return self
def __exit__(self, *args):
if not self.transferred:
self.close()
async def chunks(self, error_label):
try:
async for chunk in self.response.content.iter_chunked(65536):
yield chunk
except Exception as exc:
print(f"{error_label}: {exc}")
raise
finally:
self.close()
def streaming_response(self, content, **kwargs):
response = OwnedStreamingResponse(content, cleanup=self.close, **kwargs)
self.transferred = True
return response