| 222 | |
| 223 | |
| 224 | class StreamingResponse(Response): |
| 225 | body_iterator: AsyncContentStream |
| 226 | |
| 227 | def __init__( |
| 228 | self, |
| 229 | content: ContentStream, |
| 230 | status_code: int = 200, |
| 231 | headers: Mapping[str, str] | None = None, |
| 232 | media_type: str | None = None, |
| 233 | background: BackgroundTask | None = None, |
| 234 | ) -> None: |
| 235 | if isinstance(content, AsyncIterable): |
| 236 | self.body_iterator = content |
| 237 | else: |
| 238 | self.body_iterator = iterate_in_threadpool(content) |
| 239 | self.status_code = status_code |
| 240 | self.media_type = self.media_type if media_type is None else media_type |
| 241 | self.background = background |
| 242 | self.init_headers(headers) |
| 243 | |
| 244 | async def listen_for_disconnect(self, receive: Receive) -> None: |
| 245 | while True: |
| 246 | message = await receive() |
| 247 | if message["type"] == "http.disconnect": |
| 248 | break |
| 249 | |
| 250 | async def stream_response(self, send: Send) -> None: |
| 251 | await send({"type": "http.response.start", "status": self.status_code, "headers": self.raw_headers}) |
| 252 | async for chunk in self.body_iterator: |
| 253 | if not isinstance(chunk, bytes | memoryview): |
| 254 | chunk = chunk.encode(self.charset) |
| 255 | await send({"type": "http.response.body", "body": chunk, "more_body": True}) |
| 256 | |
| 257 | await send({"type": "http.response.body", "body": b"", "more_body": False}) |
| 258 | |
| 259 | async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: |
| 260 | if scope["type"] == "websocket": |
| 261 | send = self._wrap_websocket_denial_send(send) |
| 262 | await self.stream_response(send) |
| 263 | if self.background is not None: |
| 264 | await self.background() |
| 265 | return |
| 266 | |
| 267 | spec_version = tuple(map(int, scope.get("asgi", {}).get("spec_version", "2.0").split("."))) |
| 268 | |
| 269 | if spec_version >= (2, 4): |
| 270 | try: |
| 271 | await self.stream_response(send) |
| 272 | except OSError: |
| 273 | raise ClientDisconnect() |
| 274 | else: |
| 275 | async with create_collapsing_task_group() as task_group: |
| 276 | |
| 277 | async def wrap(func: Callable[[], Awaitable[None]]) -> None: |
| 278 | await func() |
| 279 | task_group.cancel_scope.cancel() |
| 280 | |
| 281 | task_group.start_soon(wrap, partial(self.stream_response, send)) |
no outgoing calls