Streaming responses¶
By default every backend preloads the response body into memory before
send() returns — the right thing for API payloads. For large downloads
and endless feeds pass stream=True to
RequestsBackend,
HttpxBackend,
AsyncHttpxBackend,
AiohttpBackend,
UrllibBackend or
Urllib3Backend: send()
then returns as soon as the headers arrived, and the response body is a
streaming BodyProducer over the still-open
connection — an IterableBody on the sync
backends, an AsyncIterableBody on the
async ones.
Downloading to a file¶
On a sync backend, iterate the producer’s
chunks() — each chunk is written
out as it arrives, so the download never occupies more memory than one
chunk:
from action0.client import Client
from action0.client.backends.requests import RequestsBackend
from action0.req import Request
with RequestsBackend(stream=True) as backend:
response = Client(backend).send(Request("https://example.com/big.bin"))
producer = response.body_producer()
with open("big.bin", "wb") as file:
for chunk in producer.chunks():
file.write(chunk)
Async streams¶
On the async backends the producer is consumed with async for over
achunks() (the sync accessors
chunks() / body_bytes() / body_str() raise RuntimeError there —
an async source has no synchronous view):
import asyncio
from action0.client import Client
from action0.client.backends.httpx import AsyncHttpxBackend
from action0.req import Request
async def download(url: str, path: str) -> None:
async with AsyncHttpxBackend(stream=True) as backend:
response = await Client(backend).send(Request(url))
with open(path, "wb") as file:
async for chunk in response.body_producer().achunks():
file.write(chunk)
Chunk boundaries are whatever the transport delivers — they carry no meaning. A consumer of a line-oriented feed (NDJSON, SSE and friends) buffers across chunks and splits on its own delimiter:
import json
from typing import Any
from typing import AsyncIterator
from action0.req.body import BodyProducer
async def json_lines(producer: BodyProducer) -> AsyncIterator[Any]:
"""Decode an NDJSON stream, one object per line."""
buffer = b""
async for chunk in producer.achunks():
buffer += chunk
while (newline := buffer.find(b"\n")) >= 0:
line, buffer = buffer[:newline], buffer[newline + 1 :]
if line.strip():
yield json.loads(line)
async def watch_events(url: str) -> None:
async with AsyncHttpxBackend(stream=True) as backend:
response = await Client(backend).send(Request(url))
async for event in json_lines(response.body_producer()):
print(event["type"])
Streaming through typed operations¶
Streaming composes with operations too: declare the
result type as BodyProducer and have
load() hand out the producer instead of parsing the body. check()
only inspects the status, so for a streamed response the whole
parse() step runs at headers arrival — nothing reads the body until
your code does:
from action0.client import APIClient
from action0.client import Operation
from action0.client import path_param
from action0.client.backends.requests import RequestsBackend
from action0.req import BodyProducer
from action0.req import BytesBody
from action0.req import Method
from action0.req import Response
class DownloadExport(Operation[BodyProducer]):
method = Method.GET
path = "/exports/{export_id}"
export_id: int = path_param()
def load(self, response: Response) -> BodyProducer:
# an empty (bodyless) response streams as zero chunks
return response.body_producer() or BytesBody(b"")
with RequestsBackend(stream=True) as backend:
client = APIClient(backend, "https://api.example.com/v1")
producer = client.send(DownloadExport(export_id=7)) # typed BodyProducer
with open("export-7.csv", "wb") as file:
for chunk in producer.chunks():
file.write(chunk)
Don’t combine stream=True with operations that parse the whole
payload: JsonOperation’s load() reads the
complete body (joining the chunks), so nothing is gained — and on an
async backend that synchronous read raises RuntimeError. Keep bulk
endpoints on a streaming backend and JSON endpoints on a preloading one;
backends are cheap to have two of.
Testing streaming consumers¶
The stub backends need no special support: a canned
Response carries a streamed body by wrapping
the chunks in an IterableBody (or an
AsyncIterableBody for async consumers):
from action0.client import Client
from action0.client.testing import StubBackend
from action0.req import IterableBody
from action0.req import Request
from action0.req import Response
backend = StubBackend(Response(200, body=IterableBody([b"alpha\n", b"beta\n"])))
response = Client(backend).send(Request("https://api.example.com/feed"))
for chunk in response.body_producer().chunks():
print(chunk)
# b'alpha\n'
# b'beta\n'
What to know¶
The connection is held until the body is consumed; an abandoned producer closes it when garbage-collected. Consume (or drop) bodies promptly — and consume them fully where you can: a body read to the end keeps the connection alive for reuse where the underlying library supports it.
Hooks fire at headers arrival — the
elapsedanon_response()sees excludes the body transfer.body_bytes()/body_str()still work on a sync streamed response (they join the chunks — once; a second read finds the producer drained). An async streamed body can only be read viaachunks().How timeouts guard the stream follows the underlying library: requests, httpx and urllib apply the timeout to each socket read, so it stays in force between chunks; on
AiohttpBackendthetimeout=budget is aiohttp’s total and spans the body consumption too — for long-lived streams pass a session with a tailoredClientTimeout(e.g.sock_readinstead oftotal).Error translation wraps
send(), which returned at headers arrival — an error while iterating a streamed body (a dropped connection, a read timeout) is raised by the producer as the underlying library’s own exception, not aTransportError.The caching wrappers never store streamed bodies, and a retry wrapper that throws away a transient streamed response leaves the connection cleanup to garbage collection.
TwistedBackendhas no streaming mode — itsreadBodycollects the whole body; use a customIProtocolconsumer directly on the Agent if you need that on Twisted.