Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/errors.md
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ Unlike `DecodeError`, this error is raised before the request is sent.

## `ResponseTooLargeError`

Both clients accept `max_response_body_bytes: int | None = None`. By default there is no limit. When it is set, a response body larger than the cap raises `ResponseTooLargeError` instead of being returned, whatever the status: a `200` trips it as easily as a `500`. The cap counts decoded bytes, after decompression. It applies to `send()` and the verb methods, and to the error body that `stream()` reads before raising a `StatusError`. Bytes you read yourself while iterating a `stream()` are never capped. With a cap set and `follow_redirects=True`, httpware follows the redirects itself and caps only the final response. It closes each intermediate redirect response without reading its body, so the responses in `response.history` have no content. Client `auth` is sent to the first URL only. `httpx2` keeps its `Authorization` header on a redirect within the same origin or from `http` to `https` on the same host, and drops it otherwise.
Both clients accept `max_response_body_bytes: int | None = None`. By default there is no limit. When it is set, a response body larger than the cap raises `ResponseTooLargeError` instead of being returned, whatever the status: a `200` trips it as easily as a `500`. The cap counts decoded bytes, after decompression. It applies to `send()` and the verb methods, and to the error body that `stream()` reads before raising a `StatusError`. Bytes you read yourself while iterating a `stream()` are never capped, except with an `auth` that sets `requires_response_body`: that auth needs every body, so `stream()` reads the final one under the cap before yielding it. With a cap set, httpware runs the client's `auth` flow and, with `follow_redirects=True`, follows the redirects itself, and caps only the final response. It closes each intermediate response, such as a redirect or a `DigestAuth` challenge, without reading its body, so the responses in `response.history` have no content. An `auth` that sets `requires_response_body`, like a token refresh that parses the token endpoint's JSON, gets each response read under the cap instead. Client `auth` is sent to the first URL only. `httpx2` keeps its `Authorization` header on a redirect within the same origin or from `http` to `https` on the same host, and drops it otherwise.

`ResponseTooLargeError` carries:

