Skip to content

Serialize the HTTP/2 send path under a dedicated lock (thread-safety race) - #1118

Open
Dextheking1 wants to merge 12 commits into
encode:masterfrom
Dextheking1:fix/h2-thread-safety-3566
Open

Dextheking1 wants to merge 12 commits into
encode:masterfrom
Dextheking1:fix/h2-thread-safety-3566

Conversation

@Dextheking1

Copy link
Copy Markdown

Closes encode/httpx#3566

Root cause

httpx.Client(http2=True) shares a single connection across threads (HTTP/2 multiplexing). That connection is httpcore.HTTP2Connection (plus the async twin), which wraps one shared h2.H2Connection state machine. The receive path (_read_lock) and socket writes (_write_lock) were already serialized, but the send path was not: concurrent threads could interleave

  • get_next_available_stream_id() → duplicate stream IDs handed out (StreamIDTooLowError: 2411 is lower than 2413),
  • send_headers() → HPACK encode iterating dynamic_entries while another thread mutates it (RuntimeError: deque mutated during iteration — the first traceback in the linked issue),
  • open_outbound_streams iterating h2.streams while another thread adds/removes streams (dictionary changed size during iteration, KeyError).

h2's docs require get_next_available_stream_id() to be immediately followed by the matching send_headers() — impossible under concurrency without a lock. Serializing in httpx would kill multiplexing, so the fix has to live here.

The fix

One new lock, _send_lock (threading.Lock / anyio.AsyncLock), held in handle_request around stream-ID allocation, _events[stream_id] registration, _send_request_headers, and _send_request_body. The NoAvailableStreamIDError → ConnectionNotAvailable handling and the except BaseException cleanup (_response_closed, semaphore release) keep their exact original semantics — the send block simply moved inside the existing cleanup try.

Lock-order audit: _send_lock is acquired only in handle_request, outermost relative to _read_lock/_write_lock/_state_lock; the max-streams semaphore is acquired before it and its holder never waits on a waiter. No new deadlock.

Evidence

Fail-before / pass-after with an in-memory fake HTTP/2 transport (real client- and server-side h2 state machines, server-side guarded by its own lock so only client-side httpcore code is under test):

case pre-fix post-fix
new regression test: 8 threads × 50 requests (half with 64 B bodies) 327/400 fail — LocalProtocolError(StreamIDTooLowError), KeyError 400/400 pass (×4 runs)
GET 8×300 33/2400 ok 2400/2400 (×3)
POST 16×100, half with 64 B bodies 8–1191/1600 ok 1600/1600 (×10)
POST 16×400 — 6400/6400 (×3)
existing tests: sync http2 + connection + pool (37), async http2 + connection (38) 75 passed 75 passed

New regression test: tests/_sync/test_http2_thread_safety.py.

httpx.Client(http2=True) multiplexes requests from multiple threads over a
single connection. HTTP2Connection (and its async twin) wrapped one shared
h2.H2Connection state machine, but while the receive path and socket writes
were already serialized, the send path was not. Concurrent threads could
interleave stream-ID allocation (duplicate IDs -> StreamIDTooLowError),
HPACK encoding in send_headers ("deque mutated during iteration"), and
stream-dict iteration ("dictionary changed size during iteration").

Add a _send_lock (threading.Lock / anyio.AsyncLock), held in
handle_request around stream-ID allocation, _events registration,
_send_request_headers and _send_request_body. The lock is acquired
outermost relative to the existing _read_lock/_write_lock/_state_lock,
and after the max-streams semaphore, so no new deadlock is possible.

Adds tests/_sync/test_http2_thread_safety.py: one shared HTTP2Connection
driven from 8 threads over an in-memory fake HTTP/2 server. Fails on the
pre-fix code (327/400 requests errored with LocalProtocolError /
StreamIDTooLowError / KeyError), passes with the fix.

Related to encode/httpx#3566
- tests/_sync/test_http2_thread_safety.py: add type annotations so mypy
  strict passes (FakeStream now subclasses httpcore.NetworkStream, all
  methods annotated).
- Align the new comments between httpcore/_async/http2.py (unasync source
  of truth) and httpcore/_sync/http2.py so scripts/unasync.py --check
  passes; use a magic trailing comma on the split _send_request_headers
  call so ruff format and unasync agree.
- tests/_sync/test_http2_thread_safety.py: add type annotations so mypy
  strict passes (FakeStream now subclasses httpcore.NetworkStream, all
  methods annotated).
