From 15a6bfcc84cce1f5c65d39ab86ca6b3ef81d78b3 Mon Sep 17 00:00:00 2001 From: Claude Date: Mon, 4 May 2026 22:55:34 +0000 Subject: [PATCH] =?UTF-8?q?feat(quic):=20step=206=20of=20RFC=209002=20retr?= =?UTF-8?q?ansmit=20=E2=80=94=20dispatch=20lost=20tokens=20to=20pending*?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit QuicConnection.onTokensLost(tokens) closes the loop between loss detection (step 5) and writer drain (step 4). For each lost token: - Ack: ignored (RFC 9000 §13.2.1: ACK frames not retransmittable) - MaxStreamsUni: pendingMaxStreamsUni := token.maxStreams iff token.maxStreams == advertisedMaxStreamsUni - MaxStreamsBidi / MaxData: same shape, against their advertised cap - MaxStreamData: pendingMaxStreamData[streamId] := maxData iff stream exists AND token.maxData == stream.receiveLimit The supersede check (`lost == advertised`) mirrors neqo's `fc.rs::frame_lost` line 322. If a higher extension has gone out since, the older lost frame is irrelevant — the newer value covers the receiver's grant. Without the check we'd resurrect stale extensions and waste wire bandwidth re-emitting values the peer already has. Wired into the parser's AckFrame handler immediately after detectAndRemoveLost: walk each lost packet's tokens and call onTokensLost. Caller already holds the connection lock. Tests added (8, all pass): - ackToken_doesNotPopulateAnyPending - lostMaxStreamsUni_matchingAdvertised_setsPending - lostMaxStreamsUni_supersededByHigherEmit_isDropped (the supersede-check invariant from neqo's fc.rs:322) - lostMaxStreamsBidi_matchingAdvertised_setsPending - lostMaxData_matchingAdvertised_setsPending - lostMaxData_supersededIsDropped - lostMaxStreamData_unknownStream_dropped (defensive) - multipleLostTokens_dispatchAll (one packet's worth of mixed tokens dispatched in one call) - lostTokensFromMultiplePackets_unionInPending (sequential dispatch across multiple lost packets — older stale, newer valid; the valid one survives) Mirror of neqo's `streams.rs::lost` dispatch + the `no_max_allowed_frame_after_old_loss` and `set_max_active_equal_does_not_set_frame_pending` tests from `fc.rs` deferred from step 4 per the plan. Full :quic test suite + nestsClient moq-lite tests pass. https://claude.ai/code/session_01PYYez8a6sjiakyjAxsfCEQ --- .../quic/connection/QuicConnection.kt | 50 +++++ .../quic/connection/QuicConnectionParser.kt | 10 +- .../quic/connection/OnTokensLostTest.kt | 199 ++++++++++++++++++ 3 files changed, 255 insertions(+), 4 deletions(-) create mode 100644 quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/OnTokensLostTest.kt diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt index 1754bf922..ccde73061 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt @@ -723,6 +723,56 @@ class QuicConnection( /** Caller must hold [lock]. */ internal fun streamByIdLocked(id: Long): QuicStream? = streams[id] + /** + * Step 6 of `quic/plans/2026-05-04-control-frame-retransmit.md`: + * dispatch the tokens of declared-lost packets to the matching + * `pending*` field, applying the supersede check from neqo's + * `fc.rs::frame_lost`. The supersede invariant: only re-flag for + * retransmit if the lost value still equals the connection's + * current advertised cap. If a higher extension has since gone + * out, the older lost frame is irrelevant — the newer value + * supersedes it. + * + * Called by the parser's AckFrame handler after + * [com.vitorpamplona.quic.connection.recovery.QuicLossDetection.detectAndRemoveLost] + * returns the lost set. Caller must hold [lock]. + */ + internal fun onTokensLost(tokens: List) { + for (token in tokens) { + when (token) { + com.vitorpamplona.quic.connection.recovery.RecoveryToken.Ack -> { + // ACK frames are not retransmittable per RFC 9000 + // §13.2.1; the peer's own ACKs cover newer ranges. + } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxStreamsUni -> { + if (token.maxStreams == advertisedMaxStreamsUni) { + pendingMaxStreamsUni = token.maxStreams + } + } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxStreamsBidi -> { + if (token.maxStreams == advertisedMaxStreamsBidi) { + pendingMaxStreamsBidi = token.maxStreams + } + } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxData -> { + if (token.maxData == advertisedMaxData) { + pendingMaxData = token.maxData + } + } + + is com.vitorpamplona.quic.connection.recovery.RecoveryToken.MaxStreamData -> { + val stream = streamByIdLocked(token.streamId) ?: continue + if (token.maxData == stream.receiveLimit) { + pendingMaxStreamData[token.streamId] = token.maxData + } + } + } + } + } + companion object { /** * Bound on the inbound datagram queue depth. RFC 9221 datagrams are diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt index 1b797c3eb..450c3c896 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt @@ -202,10 +202,12 @@ private fun dispatchFrames( largestAckedPn = largestAckedPn, nowMs = nowMillis, ) - - // Step 6 will dispatch tokens here. Drop for now. - @Suppress("UNUSED_VARIABLE") - val unusedForStep6 = lost + // Step 6: dispatch each lost packet's tokens to the + // matching pending* field. The supersede check (lost + // value still == advertised) lives inside onTokensLost. + for (lostPacket in lost) { + conn.onTokensLost(lostPacket.tokens) + } } } diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/OnTokensLostTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/OnTokensLostTest.kt new file mode 100644 index 000000000..f4978ff1b --- /dev/null +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/OnTokensLostTest.kt @@ -0,0 +1,199 @@ +/* + * Copyright (c) 2025 Vitor Pamplona + * + * Permission is hereby granted, free of charge, to any person obtaining a copy of + * this software and associated documentation files (the "Software"), to deal in + * the Software without restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the + * Software, and to permit persons to whom the Software is furnished to do so, + * subject to the following conditions: + * + * The above copyright notice and this permission notice shall be included in all + * copies or substantial portions of the Software. + * + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR + * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS + * FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR + * COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN + * AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION + * WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. + */ +package com.vitorpamplona.quic.connection + +import com.vitorpamplona.quic.connection.recovery.RecoveryToken +import kotlinx.coroutines.runBlocking +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull + +/** + * Step 6 of `quic/plans/2026-05-04-control-frame-retransmit.md`: + * `QuicConnection.onTokensLost` dispatches lost-packet tokens to + * the matching `pending*` field with the supersede check (lost + * value must equal current advertised). Mirrors neqo's + * `streams.rs::lost` dispatch + `fc.rs::frame_lost` flag-set logic. + */ +class OnTokensLostTest { + private fun newConn(): QuicConnection = + QuicConnection( + serverName = "example.test", + config = QuicConnectionConfig(), + tlsCertificateValidator = + com.vitorpamplona.quic.tls + .PermissiveCertificateValidator(), + ) + + @Test + fun ackToken_doesNotPopulateAnyPending() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.onTokensLost(listOf(RecoveryToken.Ack)) + } finally { + conn.lock.unlock() + } + assertNull(conn.pendingMaxStreamsUni) + assertNull(conn.pendingMaxStreamsBidi) + assertNull(conn.pendingMaxData) + assertEquals(emptyMap(), conn.pendingMaxStreamData) + } + + @Test + fun lostMaxStreamsUni_matchingAdvertised_setsPending() = + runBlocking { + val conn = newConn() + // Simulate the writer having advertised a higher cap. + conn.lock.lock() + try { + conn.advertisedMaxStreamsUni = 150L + conn.onTokensLost(listOf(RecoveryToken.MaxStreamsUni(maxStreams = 150L))) + } finally { + conn.lock.unlock() + } + assertEquals(150L, conn.pendingMaxStreamsUni) + } + + @Test + fun lostMaxStreamsUni_supersededByHigherEmit_isDropped() = + runBlocking { + val conn = newConn() + // The writer has since advertised a higher cap (200) than + // the value carried by the lost token (150). The lost + // frame is irrelevant — re-emitting 150 would not extend + // the cap. neqo's fc.rs line 322 supersede check. + conn.lock.lock() + try { + conn.advertisedMaxStreamsUni = 200L + conn.onTokensLost(listOf(RecoveryToken.MaxStreamsUni(maxStreams = 150L))) + } finally { + conn.lock.unlock() + } + assertNull(conn.pendingMaxStreamsUni, "stale lost extension must not be re-emitted") + } + + @Test + fun lostMaxStreamsBidi_matchingAdvertised_setsPending() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.advertisedMaxStreamsBidi = 200L + conn.onTokensLost(listOf(RecoveryToken.MaxStreamsBidi(maxStreams = 200L))) + } finally { + conn.lock.unlock() + } + assertEquals(200L, conn.pendingMaxStreamsBidi) + } + + @Test + fun lostMaxData_matchingAdvertised_setsPending() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.advertisedMaxData = 1_000_000L + conn.onTokensLost(listOf(RecoveryToken.MaxData(maxData = 1_000_000L))) + } finally { + conn.lock.unlock() + } + assertEquals(1_000_000L, conn.pendingMaxData) + } + + @Test + fun lostMaxData_supersededIsDropped() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.advertisedMaxData = 2_000_000L + conn.onTokensLost(listOf(RecoveryToken.MaxData(maxData = 1_000_000L))) + } finally { + conn.lock.unlock() + } + assertNull(conn.pendingMaxData) + } + + @Test + fun lostMaxStreamData_unknownStream_dropped() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.onTokensLost( + listOf(RecoveryToken.MaxStreamData(streamId = 999L, maxData = 1024L)), + ) + } finally { + conn.lock.unlock() + } + // No stream with id 999 exists ⇒ token is dropped silently. + assertEquals(emptyMap(), conn.pendingMaxStreamData) + } + + @Test + fun multipleLostTokens_dispatchAll() = + runBlocking { + val conn = newConn() + conn.lock.lock() + try { + conn.advertisedMaxStreamsUni = 150L + conn.advertisedMaxStreamsBidi = 200L + conn.advertisedMaxData = 5_000_000L + conn.onTokensLost( + listOf( + RecoveryToken.Ack, + RecoveryToken.MaxStreamsUni(maxStreams = 150L), + RecoveryToken.MaxStreamsBidi(maxStreams = 200L), + RecoveryToken.MaxData(maxData = 5_000_000L), + ), + ) + } finally { + conn.lock.unlock() + } + assertEquals(150L, conn.pendingMaxStreamsUni) + assertEquals(200L, conn.pendingMaxStreamsBidi) + assertEquals(5_000_000L, conn.pendingMaxData) + } + + @Test + fun lostTokensFromMultiplePackets_unionInPending() = + runBlocking { + // Two SentPackets each lost; their tokens are dispatched + // sequentially. Each frame type's pending field holds at + // most one value (the last setter wins; the supersede + // check filters older losses). + val conn = newConn() + conn.lock.lock() + try { + conn.advertisedMaxStreamsUni = 200L + // First lost packet had MaxStreamsUni(150) — stale, dropped. + conn.onTokensLost(listOf(RecoveryToken.MaxStreamsUni(maxStreams = 150L))) + assertNull(conn.pendingMaxStreamsUni) + // Second lost packet had MaxStreamsUni(200) — current, set. + conn.onTokensLost(listOf(RecoveryToken.MaxStreamsUni(maxStreams = 200L))) + assertEquals(200L, conn.pendingMaxStreamsUni) + } finally { + conn.lock.unlock() + } + } +}