From 783f1c79fb77f8144e0e6e4126b1c2457d9249bc Mon Sep 17 00:00:00 2001 From: Ahmed TAHRI Date: Sun, 12 Jul 2026 19:37:29 +0100 Subject: [PATCH 1/5] fix: congestion+loss against more brittle network connection --- qh3/asyncio/_transport.py | 6 +- qh3/quic/connection.py | 25 +++++-- qh3/quic/recovery.py | 109 +++++++++++++++++++++++------- src/recovery.rs | 15 +++-- src/stream_sender.rs | 14 +++- src/udp.rs | 4 +- tests/test_asyncio.py | 21 ++++++ tests/test_connection.py | 67 ++++++++++++++++--- tests/test_recovery.py | 136 +++++++++++++++++++++++++++++++++++--- 9 files changed, 340 insertions(+), 57 deletions(-) diff --git a/qh3/asyncio/_transport.py b/qh3/asyncio/_transport.py index a3141632..24f19db5 100644 --- a/qh3/asyncio/_transport.py +++ b/qh3/asyncio/_transport.py @@ -382,7 +382,11 @@ def sendto_many(self, datagrams: list[bytes], addr: typing.Any = None) -> None: target = addr if addr is not None else self._address if target is not None: try: - state.send(datagrams, str(target[0]), int(target[1])) + sent = state.send(datagrams, str(target[0]), int(target[1])) + if sent < len(datagrams): + self._register_writer() + for dgram in datagrams[sent:]: + self._queue_write(dgram, addr) return except BlockingIOError: self._register_writer() diff --git a/qh3/quic/connection.py b/qh3/quic/connection.py index 9953cecc..6910327b 100644 --- a/qh3/quic/connection.py +++ b/qh3/quic/connection.py @@ -554,7 +554,7 @@ def __init__( self._datagrams_pending: deque[bytes] = deque() self._handshake_done_pending = False self._ping_pending: list[int] = [] - self._probe_pending = False + self._probe_pending = 0 self._retire_connection_ids: list[int] = [] self._streams_blocked_pending = False @@ -708,6 +708,7 @@ def connect(self, addr: NetworkAddress, now: float) -> None: else: self._version = self._configuration.supported_versions[0] self._connect(now=now) + self._loss.start_packet_pacing(now) def datagrams_to_send(self, now: float) -> list[tuple[bytes, NetworkAddress]]: """ @@ -3237,7 +3238,7 @@ def _push_crypto_data(self) -> None: buf.seek(0) def _send_probe(self) -> None: - self._probe_pending = True + self._probe_pending += 1 def _is_stateless_reset(self, datagram: bytes) -> bool: """ @@ -3813,7 +3814,7 @@ def _write_application( # PING (probe) if self._probe_pending: self._write_ping_frame(builder, comment="probe") - self._probe_pending = False + self._probe_pending -= 1 # CRYPTO if crypto_stream is not None and not crypto_stream.sender.buffer_is_empty: @@ -3952,6 +3953,15 @@ def _write_handshake( space = self._spaces[epoch] while True: + # Handshake packets use the same path pacer as application data. + # ACKs and PTO probes bypass pacing to preserve recovery latency. + if ( + space.ack_at is None or space.ack_at >= now + ) and not self._probe_pending: + self._pacing_at = self._loss._pacer.next_send_time(now=now) + if self._pacing_at is not None: + break + if epoch == tls.Epoch.INITIAL: packet_type = QuicPacketType.INITIAL else: @@ -3964,15 +3974,18 @@ def _write_handshake( self._write_ack_frame(builder=builder, space=space, now=now) # CRYPTO + crypto_written = False if not crypto_stream.sender.buffer_is_empty: if self._write_crypto_frame( builder=builder, space=space, stream=crypto_stream ): - self._probe_pending = False + crypto_written = True + self._probe_pending = max(self._probe_pending - 1, 0) # PING (probe) if ( self._probe_pending + and not crypto_written and not self._handshake_complete and ( epoch == tls.Epoch.HANDSHAKE @@ -3980,10 +3993,12 @@ def _write_handshake( ) ): self._write_ping_frame(builder, comment="probe") - self._probe_pending = False + self._probe_pending -= 1 if builder.packet_is_empty: break + if builder._packet.in_flight: + self._loss._pacer.update_after_send(now=now) def _write_ack_frame( self, builder: QuicPacketBuilder, space: QuicPacketSpace, now: float diff --git a/qh3/quic/recovery.py b/qh3/quic/recovery.py index ce10b04c..2d510af4 100644 --- a/qh3/quic/recovery.py +++ b/qh3/quic/recovery.py @@ -7,6 +7,7 @@ from .._compat import TRACE from .._hazmat import QuicPacketPacer, QuicRttMonitor, RangeSet from .logger import QuicLoggerTrace +from .packet import QuicPacketType from .packet_builder import QuicDeliveryState, QuicSentPacket # loss detection @@ -136,9 +137,9 @@ def _start_epoch(self, now: float) -> None: cwnd_seg = self._cwnd_epoch / self._max_datagram_size self._K = _cubic_root((W_max_seg - cwnd_seg) / K_CUBIC_C) - def on_packet_acked(self, packet: QuicSentPacket) -> None: + def on_packet_acked(self, packet: QuicSentPacket, now: float) -> None: self.bytes_in_flight -= packet.sent_bytes - self._last_ack = packet.sent_time + self._last_ack = now # HyStart++ round tracking (RFC 9406 4.3): a round ends when an # ACK is received for a packet whose number is at or above the @@ -175,13 +176,13 @@ def on_packet_acked(self, packet: QuicSentPacket) -> None: # exiting slow start without a loss (HyStart triggered) self._first_slow_start = False self._W_max = self.congestion_window - self._start_epoch(packet.sent_time) + self._start_epoch(now) if self._starting_congestion_avoidance: # entering congestion avoidance after a loss self._starting_congestion_avoidance = False self._first_slow_start = False - self._start_epoch(packet.sent_time) + self._start_epoch(now) # TCP-friendly estimate (Reno-like linear growth) self._W_est = int( @@ -189,7 +190,7 @@ def on_packet_acked(self, packet: QuicSentPacket) -> None: + self._max_datagram_size * (packet.sent_bytes / self.congestion_window) ) - t = packet.sent_time - self._t_epoch + t = now - self._t_epoch W_cubic = self._W_cubic(t + self._rtt) # clamp target @@ -212,6 +213,14 @@ def on_packet_acked(self, packet: QuicSentPacket) -> None: ) def on_packet_sent(self, packet: QuicSentPacket) -> None: + # Restart before accounting the new flight so a reset cannot leave + # bytes_in_flight above a freshly reduced congestion window. + if ( + not packet.is_pmtu_probe + and self._last_ack > 0.0 + and packet.sent_time - self._last_ack >= K_CUBIC_MAX_IDLE_TIME + ): + self._reset() self.bytes_in_flight += packet.sent_bytes # Track largest sent PN and bootstrap the HyStart++ round window # on the first packet sent during slow start. @@ -223,12 +232,6 @@ def on_packet_sent(self, packet: QuicSentPacket) -> None: and self._hystart_window_end is None ): self._hystart_window_end = packet.packet_number - # reset cwnd after prolonged idle - if self._last_ack > 0.0: - elapsed_idle = packet.sent_time - self._last_ack - if elapsed_idle >= K_CUBIC_MAX_IDLE_TIME: - self._reset() - def on_packets_expired(self, packets: Iterable[QuicSentPacket]) -> None: for packet in packets: self.bytes_in_flight -= packet.sent_bytes @@ -477,6 +480,13 @@ def reset_for_new_path(self) -> None: self._rtt_smoothed = 0.0 self._rtt_variance = 0.0 + def start_packet_pacing(self, now: float) -> None: + self._pacer.start_pacing( + now=now, + congestion_window=self._cc.congestion_window, + smoothed_rtt=self._rtt_initial / 1.25, + ) + def on_ack_received( self, space: QuicPacketSpace, @@ -494,10 +504,10 @@ def on_ack_received( (RFC 9002 6.2.1). Resetting in that case would prematurely clear the PTO backoff and let a stuck handshake under-probe. """ - is_ack_eliciting = False largest_acked = ack_rangeset.bounds()[1] - 1 largest_newly_acked = None - largest_sent_time = None + largest_newly_acked_ack_eliciting = None + largest_ack_eliciting_sent_time = None if largest_acked > space.largest_acked_packet: space.largest_acked_packet = largest_acked @@ -509,13 +519,13 @@ def on_ack_received( # remove packet and update counters packet = space.sent_packets.pop(packet_number) if packet.is_ack_eliciting: - is_ack_eliciting = True + largest_newly_acked_ack_eliciting = packet_number + largest_ack_eliciting_sent_time = packet.sent_time if not packet.is_pmtu_probe: space.ack_eliciting_in_flight -= 1 if packet.in_flight: - self._cc.on_packet_acked(packet) + self._cc.on_packet_acked(packet, now=now) largest_newly_acked = packet_number - largest_sent_time = packet.sent_time # trigger callbacks dh = packet.delivery_handlers @@ -527,8 +537,8 @@ def on_ack_received( if largest_newly_acked is None: return - if largest_acked == largest_newly_acked and is_ack_eliciting: - latest_rtt = now - largest_sent_time + if largest_acked == largest_newly_acked_ack_eliciting: + latest_rtt = now - largest_ack_eliciting_sent_time log_rtt = True # limit ACK delay to max_ack_delay @@ -569,6 +579,19 @@ def on_ack_received( self._detect_loss(space, now=now) + # A later Initial CRYPTO packet can be acknowledged while an earlier + # fragment is missing. Probe that gap immediately: the peer cannot + # advance TLS until offset zero arrives, and waiting for PTO only + # prolongs a handshake that has already demonstrated peer reachability. + if any( + packet.packet_type == QuicPacketType.INITIAL + and packet.is_crypto_packet + and packet.packet_number < largest_newly_acked + for packet in space.sent_packets.values() + ): + self._queue_oldest_crypto_packet(space) + self.start_packet_pacing(now) + # reset PTO count if reset_pto_count: self._pto_count = 0 @@ -583,7 +606,46 @@ def on_loss_detection_timeout(self, now: float) -> None: else: self._pto_count += 1 self._pto_total += 1 - self.reschedule_data(now=now) + probe_count = 2 + self._queue_probe_data(probe_count, now=now) + for _ in range(probe_count): + self._send_probe() + + def _queue_probe_data(self, probe_count: int, now: float) -> None: + """Queue copies of oldest outstanding data without retiring it.""" + queued = 0 + for space in self.spaces: + for packet in space.sent_packets.values(): + if not packet.is_crypto_packet: + continue + for handler, args in packet.delivery_handlers or (): + handler(QuicDeliveryState.LOST, *args) + queued += 1 + if queued >= probe_count: + return + if queued: + return + + for space in reversed(self.spaces): + packets = [] + for packet in space.sent_packets.values(): + if not (packet.in_flight and packet.is_ack_eliciting): + continue + packets.append(packet) + if len(packets) >= probe_count: + break + if packets: + self._on_packets_rescheduled(packets, space=space, now=now) + return + + def _queue_oldest_crypto_packet(self, space: QuicPacketSpace) -> bool: + for packet in space.sent_packets.values(): + if not packet.is_crypto_packet: + continue + for handler, args in packet.delivery_handlers or (): + handler(QuicDeliveryState.LOST, *args) + return True + return False def on_packet_sent(self, packet: QuicSentPacket, space: QuicPacketSpace) -> None: space.sent_packets[packet.packet_number] = packet @@ -632,7 +694,6 @@ def reschedule_data(self, now: float) -> None: # Reschedule oldest in-flight application data (up to 2 packets) # so it is sent in the same write cycle as the PTO probe. - app_rescheduled = False if not crypto_scheduled: for space in self.spaces: if not space.sent_packets: @@ -646,11 +707,11 @@ def reschedule_data(self, now: float) -> None: break if to_reschedule: self._on_packets_rescheduled(to_reschedule, space=space, now=now) - app_rescheduled = True - # If no data was rescheduled, send a PING as the ack-eliciting probe - if not crypto_scheduled and not app_rescheduled: - self._send_probe() + # Delivery callbacks are not guaranteed to queue retransmittable data + # (a PING-only packet is one example), so every PTO explicitly queues + # a PING. It will normally be coalesced with retransmitted data. + self._send_probe() def _detect_loss(self, space: QuicPacketSpace, now: float) -> None: """ diff --git a/src/recovery.rs b/src/recovery.rs index e19b43c5..4af81503 100644 --- a/src/recovery.rs +++ b/src/recovery.rs @@ -44,16 +44,23 @@ impl QuicPacketPacer { self.packet_time } + /// Seed pacing for a new path with one datagram of immediate credit. + fn start_pacing(&mut self, now: f64, congestion_window: usize, smoothed_rtt: f64) { + self.update_rate(congestion_window, smoothed_rtt); + self.bucket_time = self.packet_time.unwrap_or(0.0); + self.evaluation_time = now; + } + /// Computes the next send time given the current time. /// - /// If a packet time is defined and the bucket has been depleted (≤ 0), - /// returns `now + packet_time`; otherwise, returns `None`. + /// If the bucket does not contain enough credit for a complete packet, + /// returns the time at which the missing credit will be available. #[inline(always)] fn next_send_time(&mut self, now: f64) -> Option { if let Some(packet_time) = self.packet_time { self.update_bucket(now); - if self.bucket_time <= 0.0 { - return Some(now + packet_time); + if self.bucket_time + K_MICRO_SECOND < packet_time { + return Some(now + packet_time - self.bucket_time); } } None diff --git a/src/stream_sender.rs b/src/stream_sender.rs index 5423477d..a93f1aab 100644 --- a/src/stream_sender.rs +++ b/src/stream_sender.rs @@ -260,6 +260,11 @@ impl QuicStreamSender { if delivery == DELIVERY_ACKED { self.ack_count += 1; if stop > start { + self.pending.subtract(start, stop); + if stop <= self.buffer_start { + return; + } + let start = std::cmp::max(start, self.buffer_start); self.acked.add(start, Some(stop)); let first_range = self.acked.get_item(0); if first_range.0 == self.buffer_start { @@ -285,7 +290,14 @@ impl QuicStreamSender { // LOST self.loss_count += 1; if stop > start { - self.pending.add(start, Some(stop)); + let start = std::cmp::max(start, self.buffer_start); + if stop > start { + self.pending.add(start, Some(stop)); + for i in 0..self.acked.len() { + let range = self.acked.get_item(i); + self.pending.subtract(range.0, range.1); + } + } } if Some(stop) == self.buffer_fin { self.buffer_is_empty = false; diff --git a/src/udp.rs b/src/udp.rs index 161d15f6..b5a6a088 100644 --- a/src/udp.rs +++ b/src/udp.rs @@ -306,7 +306,7 @@ impl PyUdpSocketState { }; let sock_ref = UdpSockRef::from(&borrowed); - match self.inner.send(sock_ref, &transmit) { + match self.inner.try_send(sock_ref, &transmit) { Ok(()) => sent += group_count, Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => break, Err(e) => { @@ -328,7 +328,7 @@ impl PyUdpSocketState { }; let sock_ref = UdpSockRef::from(&borrowed); - match self.inner.send(sock_ref, &transmit) { + match self.inner.try_send(sock_ref, &transmit) { Ok(()) => sent += 1, Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => break, Err(e) => { diff --git a/tests/test_asyncio.py b/tests/test_asyncio.py index 23514fe7..b759d0d2 100644 --- a/tests/test_asyncio.py +++ b/tests/test_asyncio.py @@ -1223,6 +1223,27 @@ def test_sendto_many_single_datagram_uses_raw_send(self): # Single datagram → _raw_send → sock.sendto sock.sendto.assert_called_once() + def test_sendto_many_queues_rust_partial_send(self): + from unittest.mock import MagicMock + + transport, loop, _, _ = self._make_transport( + gso=True, connected_addr=("::1", 9999) + ) + state = MagicMock() + state.send.return_value = 1 + transport._udp_state = state + datagrams = [b"A" * 1280, b"B" * 1280, b"C" * 1280] + + transport.sendto_many(datagrams) + + state.send.assert_called_once_with(datagrams, "::1", 9999) + loop.add_writer.assert_called_once_with(42, transport._on_write_ready) + assert list(transport._send_queue) == [ + (datagrams[1], None), + (datagrams[2], None), + ] + assert transport.get_write_buffer_size() == 2560 + def test_sendto_many_without_gso(self): transport, _, sock, _ = self._make_transport(gso=False, connected_addr=("::1", 9999)) datagrams = [b"A" * 1280, b"B" * 1280] diff --git a/tests/test_connection.py b/tests/test_connection.py index ceed6d3a..1c3c4d20 100644 --- a/tests/test_connection.py +++ b/tests/test_connection.py @@ -166,6 +166,7 @@ def create_standalone_client(self, **client_options): ) ) client._ack_delay = 0 + disable_packet_pacing(client) # kick-off handshake client.connect(SERVER_ADDR, now=time.time()) @@ -260,6 +261,9 @@ def client_and_server( def disable_packet_pacing(connection): class DummyPacketPacer(QuicPacketPacer): + def start_pacing(self, now, congestion_window, smoothed_rtt): + pass + def next_send_time(self, now): return None @@ -562,7 +566,7 @@ def test_connect_with_loss_1(self): items = client.datagrams_to_send(now=now) assert datagram_sizes(items) == [1280, 1280] assert client.get_timer() == pytest.approx(1.998) - self.assertSentPackets(client, [2, 0, 0]) + self.assertSentPackets(client, [4, 0, 0]) self.assertEvents(client, []) # server receives INITIAL, sends INITIAL + HANDSHAKE @@ -720,7 +724,7 @@ def test_connect_with_loss_3(self): items = client.datagrams_to_send(now=now) assert datagram_sizes(items) == [1280, 1280] assert client.get_timer() == pytest.approx(1.998) - self.assertSentPackets(client, [2, 0, 0]) + self.assertSentPackets(client, [4, 0, 0]) self.assertEvents(client, []) # server receives duplicate INITIAL, retransmits INITIAL + HANDSHAKE @@ -803,9 +807,9 @@ def test_connect_with_loss_4(self): now = client.get_timer() client.handle_timer(now=now) items = client.datagrams_to_send(now=now) - assert datagram_sizes(items) == [45] + assert datagram_sizes(items) == [90] assert client.get_timer() == pytest.approx(0.9) - self.assertSentPackets(client, [0, 2, 0]) + self.assertSentPackets(client, [0, 3, 0]) self.assertEvents(client, []) # server receives PING, discards INITIAL and sends ACK @@ -823,7 +827,7 @@ def test_connect_with_loss_4(self): items = server.datagrams_to_send(now=now) assert datagram_sizes(items) == [1280, 886] assert server.get_timer() == pytest.approx(2.048) - self.assertSentPackets(server, [0, 3, 0]) + self.assertSentPackets(server, [0, 5, 0]) self.assertEvents(server, []) # handshake continues normally @@ -833,7 +837,7 @@ def test_connect_with_loss_4(self): items = client.datagrams_to_send(now=now) assert datagram_sizes(items) == [313] assert client.get_timer() == pytest.approx(1.366) - self.assertSentPackets(client, [0, 3, 1]) + self.assertSentPackets(client, [0, 4, 1]) self.assertEvents(client, HANDSHAKE_COMPLETED_EVENTS) now += TICK @@ -897,20 +901,21 @@ def test_connect_with_loss_5(self): self.assertSentPackets(server, [0, 0, 1]) self.assertEvents(server, HANDSHAKE_COMPLETED_EVENTS_SINGLE_CID) - # server side PTO retransmits HANDSHAKE_DONE + NEW_CONNECTION_IDs - # (RFC 9002 6.2.4: retransmit in-flight data instead of PING) + # Server-side PTO retransmits HANDSHAKE_DONE + NEW_CONNECTION_IDs + # and includes a PING to guarantee an ack-eliciting probe. now = server.get_timer() server.handle_timer(now=now) items = server.datagrams_to_send(now=now) - assert datagram_sizes(items) == [56] + assert datagram_sizes(items) == [57, 29] assert server.get_timer() == pytest.approx(0.975) - self.assertSentPackets(server, [0, 0, 1]) + self.assertSentPackets(server, [0, 0, 2]) # Server re-emits ConnectionIdIssued when re-writing CID frames self.assertEvents(server, HANDSHAKE_COMPLETED_EVENTS_SINGLE_CID[1:]) # client receives retransmitted HANDSHAKE_DONE + CIDs, sends ACK now += TICK client.receive_datagram(items[0][0], SERVER_ADDR, now=now) + client.receive_datagram(items[1][0], SERVER_ADDR, now=now) items = client.datagrams_to_send(now=now) assert datagram_sizes(items) == [32] self.assertEvents(client, []) @@ -923,6 +928,48 @@ def test_connect_with_loss_5(self): self.assertSentPackets(server, [0, 0, 0]) self.assertEvents(server, []) + def test_connect_with_first_initial_lost_and_long_server_cid(self): + with client_and_server( + handshake=False, server_options={"connection_id_length": 20} + ) as (client, server): + now = 0.0 + client.connect(SERVER_ADDR, now=now) + initial = client.datagrams_to_send(now=now) + assert len(initial) == 2 + client_initial_space = client._spaces[tls.Epoch.INITIAL] + _, second_range = client_initial_space.sent_packets[1].delivery_handlers[0] + second_start, second_stop = second_range + + # Deliver only the ClientHello tail. The server ACKs packet 1 but + # cannot advance TLS until the prefix is retransmitted. + now += TICK + server.receive_datagram(initial[1][0], CLIENT_ADDR, now=now) + crypto_receiver = server._crypto_streams[tls.Epoch.INITIAL].receiver + assert crypto_receiver.starting_offset() == 0 + assert crypto_receiver.highest_offset == second_stop + assert list(crypto_receiver._ranges) == [(second_start, second_stop)] + ack = server.datagrams_to_send(now=now) + assert ack + + now += TICK + for data, _ in ack: + client.receive_datagram(data, SERVER_ADDR, now=now) + + # The ACK for packet 1 proves that the peer has only the tail of + # the ClientHello. Immediately probe the missing packet 0 CRYPTO + # range instead of waiting for the loss timer or PTO. + retransmission = client.datagrams_to_send(now=now) + now += TICK + for data, _ in retransmission: + server.receive_datagram(data, CLIENT_ADDR, now=now) + + assert retransmission + assert crypto_receiver.starting_offset() == second_stop + + response = server.datagrams_to_send(now=now) + assert response + self.assertEvents(server, [events.ProtocolNegotiated]) + def test_connect_with_no_transport_parameters(self): real_initialize = QuicConnection._initialize diff --git a/tests/test_recovery.py b/tests/test_recovery.py index 686cf7d0..60db11a0 100644 --- a/tests/test_recovery.py +++ b/tests/test_recovery.py @@ -47,10 +47,9 @@ def test_with_measurement(self): self.pacer.update_after_send(now=1.0) assert self.pacer.next_send_time(now=1.0) == pytest.approx(1.00005) - # 2 packets - for i in range(2): - assert self.pacer.next_send_time(now=1.00005) is None - self.pacer.update_after_send(now=1.00005) + # one packet of elapsed time grants one packet of credit + assert self.pacer.next_send_time(now=1.00005) is None + self.pacer.update_after_send(now=1.00005) assert self.pacer.next_send_time(now=1.00005) == pytest.approx(1.0001) # 1 packet @@ -58,12 +57,22 @@ def test_with_measurement(self): self.pacer.update_after_send(now=1.0001) assert self.pacer.next_send_time(now=1.0001) == pytest.approx(1.00015) - # 2 packets - for i in range(2): - assert self.pacer.next_send_time(now=1.00015) is None - self.pacer.update_after_send(now=1.00015) + assert self.pacer.next_send_time(now=1.00015) is None + self.pacer.update_after_send(now=1.00015) assert self.pacer.next_send_time(now=1.00015) == pytest.approx(1.0002) + def test_partial_credit_does_not_release_full_packet(self): + self.pacer.start_pacing( + now=1.0, congestion_window=12800, smoothed_rtt=0.2664 + ) + assert self.pacer.packet_time == pytest.approx(0.02664) + assert self.pacer.next_send_time(now=1.0) is None + + self.pacer.update_after_send(now=1.0) + assert self.pacer.next_send_time(now=1.001) == pytest.approx(1.02664) + assert self.pacer.next_send_time(now=1.01) == pytest.approx(1.02664) + assert self.pacer.next_send_time(now=1.02664) is None + class TestQuicPacketRecovery: def setup_method(self): @@ -204,6 +213,8 @@ def _handler(state): cwnd_before = self.recovery._cc.congestion_window self.recovery._pto_count = 2 + probes = [] + self.recovery._send_probe = lambda: probes.append(True) self.recovery.reschedule_data(now=5.0) assert observed == [QuicDeliveryState.LOST] @@ -211,6 +222,71 @@ def _handler(state): assert self.recovery._cc.congestion_window == cwnd_before assert self.recovery._pto_count == 2 assert len(self.ONE_RTT_SPACE.sent_packets) == 0 + assert probes == [True] + + def test_ack_only_packet_does_not_supply_rtt_sample(self): + data_packet = QuicSentPacket( + epoch=tls.Epoch.ONE_RTT, + in_flight=True, + is_ack_eliciting=True, + is_crypto_packet=False, + packet_number=0, + packet_type=QuicPacketType.ONE_RTT, + sent_bytes=1280, + sent_time=0.0, + ) + ack_packet = QuicSentPacket( + epoch=tls.Epoch.ONE_RTT, + in_flight=False, + is_ack_eliciting=False, + is_crypto_packet=False, + packet_number=1, + packet_type=QuicPacketType.ONE_RTT, + sent_bytes=30, + sent_time=5.0, + ) + self.recovery.on_packet_sent(data_packet, self.ONE_RTT_SPACE) + self.recovery.on_packet_sent(ack_packet, self.ONE_RTT_SPACE) + + ranges = RangeSet() + ranges.add(0, 2) + self.recovery.on_ack_received( + self.ONE_RTT_SPACE, ack_rangeset=ranges, ack_delay=0.0, now=10.0 + ) + + assert self.recovery._rtt_latest == 0.0 + assert not self.recovery._rtt_initialized + + def test_congestion_idle_time_uses_ack_receipt_time(self): + cc = QuicCongestionControl(max_datagram_size=1280) + first = QuicSentPacket( + epoch=tls.Epoch.ONE_RTT, + in_flight=True, + is_ack_eliciting=True, + is_crypto_packet=False, + packet_number=0, + packet_type=QuicPacketType.ONE_RTT, + sent_bytes=1280, + sent_time=1.0, + ) + cc.on_packet_sent(first) + cc.on_packet_acked(first, now=10.0) + cwnd = cc.congestion_window + + second = QuicSentPacket( + epoch=tls.Epoch.ONE_RTT, + in_flight=True, + is_ack_eliciting=True, + is_crypto_packet=False, + packet_number=1, + packet_type=QuicPacketType.ONE_RTT, + sent_bytes=1280, + sent_time=11.0, + ) + cc.on_packet_sent(second) + + assert cc.congestion_window == cwnd + assert cc.bytes_in_flight == 1280 def test_on_ack_received_initial_does_not_reset_pto_count(self): # RFC 9002 6.2.1: a client MUST NOT reset its PTO backoff on @@ -258,6 +334,44 @@ def test_on_ack_received_initial_does_not_reset_pto_count(self): ) assert self.recovery._pto_count == 0 + def test_initial_crypto_gap_is_probed_immediately(self): + observed = [] + probes = [] + self.recovery._send_probe = lambda: probes.append(True) + + for packet_number in range(2): + packet = QuicSentPacket( + epoch=tls.Epoch.INITIAL, + in_flight=True, + is_ack_eliciting=True, + is_crypto_packet=True, + packet_number=packet_number, + packet_type=QuicPacketType.INITIAL, + sent_bytes=1280, + sent_time=0.0, + ) + packet.delivery_handlers = [ + (lambda state, pn=packet_number: observed.append((pn, state)), ()) + ] + self.recovery.on_packet_sent(packet, self.INITIAL_SPACE) + + ranges = RangeSet() + ranges.add(1, 2) + self.recovery.on_ack_received( + self.INITIAL_SPACE, + ack_rangeset=ranges, + ack_delay=0.0, + now=0.05, + reset_pto_count=False, + ) + + assert observed == [ + (1, QuicDeliveryState.ACKED), + (0, QuicDeliveryState.LOST), + ] + assert probes == [] + assert list(self.INITIAL_SPACE.sent_packets) == [0] + def test_persistent_congestion_detected(self): # RFC 9002 7.6: if at least two ack-eliciting packets are # declared lost over a span longer than the persistent @@ -605,7 +719,9 @@ def _drive_round( cc.on_packet_sent(_hystart_packet(first_pn + i, send_time + i * 0.0001)) now = send_time + 0.001 for i in range(n_acks): - cc.on_packet_acked(_hystart_packet(first_pn + i, send_time + i * 0.0001)) + cc.on_packet_acked( + _hystart_packet(first_pn + i, send_time + i * 0.0001), now=now + ) cc.on_rtt_measurement(rtt, now) return first_pn + n_acks @@ -657,7 +773,7 @@ def test_css_growth_uses_divisor(self): cwnd0 = cc.congestion_window pkt = _hystart_packet(pn=0, sent_time=0.0, size=1280) cc.on_packet_sent(pkt) - cc.on_packet_acked(pkt) + cc.on_packet_acked(pkt, now=0.001) # 1280 // 4 == 320 assert cc.congestion_window == cwnd0 + 320 From 1139eb24ca36ed8a9a5d0e4f5a85abeed7cf700e Mon Sep 17 00:00:00 2001 From: Ahmed TAHRI Date: Sun, 12 Jul 2026 19:38:37 +0100 Subject: [PATCH 2/5] chore: speedup CI (tests_impl) build rust part without debug symbols --- noxfile.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/noxfile.py b/noxfile.py index 15e88b5f..bea94860 100644 --- a/noxfile.py +++ b/noxfile.py @@ -14,7 +14,7 @@ def tests_impl( session.install("-U", "pip", "maturin", silent=False) session.install("-r", "dev-requirements.txt", silent=False) - session.run("maturin", "develop") + session.run("maturin", "develop", "--release") # Show the pip version. session.run("pip", "--version") From 97da0700640f7bfee8d89c81849086d79768f878 Mon Sep 17 00:00:00 2001 From: Ahmed TAHRI Date: Sun, 12 Jul 2026 20:13:03 +0100 Subject: [PATCH 3/5] chore: update lsqpack to v2.6.5 --- Cargo.lock | 34 +++++++++++++++++----------------- Cargo.toml | 2 +- qh3/__init__.py | 2 +- 3 files changed, 19 insertions(+), 19 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e5c033ed..b947c6ec 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -177,9 +177,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.2.66" +version = "1.2.67" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f5d6cac793997bd970000024b2934968efe83b382de4fdcf4fcb46b6ee4ad996" +checksum = "e17dd265a7d0f31ef544e1b20e03add05d3b45b491b633b10d67145d2acc1a38" dependencies = [ "find-msvc-tools", "jobserver", @@ -669,9 +669,9 @@ checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" [[package]] name = "ls-qpack-rs" -version = "0.3.1" +version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "67ca13a214fd4d18efa614633fd38697117cb3728d131afadc2268e7f9f5e8e7" +checksum = "06a286661393cc63e0ace26270d7f2d643bdb6506c60bf899d7851e162f25827" dependencies = [ "libc", "ls-qpack-rs-sys", @@ -679,9 +679,9 @@ dependencies = [ [[package]] name = "ls-qpack-rs-sys" -version = "0.3.1" +version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62575b97833935501a543e5bcc43af9fec5042015f15dcbc3c3aee6e47d7f205" +checksum = "349b37abcdf89031e8730e233df7b3f6adf4f2ce4df81497daf97101e12fe406" dependencies = [ "bindgen", "cmake", @@ -974,7 +974,7 @@ dependencies = [ [[package]] name = "qh3" -version = "1.9.3" +version = "1.9.4" dependencies = [ "aws-lc-rs", "bincode", @@ -1058,9 +1058,9 @@ dependencies = [ [[package]] name = "regex" -version = "1.12.4" +version = "1.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f1292b7759ae1cb9ec195452d1390a074f0cd8541ab7a5a8c31cd6db45d4a6ba" +checksum = "2a0e75113e14dc5acb068cd0786884f214f1312650a3d36d269f5c4f3cdee8a2" dependencies = [ "aho-corasick", "memchr", @@ -1070,9 +1070,9 @@ dependencies = [ [[package]] name = "regex-automata" -version = "0.4.14" +version = "0.4.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6e1dd4122fc1595e8162618945476892eefca7b88c52820e74af6262213cae8f" +checksum = "1f388202e4b80542a0921078cc23b6333bcf1409c1e3f86404cae4766a6131db" dependencies = [ "aho-corasick", "memchr", @@ -1257,9 +1257,9 @@ dependencies = [ [[package]] name = "sha1" -version = "0.10.6" +version = "0.10.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" dependencies = [ "cfg-if", "cpufeatures", @@ -1671,18 +1671,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.53" +version = "0.8.54" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75726053136156d419e285b9b7eddaaea9e3fea6ce32eed44a89901f0bd98de1" +checksum = "b7cbbc0a705a0fd05cc3676525980d2bf5a9bc4adac6d6475209a7887cf59d19" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.53" +version = "0.8.54" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4714fd92cf900833d49538023a9b3915155210801d1c1169eba513b2addefd71" +checksum = "e2e817b7b52d0c7358d3246da9d69935ebb18116b2b102b4230dac079b4862f5" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml index 8ec9e2d5..3fe0bee5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "qh3" -version = "1.9.3" +version = "1.9.4" edition = "2021" rust-version = "1.75" license = "BSD-3-Clause" diff --git a/qh3/__init__.py b/qh3/__init__.py index acf2347c..c7b8cbb0 100644 --- a/qh3/__init__.py +++ b/qh3/__init__.py @@ -13,7 +13,7 @@ from .quic.packet import QuicProtocolVersion from .tls import CipherSuite, SessionTicket -__version__ = "1.9.3" +__version__ = "1.9.4" __all__ = ( "connect", From e88731cbde4e78bb118d73c7905b92283036bfb2 Mon Sep 17 00:00:00 2001 From: Ahmed TAHRI Date: Sun, 12 Jul 2026 20:15:44 +0100 Subject: [PATCH 4/5] docs: write changelog for 1.9.4 --- CHANGELOG.rst | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/CHANGELOG.rst b/CHANGELOG.rst index cef18622..688c2d4b 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -1,3 +1,12 @@ +1.9.4 (2026-07-12) +================== + +**Fixed** +- Congestion and loss algorithms against more brittle network connections + +**Changed** +- Updated lsqpack to v2.6.5 via our ls-qpack-rs crate v0.3.2 + 1.9.3 (2026-07-08) ================== From aa7e82f6afc08346639120aa3445b3fdb183e101 Mon Sep 17 00:00:00 2001 From: Ahmed TAHRI Date: Sun, 12 Jul 2026 20:18:53 +0100 Subject: [PATCH 5/5] chore: fix linter errors --- qh3/_hazmat.pyi | 3 +++ qh3/quic/recovery.py | 1 + 2 files changed, 4 insertions(+) diff --git a/qh3/_hazmat.pyi b/qh3/_hazmat.pyi index fe2f6577..90e14b4c 100644 --- a/qh3/_hazmat.pyi +++ b/qh3/_hazmat.pyi @@ -424,6 +424,9 @@ def decode_packet_number(truncated: int, num_bits: int, expected: int) -> int: class QuicPacketPacer: def __init__(self, max_datagram_size: int) -> None: ... + def start_pacing( + self, now: float, congestion_window: int, smoothed_rtt: float + ) -> None: ... def next_send_time(self, now: float) -> float | None: ... def update_after_send(self, now: float) -> None: ... def update_bucket(self, now: float) -> None: ... diff --git a/qh3/quic/recovery.py b/qh3/quic/recovery.py index 2d510af4..2143d2fc 100644 --- a/qh3/quic/recovery.py +++ b/qh3/quic/recovery.py @@ -232,6 +232,7 @@ def on_packet_sent(self, packet: QuicSentPacket) -> None: and self._hystart_window_end is None ): self._hystart_window_end = packet.packet_number + def on_packets_expired(self, packets: Iterable[QuicSentPacket]) -> None: for packet in packets: self.bytes_in_flight -= packet.sent_bytes