- Align the new comments between httpcore/_async/http2.py (unasync source
  of truth) and httpcore/_sync/http2.py so scripts/unasync.py --check
  passes; use a magic trailing comma on the split _send_request_headers
  call so ruff format and unasync agree.
- tests/_sync/test_http2_thread_safety.py: add type annotations so mypy
  strict passes (FakeStream now subclasses httpcore.NetworkStream, all
  methods annotated).
- Align the new comments between httpcore/_async/http2.py (unasync source
  of truth) and httpcore/_sync/http2.py so scripts/unasync.py --check
  passes; use a magic trailing comma on the split _send_request_headers
  call so ruff format and unasync agree.
@Priyankm23

Copy link
Copy Markdown

Hi @Dextheking1 / maintainers!

A quick note regarding the CI failure in this PR:
The failure at Build package & docs is caused by twine==6.1.0 in requirements.txt, which fails to recognize wheel Metadata-Version: 2.5 generated by newer hatchling. Bumping twine to 7.0.0 resolves the build failure.

Also, regarding the concurrency design:
Holding AsyncLock across _send_request_body means that if one async request is uploading or streaming a large body, all other async requests on that connection are blocked from sending until that upload completes. This creates a bottleneck and defeats HTTP/2 multiplexing in async mode.

I have implemented an alternative approach on my fork that:

  1. Uses AsyncThreadLock / ThreadLock (threading.Lock in sync, zero-cost no-op in async) so async requests can continue to multiplex and upload bodies freely without any lock contention.
  2. Synchronizes all shared in-memory _h2_state operations (stream ID allocation, headers, data, end_stream, acknowledge_received_data, and buffer draining) with zero network I/O under lock.
  3. Achieves 100% statement coverage across all 3,656 statements with all checks and tests passing.

Branch reference: https://github.com/Priyankm23/httpcore/tree/fix-http2-header-encoding-race (commit 7d48dcc).

I plan to open this as a PR against encode/httpcore as soon as the repository's temporary interaction limits are lifted, or maintainers are welcome to review the branch directly.

@adenzhou1350

Copy link
Copy Markdown

I reproduced a cancellation regression in the async path at aaf674a: cancelling a request while it waits for _send_lock reaches the cleanup before stream_id is assigned. The resulting UnboundLocalError masks the original CancelledError.

This uses a real TCP loopback HTTP/2 server and the public AsyncConnectionPool API, with one paused upload and a second request. Native Windows Python 3.12.13 and 3.14.3 both give:

  • Base 10a6582: CancelledError.
  • PR head aaf674a: UnboundLocalError: cannot access local variable 'stream_id' where it is not associated with a value (context: CancelledError).

Warmup, resumed upload and a subsequent request all return 200 on the same connection. The probe confirms the second task is suspended at async with self._send_lock before cancelling it; it does not replace client locks, streams or h2 state. Dependencies: anyio 4.11.0, h2 4.3.0, h11 0.16.0. This check only covers this cancellation case, not the overall thread-safety fix.

Could the pre-stream cancellation path release its acquired stream permit without invoking stream-specific cleanup? The existing cleanup should still run for errors after stream allocation.

Standalone reproduction (run against each checkout with its asyncio/http2 extras installed)
import asyncio
import json
from pathlib import Path

def await_frames(task):
    current = task.get_coro()
    frames = []
    while current is not None:
        frame = getattr(current, "cr_frame", getattr(current, "gi_frame", None))
        if frame is not None:
            frames.append((frame, {
                "file": frame.f_code.co_filename,
                "function": frame.f_code.co_name,
                "line": frame.f_lineno,
            }))
        current = getattr(current, "cr_await", getattr(current, "gi_yieldfrom", None))
    return frames

