Skip to content
Open
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
9 changes: 8 additions & 1 deletion benchmark/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,14 @@ git worktree add ../httpx2-main main && (cd ../httpx2-main && uv sync)
scripts/benchmark --python main=../httpx2-main/.venv/bin/python --python work=.venv/bin/python --lib httpx2
```

`punkreq` is supported as an extra reference point (`--lib punkreq`) if it is installed.
The original `httpx` and `punkreq` are supported as reference points (`--lib httpx`, `--lib punkreq`).
Neither is part of the project's environment (installing `httpx` next to `httpx2` interferes with the
alias tests), so provision them in a separate interpreter and pass it with `--python`:

```
uv venv --python 3.14 /tmp/refs && uv pip install --python /tmp/refs/bin/python httpx punkreq zuvloop
scripts/benchmark --python refs=/tmp/refs/bin/python --lib httpx --lib punkreq
```

## Reading the results

Expand Down
26 changes: 21 additions & 5 deletions benchmark/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,19 @@ def check_size(received: int, expected: int) -> None:
def build_httpx2(scenario: Scenario, payload: bytes) -> tuple[RequestFn, CloseFn]:
import httpx2

return _build_httpx_like(httpx2, scenario, payload)


def build_httpx(scenario: Scenario, payload: bytes) -> tuple[RequestFn, CloseFn]:
# The original httpx, for reference; it shares the httpx2 API. Not part of the
# project environment: see the README for how to provision it.
import httpx

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: When --lib httpx is selected through the documented benchmark environment, this import fails because the original httpx package is not declared or locked by the project. Add httpx to the benchmark dependency group, or document and provision it separately.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At benchmark/client.py, line 73:

<comment>When `--lib httpx` is selected through the documented benchmark environment, this import fails because the original `httpx` package is not declared or locked by the project. Add `httpx` to the benchmark dependency group, or document and provision it separately.</comment>

<file context>
@@ -65,6 +65,18 @@ def check_size(received: int, expected: int) -> None:
+
+def build_httpx(scenario: Scenario, payload: bytes) -> tuple[RequestFn, CloseFn]:
+    # The original httpx, for reference; it shares the httpx2 API.
+    import httpx
+
+    return _build_httpx_like(httpx, scenario, payload)
</file context>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in c6d2d00, by documenting it: httpx cannot join the bench group because installing it next to httpx2 breaks the alias tests' type-checks, so the README now shows how to provision httpx/punkreq in a separate interpreter passed via --python, and the import site says so.


return _build_httpx_like(httpx, scenario, payload)


def _build_httpx_like(httpx2: Any, scenario: Scenario, payload: bytes) -> tuple[RequestFn, CloseFn]:

client = httpx2.AsyncClient(
limits=httpx2.Limits(
max_connections=scenario.max_connections,
Expand All @@ -79,13 +92,15 @@ def build_httpx2(scenario: Scenario, payload: bytes) -> tuple[RequestFn, CloseFn
if scenario.post:
headers.append(("content-length", str(len(payload))))

class RequestBody(httpx2.AsyncByteStream):
async def __aiter__(self) -> AsyncIterator[bytes]:
if scenario.post:
yield payload
async def request_body(self: object) -> AsyncIterator[bytes]:
if scenario.post:
yield payload

# Built dynamically so the helper can serve both httpx2 and httpx.
request_body_stream = type("RequestBody", (httpx2.AsyncByteStream,), {"__aiter__": request_body})

async def stream_one() -> None:
request = httpx2.Request(scenario.method, scenario.url, headers=headers, stream=RequestBody())
request = httpx2.Request(scenario.method, scenario.url, headers=headers, stream=request_body_stream())
response = await client.send(request, stream=True, follow_redirects=False)
received = 0
async for chunk in response.aiter_raw(scenario.chunk_size):
Expand Down Expand Up @@ -198,6 +213,7 @@ async def read_one() -> None:

BUILDERS: dict[str, Builder] = {
"httpx2": build_httpx2,
"httpx": build_httpx,
"httpcore2": build_httpcore2,
"aiohttp": build_aiohttp,
"punkreq": build_punkreq,
Expand Down
1 change: 1 addition & 0 deletions scripts/unasync.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
),
("import trio as concurrency", "from tests.httpcore2 import concurrency"),
("anyio.sleep", "concurrency.sleep"),
("BACKENDS = \\[None, httpcore2.AnyIOBackend\\(\\)\\]", "BACKENDS = [None]"),
("AsyncIterator", "Iterator"),
("Async([A-Z][A-Za-z0-9_]*)", r"\2"),
("async def", "def"),
Expand Down
2 changes: 2 additions & 0 deletions src/httpcore2/httpcore2/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
AsyncHTTPProxy,
AsyncSOCKSProxy,
)
from ._backends.asyncio import AsyncioBackend
from ._backends.base import (
SOCKET_OPTION,
AsyncNetworkBackend,
Expand Down Expand Up @@ -99,6 +100,7 @@ def __init__(self, *args, **kwargs): # type: ignore
# network backends, implementations
"SyncBackend",
"AnyIOBackend",
"AsyncioBackend",
"TrioBackend",
# network backends, mock implementations
"AsyncMockBackend",
Expand Down
299 changes: 299 additions & 0 deletions src/httpcore2/httpcore2/_backends/asyncio.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,299 @@
from __future__ import annotations

import asyncio
import collections
import inspect
import ssl
import typing

from .._exceptions import (
ConnectError,
ConnectTimeout,
ReadError,
ReadTimeout,
WriteError,
WriteTimeout,
)
from .base import SOCKET_OPTION, AsyncNetworkBackend, AsyncNetworkStream

# Stop reading from the socket once this much data is buffered but unread,
# and start again once the buffer drains below the low-water mark.
RECEIVE_HIGH_WATER = 256 * 1024
RECEIVE_LOW_WATER = 64 * 1024

# Stagger connection attempts across resolved addresses, as in RFC 8305.
HAPPY_EYEBALLS_DELAY = 0.25
# The event loop insists on a TLS handshake timeout; this stands in for "none".
NO_HANDSHAKE_TIMEOUT = 365 * 24 * 60 * 60.0

_happy_eyeballs_support: dict[type[asyncio.AbstractEventLoop], bool] = {}


def _connection_kwargs(loop: asyncio.AbstractEventLoop) -> dict[str, typing.Any]:
# Not every event loop implements Happy Eyeballs; use it where available.
supported = _happy_eyeballs_support.get(type(loop))
if supported is None:
try:
supported = "happy_eyeballs_delay" in inspect.signature(loop.create_connection).parameters
except (TypeError, ValueError):
# Some extension-implemented loops expose no signature to inspect.
supported = False
_happy_eyeballs_support[type(loop)] = supported
return {"happy_eyeballs_delay": HAPPY_EYEBALLS_DELAY} if supported else {}


class _Timeout(Exception):
"""
Raised on a waiter future when its deadline passes.
"""


def _timeout(waiter: asyncio.Future[None]) -> None:
if not waiter.done():
waiter.set_exception(_Timeout())


def _wake(waiter: asyncio.Future[None] | None) -> None:
if waiter is not None and not waiter.done():
waiter.set_result(None)


class AsyncioStreamProtocol(asyncio.Protocol):
"""
Buffers received data for `AsyncioStream`, applying backpressure to the
transport once too much is buffered, and wakes up pending reads and writes.
"""

def __init__(self) -> None:
self.transport: asyncio.Transport | None = None
self.chunks: collections.deque[bytes] = collections.deque()
self.buffered = 0
self.reading_paused = False
self.writing_paused = False
self.eof = False
self.closed = False
self.exception: Exception | None = None
self.read_waiter: asyncio.Future[None] | None = None
self.write_waiter: asyncio.Future[None] | None = None

def connection_made(self, transport: asyncio.BaseTransport) -> None:
# The transport implements the interface without necessarily subclassing it.
self.transport = typing.cast(asyncio.Transport, transport)

def data_received(self, data: bytes) -> None:
self.chunks.append(data)
self.buffered += len(data)
if self.buffered >= RECEIVE_HIGH_WATER and not self.reading_paused:
assert self.transport is not None
self.transport.pause_reading()
self.reading_paused = True
_wake(self.read_waiter)

def eof_received(self) -> bool:
self.eof = True
_wake(self.read_waiter)
# Let the transport close: an HTTP peer that has sent a FIN is done.
return False

def connection_lost(self, exc: Exception | None) -> None:
self.closed = True
self.exception = exc
_wake(self.read_waiter)
_wake(self.write_waiter)

def pause_writing(self) -> None:
self.writing_paused = True

def resume_writing(self) -> None:
self.writing_paused = False
_wake(self.write_waiter)


class AsyncioStream(AsyncNetworkStream):
def __init__(self, transport: asyncio.Transport, protocol: AsyncioStreamProtocol) -> None:
self._transport = transport
self._protocol = protocol

async def read(self, max_bytes: int, timeout: float | None = None) -> bytes:
protocol = self._protocol
if not protocol.chunks:
if protocol.eof or protocol.closed:
return self._read_at_end()
try:
await self._wait("read", timeout)
except _Timeout:
raise ReadTimeout("timed out") from None
if not protocol.chunks:
return self._read_at_end()

chunk = protocol.chunks[0]
if len(chunk) <= max_bytes:
protocol.chunks.popleft()
else:
protocol.chunks[0] = chunk[max_bytes:]
chunk = chunk[:max_bytes]
protocol.buffered -= len(chunk)
if protocol.reading_paused and protocol.buffered <= RECEIVE_LOW_WATER:
protocol.reading_paused = False
self._transport.resume_reading()
return chunk

def _read_at_end(self) -> bytes:
# No buffered data and no more coming: a clean EOF reads as empty,
# a connection dropped by an error is a read error.
if self._protocol.exception is not None:
raise ReadError(str(self._protocol.exception)) from self._protocol.exception
return b""

async def write(self, buffer: bytes, timeout: float | None = None) -> None:
if not buffer:
return
protocol = self._protocol
if protocol.closed or self._transport.is_closing():
raise WriteError("Connection closed")
self._transport.write(buffer)
if protocol.writing_paused:
# The transport's send buffer is full; wait for it to drain.
try:
await self._wait("write", timeout)
except _Timeout:
raise WriteTimeout("timed out") from None
if protocol.closed:
raise WriteError(str(protocol.exception or "Connection closed"))

async def _wait(self, kind: str, timeout: float | None) -> None:
loop = asyncio.get_running_loop()
waiter: asyncio.Future[None] = loop.create_future()
protocol = self._protocol
if kind == "read":
protocol.read_waiter = waiter
else:
protocol.write_waiter = waiter
handle = None if timeout is None else loop.call_later(timeout, _timeout, waiter)
try:
await waiter
finally:
if handle is not None:
handle.cancel()
if kind == "read":
protocol.read_waiter = None
else:
protocol.write_waiter = None

async def aclose(self) -> None:
if self._protocol.closed:
return
self._transport.close()
# Closing only schedules the socket close on the event loop. Yield once
# so it runs now, then force it if unsent data is still holding it up.
await asyncio.sleep(0)
if not self._protocol.closed:
self._transport.abort()

async def start_tls(
self,
ssl_context: ssl.SSLContext,
server_hostname: str | None = None,
timeout: float | None = None,
) -> AsyncNetworkStream:
protocol = self._protocol
if protocol.chunks or protocol.eof or protocol.closed:
# Nothing may arrive before the handshake: anything already buffered
# is plaintext that must not be mistaken for data received over TLS.
await self.aclose()
raise ConnectError("Received unexpected data before the TLS handshake")

loop = asyncio.get_running_loop()
# The loop reports its own handshake timeout as a connection error, so
# the deadline is applied here to raise a timeout, with the loop's own
# deadline kept out of the way.
handshake = loop.start_tls(
self._transport,
protocol,
ssl_context,
server_hostname=server_hostname,
ssl_handshake_timeout=NO_HANDSHAKE_TIMEOUT if timeout is None else timeout + 1.0,
)
try:
transport = await asyncio.wait_for(handshake, timeout)
except (TimeoutError, asyncio.TimeoutError):
self._transport.close()
raise ConnectTimeout("timed out") from None
except (OSError, ssl.SSLError) as exc:
self._transport.close()
raise ConnectError(str(exc)) from exc
if transport is None: # pragma: no cover
raise ConnectError("TLS handshake failed")
protocol.transport = transport
return AsyncioStream(transport, protocol)

def get_extra_info(self, info: str) -> typing.Any:
if info == "ssl_object":
return self._transport.get_extra_info("ssl_object")
if info == "client_addr":
return self._transport.get_extra_info("sockname")
if info == "server_addr":
return self._transport.get_extra_info("peername")
if info == "socket":
return self._transport.get_extra_info("socket")
if info == "is_readable":
# The event loop keeps reading while the connection is idle, so a
# FIN or stray data from the server is already known here without
# touching the socket.
protocol = self._protocol
return bool(protocol.chunks) or protocol.eof or protocol.closed
return None


class AsyncioBackend(AsyncNetworkBackend):
async def connect_tcp(
self,
host: str,
port: int,
timeout: float | None = None,
local_address: str | None = None,
socket_options: typing.Iterable[SOCKET_OPTION] | None = None,
) -> AsyncNetworkStream:
loop = asyncio.get_running_loop()
local_addr = None if local_address is None else (local_address, 0)
# By default TCP sockets opened in `asyncio` include TCP_NODELAY.
connect = loop.create_connection(
AsyncioStreamProtocol, host, port, local_addr=local_addr, **_connection_kwargs(loop)
)
return await self._connect(connect, timeout, socket_options)

async def connect_unix_socket(
self,
path: str,
timeout: float | None = None,
socket_options: typing.Iterable[SOCKET_OPTION] | None = None,
) -> AsyncNetworkStream:
loop = asyncio.get_running_loop()
connect = loop.create_unix_connection(AsyncioStreamProtocol, path)
return await self._connect(connect, timeout, socket_options)

async def _connect(
self,
connect: typing.Coroutine[typing.Any, typing.Any, tuple[asyncio.BaseTransport, AsyncioStreamProtocol]],
timeout: float | None,
socket_options: typing.Iterable[SOCKET_OPTION] | None,
) -> AsyncioStream:
try:
transport, protocol = await asyncio.wait_for(connect, timeout)
except (TimeoutError, asyncio.TimeoutError):
raise ConnectTimeout("timed out") from None
except OSError as exc:
raise ConnectError(str(exc)) from exc
stream = AsyncioStream(typing.cast(asyncio.Transport, transport), protocol)
if socket_options:
sock = transport.get_extra_info("socket")
try:
for option in socket_options:
sock.setsockopt(*option)
except OSError as exc:
await stream.aclose()
raise ConnectError(str(exc)) from exc
return stream

async def sleep(self, seconds: float) -> None:
await asyncio.sleep(seconds)
Loading