From 9a4454fe9e14f19af1214b3765bef45e92dccd86 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Tue, 22 Sep 2026 23:17:00 +0200 Subject: [PATCH 01/12] Serialize HTTP/2 send path under a dedicated lock 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 --- httpcore/_async/http2.py | 38 +++--- httpcore/_sync/http2.py | 44 ++++--- tests/_sync/test_http2_thread_safety.py | 149 ++++++++++++++++++++++++ 3 files changed, 205 insertions(+), 26 deletions(-) create mode 100644 tests/_sync/test_http2_thread_safety.py diff --git a/httpcore/_async/http2.py b/httpcore/_async/http2.py index dbd0beeb4..456e13403 100644 --- a/httpcore/_async/http2.py +++ b/httpcore/_async/http2.py @@ -60,6 +60,10 @@ def __init__( self._state_lock = AsyncLock() self._read_lock = AsyncLock() self._write_lock = AsyncLock() + # Guards the shared `h2` state machine on the send path, plus + # stream-ID allocation and the `_events` mapping. See the sync + # `HTTP2Connection` for details. (encode/httpx#3566) + self._send_lock = AsyncLock() self._sent_connection_init = False self._used_all_stream_ids = False self._connection_error = False @@ -131,19 +135,23 @@ async def handle_async_request(self, request: Request) -> Response: await self._max_streams_semaphore.acquire() try: - stream_id = self._h2_state.get_next_available_stream_id() - self._events[stream_id] = [] - except h2.exceptions.NoAvailableStreamIDError: # pragma: nocover - self._used_all_stream_ids = True - self._request_count -= 1 - raise ConnectionNotAvailable() - - try: - kwargs = {"request": request, "stream_id": stream_id} - async with Trace("send_request_headers", logger, request, kwargs): - await self._send_request_headers(request=request, stream_id=stream_id) - async with Trace("send_request_body", logger, request, kwargs): - await self._send_request_body(request=request, stream_id=stream_id) + # The send path mutates the shared `h2` state machine, which is + # not safe for concurrent use, so stream ID allocation and the + # send itself are serialized under a single lock. Note that h2 + # requires `get_next_available_stream_id()` to be immediately + # followed by the matching `send_headers()` call, otherwise + # concurrent tasks may be handed duplicate stream IDs. + # (encode/httpx#3566) + async with self._send_lock: + stream_id = self._h2_state.get_next_available_stream_id() + self._events[stream_id] = [] + kwargs = {"request": request, "stream_id": stream_id} + async with Trace("send_request_headers", logger, request, kwargs): + await self._send_request_headers( + request=request, stream_id=stream_id + ) + async with Trace("send_request_body", logger, request, kwargs): + await self._send_request_body(request=request, stream_id=stream_id) async with Trace( "receive_response_headers", logger, request, kwargs ) as trace: @@ -162,6 +170,10 @@ async def handle_async_request(self, request: Request) -> Response: "stream_id": stream_id, }, ) + except h2.exceptions.NoAvailableStreamIDError: # pragma: nocover + self._used_all_stream_ids = True + self._request_count -= 1 + raise ConnectionNotAvailable() except BaseException as exc: # noqa: PIE786 with AsyncShieldCancellation(): kwargs = {"stream_id": stream_id} diff --git a/httpcore/_sync/http2.py b/httpcore/_sync/http2.py index ddcc18900..9e96b944a 100644 --- a/httpcore/_sync/http2.py +++ b/httpcore/_sync/http2.py @@ -60,6 +60,15 @@ def __init__( self._state_lock = Lock() self._read_lock = Lock() self._write_lock = Lock() + # Guards the shared `h2` state machine on the send path, plus + # stream-ID allocation and the `_events` mapping. The same + # `HTTP2Connection` is handed to multiple threads by the connection + # pool (HTTP/2 multiplexing), and `h2` is not thread-safe: concurrent + # `send_headers` calls corrupt the HPACK encoder table ("deque mutated + # during iteration") and the streams dict ("dictionary changed size + # during iteration"), and concurrent `get_next_available_stream_id` + # calls can hand out duplicate stream IDs. (encode/httpx#3566) + self._send_lock = Lock() self._sent_connection_init = False self._used_all_stream_ids = False self._connection_error = False @@ -131,19 +140,24 @@ def handle_request(self, request: Request) -> Response: self._max_streams_semaphore.acquire() try: - stream_id = self._h2_state.get_next_available_stream_id() - self._events[stream_id] = [] - except h2.exceptions.NoAvailableStreamIDError: # pragma: nocover - self._used_all_stream_ids = True - self._request_count -= 1 - raise ConnectionNotAvailable() - - try: - kwargs = {"request": request, "stream_id": stream_id} - with Trace("send_request_headers", logger, request, kwargs): - self._send_request_headers(request=request, stream_id=stream_id) - with Trace("send_request_body", logger, request, kwargs): - self._send_request_body(request=request, stream_id=stream_id) + # The send path mutates the shared `h2` state machine, which is + # not thread-safe, so stream ID allocation and the send itself + # are serialized under a single lock. Note that h2 requires + # `get_next_available_stream_id()` to be immediately followed by + # the matching `send_headers()` call, otherwise concurrent + # threads may be handed duplicate stream IDs. Without this, + # multithreaded clients sharing one connection hit errors such as + # "deque mutated during iteration", "dictionary changed size + # during iteration", and `StreamIDTooLowError`. + # (encode/httpx#3566) + with self._send_lock: + stream_id = self._h2_state.get_next_available_stream_id() + self._events[stream_id] = [] + kwargs = {"request": request, "stream_id": stream_id} + with Trace("send_request_headers", logger, request, kwargs): + self._send_request_headers(request=request, stream_id=stream_id) + with Trace("send_request_body", logger, request, kwargs): + self._send_request_body(request=request, stream_id=stream_id) with Trace( "receive_response_headers", logger, request, kwargs ) as trace: @@ -162,6 +176,10 @@ def handle_request(self, request: Request) -> Response: "stream_id": stream_id, }, ) + except h2.exceptions.NoAvailableStreamIDError: # pragma: nocover + self._used_all_stream_ids = True + self._request_count -= 1 + raise ConnectionNotAvailable() except BaseException as exc: # noqa: PIE786 with ShieldCancellation(): kwargs = {"stream_id": stream_id} diff --git a/tests/_sync/test_http2_thread_safety.py b/tests/_sync/test_http2_thread_safety.py new file mode 100644 index 000000000..4c1e48693 --- /dev/null +++ b/tests/_sync/test_http2_thread_safety.py @@ -0,0 +1,149 @@ +"""Regression test for encode/httpx#3566. + +`httpx.Client(http2=True)` shares a single connection across threads (HTTP/2 +multiplexing). That connection is `httpcore.HTTP2Connection`, which wraps one +`h2.H2Connection` state machine. The send path — stream-ID allocation, the +`_events` mapping, and the HPACK encoding in `send_headers` — used to be +completely unserialized, so concurrent threads could corrupt the state +machine ("deque mutated during iteration", "dictionary changed size during +iteration", `StreamIDTooLowError`). + +This test drives one shared `HTTP2Connection` from multiple threads over a +fully in-memory fake HTTP/2 server (a real server-side `h2` connection +guarded by its own lock, so only the *client-side* httpcore code is under +test) and asserts every request succeeds. +""" + +import threading +import time + +import h2.config +import h2.connection +import h2.events + +import httpcore + + +class FakeStream: + """In-memory full-duplex socket backed by a real server-side h2 connection.""" + + def __init__(self): + self._server = h2.connection.H2Connection( + config=h2.config.H2Configuration(client_side=False) + ) + self._server.initiate_connection() + self._lock = threading.Lock() + self._cond = threading.Condition(self._lock) + self._closed = False + self._responded = set() + + def _respond(self, stream_id): + self._server.send_headers( + stream_id, + [(":status", "200"), ("content-length", "2")], + end_stream=False, + ) + self._server.send_data(stream_id, b"ok", end_stream=True) + + # -- NetworkStream interface -- + def write(self, data, timeout=None): + with self._lock: + if data: + for event in self._server.receive_data(data): + if isinstance(event, h2.events.RequestReceived): + if event.stream_ended is not None: + self._responded.add(event.stream_id) + self._respond(event.stream_id) + elif isinstance(event, h2.events.DataReceived): + # Replenish the server's inbound flow-control window, + # like any real HTTP/2 server does. + self._server.acknowledge_received_data( + event.flow_controlled_length, event.stream_id + ) + elif isinstance(event, h2.events.StreamEnded): + if event.stream_id not in self._responded: + self._responded.add(event.stream_id) + self._respond(event.stream_id) + self._cond.notify_all() + + def read(self, n, timeout=None): + deadline = None if timeout is None else time.monotonic() + timeout + with self._lock: + while True: + data = self._server.data_to_send(n) + if data: + return data + if self._closed: + return b"" + remaining = None if deadline is None else deadline - time.monotonic() + if remaining is not None and remaining <= 0: + return b"" + self._cond.wait(timeout=1.0) + + def close(self): + with self._lock: + self._closed = True + self._cond.notify_all() + + def get_extra_info(self, name): + return None + + +def test_http2_connection_is_thread_safe(): + """ + One shared HTTP2Connection must survive N threads racing the send path. + + Without the `_send_lock` fix, this fails with errors such as + "deque mutated during iteration", "dictionary changed size during + iteration", `StreamIDTooLowError`, or `KeyError`. + """ + n_threads = 8 + n_requests = 50 + + origin = httpcore.Origin(b"https", b"example.org", 443) + connection = httpcore.HTTP2Connection(origin=origin, stream=FakeStream()) + + errors = [] + errors_lock = threading.Lock() + + def worker(worker_id): + for i in range(n_requests): + body = b"x" * 64 if i % 2 else b"" + headers = [ + (b"host", b"example.org"), + (b"user-agent", b"thread-safety-test"), + (b"x-request", f"{worker_id}-{i}".encode()), + ] + if body: + headers.append((b"content-length", str(len(body)).encode())) + request = httpcore.Request( + "POST" if body else "GET", + f"https://example.org/{worker_id}/{i}", + headers=headers, + content=body, + extensions={"timeout": {"read": 10, "write": 10}}, + ) + try: + response = connection.handle_request(request) + content = response.read() + response.close() + assert response.status == 200 + assert content == b"ok" + except Exception as exc: # noqa: BLE001 + with errors_lock: + errors.append(exc) + + threads = [ + threading.Thread(target=worker, args=(w,), name=f"h2-worker-{w}") + for w in range(n_threads) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=120) + + assert not any(thread.is_alive() for thread in threads), "worker threads hung" + assert errors == [], ( + f"{len(errors)} request(s) failed out of {n_threads * n_requests}: " + f"{sorted({type(e).__name__ for e in errors})}" + ) From bfa0a7179ee5739f82fe6e04d077ca4375c44e14 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 23 Sep 2026 14:37:54 +0200 Subject: [PATCH 02/12] Fix lint failures (mypy, unasync) on the h2 send-lock change - 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 | 23 ++++++++++++----------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/tests/_sync/test_http2_thread_safety.py b/tests/_sync/test_http2_thread_safety.py index 4c1e48693..094c77892 100644 --- a/tests/_sync/test_http2_thread_safety.py +++ b/tests/_sync/test_http2_thread_safety.py @@ -16,6 +16,7 @@ import threading import time +import typing import h2.config import h2.connection @@ -24,10 +25,10 @@ import httpcore -class FakeStream: +class FakeStream(httpcore.NetworkStream): """In-memory full-duplex socket backed by a real server-side h2 connection.""" - def __init__(self): + def __init__(self) -> None: self._server = h2.connection.H2Connection( config=h2.config.H2Configuration(client_side=False) ) @@ -35,9 +36,9 @@ def __init__(self): self._lock = threading.Lock() self._cond = threading.Condition(self._lock) self._closed = False - self._responded = set() + self._responded: set[int] = set() - def _respond(self, stream_id): + def _respond(self, stream_id: int) -> None: self._server.send_headers( stream_id, [(":status", "200"), ("content-length", "2")], @@ -46,10 +47,10 @@ def _respond(self, stream_id): self._server.send_data(stream_id, b"ok", end_stream=True) # -- NetworkStream interface -- - def write(self, data, timeout=None): + def write(self, buffer: bytes, timeout: float | None = None) -> None: with self._lock: - if data: - for event in self._server.receive_data(data): + if buffer: + for event in self._server.receive_data(buffer): if isinstance(event, h2.events.RequestReceived): if event.stream_ended is not None: self._responded.add(event.stream_id) @@ -66,11 +67,11 @@ def write(self, data, timeout=None): self._respond(event.stream_id) self._cond.notify_all() - def read(self, n, timeout=None): + def read(self, max_bytes: int, timeout: float | None = None) -> bytes: deadline = None if timeout is None else time.monotonic() + timeout with self._lock: while True: - data = self._server.data_to_send(n) + data = self._server.data_to_send(max_bytes) if data: return data if self._closed: @@ -80,12 +81,12 @@ def read(self, n, timeout=None): return b"" self._cond.wait(timeout=1.0) - def close(self): + def close(self) -> None: with self._lock: self._closed = True self._cond.notify_all() - def get_extra_info(self, name): + def get_extra_info(self, info: str) -> typing.Any: return None From 32b8974fe8e2a5a5464d7b59456f6fd6da79b060 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 23 Sep 2026 14:37:57 +0200 Subject: [PATCH 03/12] Fix lint failures (mypy, unasync) on the h2 send-lock change - 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. --- httpcore/_async/http2.py | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/httpcore/_async/http2.py b/httpcore/_async/http2.py index 456e13403..fe047971f 100644 --- a/httpcore/_async/http2.py +++ b/httpcore/_async/http2.py @@ -61,8 +61,13 @@ def __init__( self._read_lock = AsyncLock() self._write_lock = AsyncLock() # Guards the shared `h2` state machine on the send path, plus - # stream-ID allocation and the `_events` mapping. See the sync - # `HTTP2Connection` for details. (encode/httpx#3566) + # stream-ID allocation and the `_events` mapping. The same + # `HTTP2Connection` is handed to multiple threads by the connection + # pool (HTTP/2 multiplexing), and `h2` is not thread-safe: concurrent + # `send_headers` calls corrupt the HPACK encoder table ("deque mutated + # during iteration") and the streams dict ("dictionary changed size + # during iteration"), and concurrent `get_next_available_stream_id` + # calls can hand out duplicate stream IDs. (encode/httpx#3566) self._send_lock = AsyncLock() self._sent_connection_init = False self._used_all_stream_ids = False @@ -140,7 +145,10 @@ async def handle_async_request(self, request: Request) -> Response: # send itself are serialized under a single lock. Note that h2 # requires `get_next_available_stream_id()` to be immediately # followed by the matching `send_headers()` call, otherwise - # concurrent tasks may be handed duplicate stream IDs. + # concurrent callers may be handed duplicate stream IDs. Without + # this, clients sharing one connection hit errors such as + # "deque mutated during iteration", "dictionary changed size + # during iteration", and `StreamIDTooLowError`. # (encode/httpx#3566) async with self._send_lock: stream_id = self._h2_state.get_next_available_stream_id() @@ -148,7 +156,8 @@ async def handle_async_request(self, request: Request) -> Response: kwargs = {"request": request, "stream_id": stream_id} async with Trace("send_request_headers", logger, request, kwargs): await self._send_request_headers( - request=request, stream_id=stream_id + request=request, + stream_id=stream_id, ) async with Trace("send_request_body", logger, request, kwargs): await self._send_request_body(request=request, stream_id=stream_id) From 6dab40264bd4d99aad2a4f1adfac9aeb36bea863 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 23 Sep 2026 14:38:00 +0200 Subject: [PATCH 04/12] Fix lint failures (mypy, unasync) on the h2 send-lock change - 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. --- httpcore/_sync/http2.py | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/httpcore/_sync/http2.py b/httpcore/_sync/http2.py index 9e96b944a..0689764c7 100644 --- a/httpcore/_sync/http2.py +++ b/httpcore/_sync/http2.py @@ -141,12 +141,12 @@ def handle_request(self, request: Request) -> Response: try: # The send path mutates the shared `h2` state machine, which is - # not thread-safe, so stream ID allocation and the send itself - # are serialized under a single lock. Note that h2 requires - # `get_next_available_stream_id()` to be immediately followed by - # the matching `send_headers()` call, otherwise concurrent - # threads may be handed duplicate stream IDs. Without this, - # multithreaded clients sharing one connection hit errors such as + # not safe for concurrent use, so stream ID allocation and the + # send itself are serialized under a single lock. Note that h2 + # requires `get_next_available_stream_id()` to be immediately + # followed by the matching `send_headers()` call, otherwise + # concurrent callers may be handed duplicate stream IDs. Without + # this, clients sharing one connection hit errors such as # "deque mutated during iteration", "dictionary changed size # during iteration", and `StreamIDTooLowError`. # (encode/httpx#3566) @@ -155,7 +155,10 @@ def handle_request(self, request: Request) -> Response: self._events[stream_id] = [] kwargs = {"request": request, "stream_id": stream_id} with Trace("send_request_headers", logger, request, kwargs): - self._send_request_headers(request=request, stream_id=stream_id) + self._send_request_headers( + request=request, + stream_id=stream_id, + ) with Trace("send_request_body", logger, request, kwargs): self._send_request_body(request=request, stream_id=stream_id) with Trace( From 21731aea74c595ac891d9819d9e9408c6b4c40c0 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 22:28:01 +0200 Subject: [PATCH 05/12] Bump twine to 7.0.0 for Metadata-Version 2.5 support --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/requirements.txt b/requirements.txt index 90219a8ce..397452f08 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,7 +10,7 @@ jinja2==3.1.6 # Packaging build==1.2.2.post1 -twine==6.1.0 +twine==7.0.0 # Tests & Linting coverage[toml]==7.5.4 From 0107d3b8cf29a8d1705ca533fdfbdec18abab6da Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 23:52:59 +0200 Subject: [PATCH 06/12] Scope twine 7 pin to Python 3.9+ (twine 7 requires it) --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/requirements.txt b/requirements.txt index 397452f08..d58db6805 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,7 +10,7 @@ jinja2==3.1.6 # Packaging build==1.2.2.post1 -twine==7.0.0 +twine==7.0.0; python_version >= "3.9" # twine 7 needs py3.9+ (Metadata-Version 2.5) # Tests & Linting coverage[toml]==7.5.4 From b7029ca7c7a7f84f63ea2079519605868ad71f06 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 23:53:01 +0200 Subject: [PATCH 07/12] Skip twine check on Python 3.8 (twine 7 requires 3.9+) --- scripts/build | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/scripts/build b/scripts/build index 657ded044..97ba4b247 100755 --- a/scripts/build +++ b/scripts/build @@ -11,5 +11,9 @@ fi set -x ${PREFIX}python -m build -${PREFIX}twine check dist/* +# twine>=7 (needed for Metadata-Version 2.5) requires Python 3.9+, +# so only check the distribution where twine is installed. +if ${PREFIX}python -c "import twine" 2>/dev/null; then + ${PREFIX}twine check dist/* +fi ${PREFIX}mkdocs build From 07448c6f720664934938851b23291b21f6876427 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 23:54:36 +0200 Subject: [PATCH 08/12] Fix twine marker: twine 7 requires Python 3.10+, not 3.9 --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/requirements.txt b/requirements.txt index d58db6805..84e96dce4 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,7 +10,7 @@ jinja2==3.1.6 # Packaging build==1.2.2.post1 -twine==7.0.0; python_version >= "3.9" # twine 7 needs py3.9+ (Metadata-Version 2.5) +twine==7.0.0; python_version >= "3.10" # twine 7 needs py3.10+ (Metadata-Version 2.5) # Tests & Linting coverage[toml]==7.5.4 From e729bc71a3b9169eb4c2b407974a0d6ece779da9 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 23:54:39 +0200 Subject: [PATCH 09/12] Fix comment: twine 7 requires Python 3.10+ --- scripts/build | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/build b/scripts/build index 97ba4b247..8737b8a49 100755 --- a/scripts/build +++ b/scripts/build @@ -11,7 +11,7 @@ fi set -x ${PREFIX}python -m build -# twine>=7 (needed for Metadata-Version 2.5) requires Python 3.9+, +# twine>=7 (needed for Metadata-Version 2.5) requires Python 3.10+, # so only check the distribution where twine is installed. if ${PREFIX}python -c "import twine" 2>/dev/null; then ${PREFIX}twine check dist/* From 9ae267a6f1bb0122daa312a2cc718cab9227f420 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Wed, 30 Sep 2026 23:58:54 +0200 Subject: [PATCH 10/12] Mark unreachable defensive branches no-cover (100% coverage gate) --- tests/_sync/test_http2_thread_safety.py | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/tests/_sync/test_http2_thread_safety.py b/tests/_sync/test_http2_thread_safety.py index 094c77892..fb05ee057 100644 --- a/tests/_sync/test_http2_thread_safety.py +++ b/tests/_sync/test_http2_thread_safety.py @@ -74,19 +74,19 @@ def read(self, max_bytes: int, timeout: float | None = None) -> bytes: data = self._server.data_to_send(max_bytes) if data: return data - if self._closed: + if self._closed: # pragma: no cover return b"" - remaining = None if deadline is None else deadline - time.monotonic() - if remaining is not None and remaining <= 0: + remaining = None if deadline is None else deadline - time.monotonic() # pragma: no cover + if remaining is not None and remaining <= 0: # pragma: no cover return b"" - self._cond.wait(timeout=1.0) + self._cond.wait(timeout=1.0) # pragma: no cover - def close(self) -> None: + def close(self) -> None: # pragma: no cover with self._lock: self._closed = True self._cond.notify_all() - def get_extra_info(self, info: str) -> typing.Any: + def get_extra_info(self, info: str) -> typing.Any: # pragma: no cover return None @@ -130,7 +130,7 @@ def worker(worker_id): response.close() assert response.status == 200 assert content == b"ok" - except Exception as exc: # noqa: BLE001 + except Exception as exc: # noqa: BLE001 # pragma: no cover with errors_lock: errors.append(exc) From dd586884c0f8f46ba5478d0f2e705823379efafd Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Thu, 1 Oct 2026 00:03:27 +0200 Subject: [PATCH 11/12] Use typing.Optional for 3.8/3.9 compat (X | Y needs 3.10) --- tests/_sync/test_http2_thread_safety.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/_sync/test_http2_thread_safety.py b/tests/_sync/test_http2_thread_safety.py index fb05ee057..53aa1f41a 100644 --- a/tests/_sync/test_http2_thread_safety.py +++ b/tests/_sync/test_http2_thread_safety.py @@ -47,7 +47,7 @@ def _respond(self, stream_id: int) -> None: self._server.send_data(stream_id, b"ok", end_stream=True) # -- NetworkStream interface -- - def write(self, buffer: bytes, timeout: float | None = None) -> None: + def write(self, buffer: bytes, timeout: typing.Optional[float] = None) -> None: with self._lock: if buffer: for event in self._server.receive_data(buffer): @@ -67,7 +67,7 @@ def write(self, buffer: bytes, timeout: float | None = None) -> None: self._respond(event.stream_id) self._cond.notify_all() - def read(self, max_bytes: int, timeout: float | None = None) -> bytes: + def read(self, max_bytes: int, timeout: typing.Optional[float] = None) -> bytes: deadline = None if timeout is None else time.monotonic() + timeout with self._lock: while True: From aaf674a21e037c25cc5616132ab25b15869ef000 Mon Sep 17 00:00:00 2001 From: Dextheking1 Date: Thu, 1 Oct 2026 00:09:21 +0200 Subject: [PATCH 12/12] Fix mypy strict errors on Python 3.8 (typing.Set, bytes annotation) --- tests/_sync/test_http2_thread_safety.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/_sync/test_http2_thread_safety.py b/tests/_sync/test_http2_thread_safety.py index 53aa1f41a..e6aaa04a0 100644 --- a/tests/_sync/test_http2_thread_safety.py +++ b/tests/_sync/test_http2_thread_safety.py @@ -36,7 +36,7 @@ def __init__(self) -> None: self._lock = threading.Lock() self._cond = threading.Condition(self._lock) self._closed = False - self._responded: set[int] = set() + self._responded: typing.Set[int] = set() def _respond(self, stream_id: int) -> None: self._server.send_headers( @@ -71,7 +71,7 @@ def read(self, max_bytes: int, timeout: typing.Optional[float] = None) -> bytes: deadline = None if timeout is None else time.monotonic() + timeout with self._lock: while True: - data = self._server.data_to_send(max_bytes) + data: bytes = self._server.data_to_send(max_bytes) if data: return data if self._closed: # pragma: no cover