async def trial():
    import h2.config
    import h2.connection
    import h2.events
    import h2.settings
    import httpcore

    body_entered = asyncio.Event()
    release_body = asyncio.Event()
    second_received = asyncio.Event()
    handlers = set()
    peer_requests = []
    server_errors = []
    accepted_connections = 0

    async def peer(reader, writer):
        nonlocal accepted_connections
        accepted_connections += 1
        task = asyncio.current_task()
        handlers.add(task)
        connection = h2.connection.H2Connection(
            config=h2.config.H2Configuration(client_side=False)
        )
        connection.initiate_connection()
        connection.update_settings({h2.settings.SettingCodes.MAX_CONCURRENT_STREAMS: 10})
        writer.write(connection.data_to_send())
        await writer.drain()
        paths = {}
        try:
            while data := await reader.read(65536):
                for event in connection.receive_data(data):
                    if isinstance(event, h2.events.RequestReceived):
                        path = dict(event.headers)[b":path"].decode()
                        paths[event.stream_id] = path
                        peer_requests.append({"stream_id": event.stream_id, "path": path})
                        if path == "/second":
                            # Keep a baseline request pending on real response I/O.
                            second_received.set()
                    if isinstance(event, h2.events.DataReceived):
                        connection.acknowledge_received_data(
                            event.flow_controlled_length, event.stream_id
                        )
                    if isinstance(event, h2.events.StreamEnded):
                        if paths[event.stream_id] != "/second":
                            connection.send_headers(event.stream_id, [(b":status", b"200")])
                            connection.send_data(event.stream_id, b"ok", end_stream=True)
                writer.write(connection.data_to_send())
                await writer.drain()
        except Exception as exc:
            server_errors.append(type(exc).__name__ + ": " + str(exc))
        finally:
            writer.close()
            await writer.wait_closed()
            handlers.discard(task)

    async def slow_body():
        body_entered.set()
        await release_body.wait()
        yield b"x"

    server = await asyncio.start_server(peer, "127.0.0.1", 0)
    url = f"http://127.0.0.1:{server.sockets[0].getsockname()[1]}"
    first = second = None
    phase = None
    observation = None
    cancelled_result = None
    controls = []
    try:
        async with httpcore.AsyncConnectionPool(
            http1=False, http2=True, max_connections=1
        ) as pool:
            warmup = await asyncio.wait_for(pool.request("GET", url + "/warmup"), 2)
            controls.append([warmup.status, warmup.content.decode()])
            first = asyncio.create_task(pool.request(
                "POST", url + "/slow", headers={"content-length": "1"},
                content=slow_body(),
            ))
            await asyncio.wait_for(body_entered.wait(), 2)
            second = asyncio.create_task(pool.request("GET", url + "/second"))
            deadline = asyncio.get_running_loop().time() + 2
            while asyncio.get_running_loop().time() < deadline:
                if second.done():
                    raise AssertionError("Second request ended before cancellation")
                if second_received.is_set():
                    phase = "peer_received_second_headers_waiting_for_response"
                    break
                for frame, info in await_frames(second):
                    if (frame.f_code.co_name == "handle_async_request"
                            and Path(frame.f_code.co_filename).name == "http2.py"
                            and "stream_id" not in frame.f_locals):
                        source_line = Path(frame.f_code.co_filename).read_text().splitlines()[
                            frame.f_lineno - 1
                        ].strip()
                        if source_line == "async with self._send_lock:":
                            phase = "waiting_send_lock_before_stream_id_assignment"
                            observation = {**info, "source_line": source_line}
                            break
                if phase:
                    break
                await asyncio.sleep(0.001)
            assert phase is not None, "Could not prove the cancellation phase"
            second.cancel()
            try:
                await second
            except BaseException as exc:
                cancelled_result = {"type": type(exc).__name__, "message": str(exc)}
                if exc.__context__ is not None:
                    cancelled_result["context_type"] = type(exc.__context__).__name__
            else:
                raise AssertionError("Cancellation unexpectedly returned a response")
            release_body.set()
            response = await asyncio.wait_for(first, 2)
            controls.append([response.status, response.content.decode()])
            after = await asyncio.wait_for(pool.request("GET", url + "/after"), 2)
            controls.append([after.status, after.content.decode()])
    finally:
        release_body.set()
        for task in (first, second):
            if task is not None and not task.done():
                task.cancel()
        await asyncio.gather(*(t for t in (first, second) if t is not None),
                             return_exceptions=True)
        server.close()
        await server.wait_closed()
        if handlers:
            await asyncio.wait_for(asyncio.gather(*tuple(handlers)), 2)
    assert controls == [[200, "ok"]] * 3, controls
    assert accepted_connections == 1, accepted_connections
    assert not server_errors, server_errors
    return {"phase": phase, "cancellation": cancelled_result,
            "await_observation": observation, "controls": controls,
            "accepted_connections": accepted_connections,
            "peer_requests": peer_requests}

print(json.dumps(asyncio.run(asyncio.wait_for(trial(), 12)), indent=2))

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Development

Successfully merging this pull request may close these issues.

deque mutated / dictionary changed during iteration

3 participants