MCPcopy Create free account
hub / github.com/encode/starlette / StreamingResponse

Class StreamingResponse

starlette/responses.py:224–285  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

222
223
224class 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))

Calls

no outgoing calls