Expand Down
170 changes: 130 additions & 40 deletions src/httpware/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
"event_hooks": "event_hooks=... is not supported; use middleware=... instead.",
}
_TOO_MANY_REDIRECTS_MESSAGE = "Exceeded maximum allowed redirects."
_NO_AUTH = httpx2.Auth()
_BASE_URL_QUERY_MESSAGE = (
"base_url must not contain a query string: httpx2 appends request paths after it, "
"producing malformed URLs. Pass the query as params=... instead."
Expand Down Expand Up @@ -133,58 +134,153 @@ def _select_httpx2_options(
return forwarded


async def _send_following_redirects_async(client: httpx2.AsyncClient, request: httpx2.Request) -> httpx2.Response:
"""Send `request` streaming, following redirects hop by hop without reading intermediate bodies."""
history: list[httpx2.Response] = []
response = await client.send(request, stream=True, follow_redirects=False)
while client.follow_redirects and response.next_request is not None:
def _request_auth(client: httpx2.Client | httpx2.AsyncClient, request: httpx2.Request) -> httpx2.Auth:
"""Return the auth httpx2 applies to `request`: the client's, else Basic from URL credentials, else none."""
if client.auth is not None:
return client.auth
if request.url.username or request.url.password:
return httpx2.BasicAuth(request.url.username, request.url.password)
return _NO_AUTH


async def _read_capped_and_close_async(streaming: httpx2.Response, cap: int) -> httpx2.Response:
"""Buffer `streaming` under `cap` via `_read_capped_async`, closing it either way."""
try:
return await _read_capped_async(streaming, cap, streaming.request)
finally:
await streaming.aclose()


async def _send_redirect_hops_async(
client: httpx2.AsyncClient,
request: httpx2.Request,
prior_history: list[httpx2.Response],
) -> httpx2.Response:
"""Send `request` streaming, following redirects hop by hop without reading intermediate bodies.

Histories and the `max_redirects` count include `prior_history`, as in httpx2.
"""
hops = list(prior_history)
while True:
if len(hops) > client.max_redirects:
raise httpx2.TooManyRedirects(_TOO_MANY_REDIRECTS_MESSAGE, request=request)
response = await client.send(request, stream=True, follow_redirects=False, auth=_NO_AUTH)
response.history = list(hops)
if not client.follow_redirects or response.next_request is None:
return response
await response.aclose()
history.append(response)
if len(history) > client.max_redirects:
raise httpx2.TooManyRedirects(_TOO_MANY_REDIRECTS_MESSAGE, request=response.next_request)
response = await client.send(response.next_request, stream=True, follow_redirects=False, auth=None)
response.history = history
return response
hops.append(response)
request = response.next_request


async def _send_capped_async(client: httpx2.AsyncClient, request: httpx2.Request, cap: int) -> httpx2.Response:
"""Send `request` streaming, driving the client's auth flow and redirects without reading intermediate bodies.

An auth that sets `requires_response_body` gets each response buffered under `cap` instead.
"""
auth = _request_auth(client, request)
flow = auth.async_auth_flow(request)
history: list[httpx2.Response] = []
try:
request = await anext(flow)
while True:
response = await _send_redirect_hops_async(client, request, history)
if auth.requires_response_body:
response = await _read_capped_and_close_async(response, cap)
try:
next_request = await flow.asend(response)
except StopAsyncIteration:
return response
except BaseException:
await response.aclose()
raise
await response.aclose()
response.history = list(history)
history.append(response)
request = next_request
finally:
await flow.aclose()


@contextlib.asynccontextmanager
async def _stream_following_redirects_async(
async def _stream_capped_async(
client: httpx2.AsyncClient,
method: str,
url: httpx2.URL | str,
kwargs: dict[str, typing.Any],
cap: int,
) -> AsyncIterator[httpx2.Response]:
"""Async mirror of `httpx2.AsyncClient.stream` that follows redirects via `_send_following_redirects_async`."""
response = await _send_following_redirects_async(client, client.build_request(method, url, **kwargs))
"""Async mirror of `httpx2.AsyncClient.stream` that sends via `_send_capped_async`."""
response = await _send_capped_async(client, client.build_request(method, url, **kwargs), cap)
try:
yield response
finally:
await response.aclose()


def _send_following_redirects(client: httpx2.Client, request: httpx2.Request) -> httpx2.Response:
"""Sync mirror of `_send_following_redirects_async`."""
history: list[httpx2.Response] = []
response = client.send(request, stream=True, follow_redirects=False)
while client.follow_redirects and response.next_request is not None:
def _read_capped_and_close(streaming: httpx2.Response, cap: int) -> httpx2.Response:
"""Sync mirror of `_read_capped_and_close_async`."""
try:
return _read_capped(streaming, cap, streaming.request)
finally:
streaming.close()


def _send_redirect_hops(
client: httpx2.Client,
request: httpx2.Request,
prior_history: list[httpx2.Response],
) -> httpx2.Response:
"""Sync mirror of `_send_redirect_hops_async`."""
hops = list(prior_history)
while True:
if len(hops) > client.max_redirects:
raise httpx2.TooManyRedirects(_TOO_MANY_REDIRECTS_MESSAGE, request=request)
response = client.send(request, stream=True, follow_redirects=False, auth=_NO_AUTH)
response.history = list(hops)
if not client.follow_redirects or response.next_request is None:
return response
response.close()
history.append(response)
if len(history) > client.max_redirects:
raise httpx2.TooManyRedirects(_TOO_MANY_REDIRECTS_MESSAGE, request=response.next_request)
response = client.send(response.next_request, stream=True, follow_redirects=False, auth=None)
response.history = history
return response
hops.append(response)
request = response.next_request


def _send_capped(client: httpx2.Client, request: httpx2.Request, cap: int) -> httpx2.Response:
"""Sync mirror of `_send_capped_async`."""
auth = _request_auth(client, request)
flow = auth.sync_auth_flow(request)
history: list[httpx2.Response] = []
try:
request = next(flow)
while True:
response = _send_redirect_hops(client, request, history)
if auth.requires_response_body:
response = _read_capped_and_close(response, cap)
try:
next_request = flow.send(response)
except StopIteration:
return response
except BaseException:
response.close()
raise
response.close()
response.history = list(history)
history.append(response)
request = next_request
finally:
flow.close()


@contextlib.contextmanager
def _stream_following_redirects(
def _stream_capped(
client: httpx2.Client,
method: str,
url: httpx2.URL | str,
kwargs: dict[str, typing.Any],
cap: int,
) -> Iterator[httpx2.Response]:
"""Sync mirror of `_stream_following_redirects_async`."""
response = _send_following_redirects(client, client.build_request(method, url, **kwargs))
"""Sync mirror of `_stream_capped_async`."""
response = _send_capped(client, client.build_request(method, url, **kwargs), cap)
try:
yield response
finally:
Expand Down Expand Up @@ -284,11 +380,8 @@ async def _terminal(self, request: httpx2.Request) -> httpx2.Response:
if cap is None:
response = await self._httpx2_client.send(request)
else:
streaming = await _send_following_redirects_async(self._httpx2_client, request)
try:
response = await _read_capped_async(streaming, cap, streaming.request)
finally:
await streaming.aclose()
streaming = await _send_capped_async(self._httpx2_client, request, cap)
response = await _read_capped_and_close_async(streaming, cap)
except RuntimeError as exc:
if self._httpx2_client.is_closed:
raise TransportError(str(exc)) from exc
Expand Down Expand Up @@ -1140,7 +1233,7 @@ async def stream( # noqa: PLR0913 — mirrors httpx2 per-method signatures; kwa
opened = (
self._httpx2_client.stream(method, merged_url, **kwargs)
if cap is None
else _stream_following_redirects_async(self._httpx2_client, method, merged_url, kwargs)
else _stream_capped_async(self._httpx2_client, method, merged_url, kwargs, cap)
)
async with _httpx2_exception_mapper(), opened as response:
if HTTPStatus.BAD_REQUEST <= response.status_code < 600: # noqa: PLR2004 — 600 is the synthetic upper bound for 5xx
Expand Down Expand Up @@ -1225,11 +1318,8 @@ def _terminal(self, request: httpx2.Request) -> httpx2.Response:
if cap is None:
response = self._httpx2_client.send(request)
else:
streaming = _send_following_redirects(self._httpx2_client, request)
try:
response = _read_capped(streaming, cap, streaming.request)
finally:
streaming.close()
streaming = _send_capped(self._httpx2_client, request, cap)
response = _read_capped_and_close(streaming, cap)
except RuntimeError as exc:
if self._httpx2_client.is_closed:
raise TransportError(str(exc)) from exc
Expand Down Expand Up @@ -2102,7 +2192,7 @@ def stream( # noqa: PLR0913 — mirrors httpx2 per-method signatures; kwargs-fo
opened = (
self._httpx2_client.stream(method, merged_url, **kwargs)
if cap is None
else _stream_following_redirects(self._httpx2_client, method, merged_url, kwargs)
else _stream_capped(self._httpx2_client, method, merged_url, kwargs, cap)
)
with _httpx2_exception_mapper_sync(), opened as response:
if HTTPStatus.BAD_REQUEST <= response.status_code < 600: # noqa: PLR2004 — 600 is the synthetic upper bound for 5xx
Expand Down
Loading
Loading