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 4ba2ef011..46c86962c 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt @@ -762,6 +762,38 @@ class QuicConnection( EncryptionLevel.APPLICATION -> application } + /** + * RFC 9002 §6.2.4 PTO probe — spec-correct retransmit path. Move + * every byte currently sent-but-not-yet-ACK'd in the [level]'s + * CRYPTO send buffer back to its retransmit queue, so the next + * [com.vitorpamplona.quic.connection.drainOutbound] re-emits the + * same bytes (at the same offsets) inside a fresh CRYPTO frame on + * a new packet number. + * + * The driver calls this from its PTO branch when 1-RTT keys + * aren't yet installed — i.e. the handshake hasn't finished, so + * the only thing the peer could be missing is our ClientHello / + * ClientFinished. A bare PING is insufficient because if the + * server never saw our original Initial it has no DCID state to + * correlate a PING against (it'll be dropped). Retransmitting the + * CRYPTO actually advances the handshake. + * + * Idempotent: a second consecutive call is a no-op because the + * first call moved everything out of inFlight. Old `RecoveryToken.Crypto` + * entries in [LevelState.sentPackets] for the still-tracked + * original PNs remain harmless — when loss detection eventually + * declares them lost, [onTokensLost] re-runs `markLost` on the + * same offset/length range, which is itself idempotent (the bytes + * are already in retransmit or already ACK'd by then). + * + * Caller must hold [lock] (or call from inside an existing locked + * region — typically the driver's PTO branch under + * [QuicConnectionDriver.sendLoop]). + */ + internal fun requeueAllInflightCrypto(level: EncryptionLevel) { + levelState(level).cryptoSend.requeueAllInflight() + } + /** Caller must hold [lock]. Snapshot of streams for the driver's send loop. */ internal fun streamsLocked(): Map = streams diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt index 48f198607..001bb8b5b 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionDriver.kt @@ -145,12 +145,37 @@ class QuicConnectionDriver( Unit } if (woke == null) { - // PTO fired. Set pendingPing so the writer emits a - // PING on the next drain (RFC 9002 §6.2.4 probe - // packet). The peer's ACK feeds loss detection + - // retransmit (steps 5–6). + // PTO fired. RFC 9002 §6.2.4: the probe packet MUST + // be ack-eliciting at the encryption level with + // unacknowledged data, and SHOULD retransmit lost + // data rather than just emit a PING. + // + // Pre-1-RTT we have a concrete thing to retransmit: + // the unacknowledged ClientHello (at Initial) or + // ClientFinished (at Handshake). We requeue ALL + // currently-inflight CRYPTO bytes for the highest + // active pre-application level so the next drain's + // takeChunk emits a fresh CRYPTO frame at the + // original offset. The PING flag is still set as a + // belt-and-suspenders fallback — `collectHandshakeLevelFrames` + // emits the PING only when no CRYPTO retransmit ends + // up in the same frame list, so once we have CRYPTO + // to send we don't waste a frame on a redundant PING. + // + // Post-handshake (1-RTT installed) we retain the + // bare-PING behavior; STREAM retransmit is driven by + // packet-number-threshold loss detection from the + // ACK that the PING elicits, which is a richer signal + // than blindly requeueing every inflight byte across + // every open stream. connection.lock.withLock { connection.pendingPing = true + if (connection.application.sendProtection == null) { + val level = highestPreApplicationLevel(connection) + if (level != null) { + connection.requeueAllInflightCrypto(level) + } + } connection.consecutivePtoCount = (connection.consecutivePtoCount + 1).coerceAtMost(6) } @@ -223,3 +248,19 @@ class QuicConnectionDriver( private const val CLOSE_FLUSH_TIMEOUT_MILLIS = 250L } } + +/** + * Highest encryption level for which `conn` currently holds send keys + * AND hasn't yet discarded them, given that 1-RTT keys are NOT + * installed. Returns null when the level state has been completely + * cleared (e.g. CLOSED after a CONNECTION_CLOSE was sent). Mirrors the + * private helper in [com.vitorpamplona.quic.connection.QuicConnectionWriter] + * — kept in lockstep so the driver's PTO branch and the writer's PING + * placement target the same level. + */ +private fun highestPreApplicationLevel(conn: QuicConnection): EncryptionLevel? = + when { + conn.handshake.sendProtection != null -> EncryptionLevel.HANDSHAKE + conn.initial.sendProtection != null && !conn.initial.keysDiscarded -> EncryptionLevel.INITIAL + else -> null + } diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt index 71e3781fb..6bffeb129 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionWriter.kt @@ -395,18 +395,32 @@ private fun collectHandshakeLevelFrames( length = cryptoChunk.data.size.toLong(), ) } - // RFC 9002 §6.2.4: PTO probe MUST be emitted at the encryption level - // that has unacknowledged data. `buildApplicationPacket` consumes - // [pendingPing] preferentially when 1-RTT keys exist; if they don't, - // the probe lives or dies here at Handshake / Initial level. Pre-fix - // the flag was set but never honored pre-handshake — the connection - // sat silent through every PTO, never retransmitting the ClientHello - // even when the first datagram was lost. This caused the - // "1-packet-on-the-wire-then-timeout" symptom against aioquic in the - // quic-interop-runner sim where the first ClientHello was dropped - // because the server hadn't finished startup yet. - if (conn.pendingPing && conn.application.sendProtection == null && level == highestPreApplicationLevel(conn)) { - frames += PingFrame + // RFC 9002 §6.2.4: PTO probe MUST be ack-eliciting at the + // encryption level with unacknowledged data. `buildApplicationPacket` + // consumes [pendingPing] when 1-RTT keys exist; pre-1-RTT we + // honor it here at the highest active level (Handshake > Initial). + // + // The driver's PTO branch also calls + // [QuicConnection.requeueAllInflightCrypto], which moves any + // unacknowledged ClientHello / ClientFinished bytes back onto + // the cryptoSend retransmit queue — so the takeChunk above + // already produced a CRYPTO frame in that case, and the PING + // would be redundant. We only emit a bare PING when there's no + // CRYPTO retransmit available (e.g. the unacknowledged data was + // already sitting in cryptoSend's retransmit queue from a + // previous PTO and got drained, or the original Initial was + // never sent at all). This preserves the "send something + // ack-eliciting on every PTO" contract without wasting a frame. + if (conn.pendingPing && + conn.application.sendProtection == null && + level == highestPreApplicationLevel(conn) + ) { + if (frames.none { it is CryptoFrame }) { + frames += PingFrame + } + // Either the CRYPTO frame already covers the ack-eliciting + // requirement at this level, or we just appended a PING. + // Clear the flag so the next drain doesn't double-fire. conn.pendingPing = false } if (frames.isEmpty()) return null @@ -414,16 +428,16 @@ private fun collectHandshakeLevelFrames( } /** - * The highest encryption level for which we currently hold send keys, given - * that 1-RTT keys are NOT yet installed. Used by [collectHandshakeLevelFrames] - * to decide where a PTO PING probe should ride when the application level - * isn't usable yet. + * Highest encryption level for which we currently hold send keys, given + * 1-RTT keys are NOT installed. Used by [collectHandshakeLevelFrames] + * to place a PTO PING (RFC 9002 §6.2.4) at the right level when the + * application path can't carry it. */ private fun highestPreApplicationLevel(conn: QuicConnection): EncryptionLevel = when { conn.handshake.sendProtection != null -> EncryptionLevel.HANDSHAKE conn.initial.sendProtection != null && !conn.initial.keysDiscarded -> EncryptionLevel.INITIAL - else -> EncryptionLevel.INITIAL // fallback; collectHandshakeLevelFrames already early-returns on no keys + else -> EncryptionLevel.INITIAL // collectHandshakeLevelFrames already returned null on no keys } /** diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/SendBuffer.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/SendBuffer.kt index 2b9d7807d..ec56ee9ad 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/SendBuffer.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/SendBuffer.kt @@ -315,6 +315,54 @@ class SendBuffer( } } + /** + * Re-queue every byte range currently sent-but-not-yet-ACK'd for + * retransmission. The writer's next [takeChunk] drains the + * retransmit queue before any fresh bytes, at the original + * offsets — so the peer sees the same CRYPTO/STREAM bytes again + * with the same offset, only at a new packet number. + * + * Used by the RFC 9002 §6.2.4 PTO probe path: when the timer + * fires before the handshake completes, the client MUST + * retransmit the unacknowledged ClientHello (the CRYPTO bytes + * sitting in [inFlight] for the Initial level). A PING alone + * isn't enough — the peer may never have seen the original + * datagram and therefore has no state to correlate the PING + * against. + * + * Idempotent: ranges already in [retransmit] are not affected + * (only [inFlight] is walked); calling twice in a row is a no-op + * the second time because the first call moved everything out of + * [inFlight]. FIN bit is preserved per range — a lost FIN-bearing + * range will re-emit FIN. + * + * Best-effort buffers ([bestEffort] = true) drop the inflight + * ranges instead of retransmitting them, matching [markLost]'s + * semantics. The data buffer can compact as if those ranges had + * been ACK'd. + */ + fun requeueAllInflight() { + synchronized(this) { + if (inFlight.isEmpty()) return + if (bestEffort) { + inFlight.clear() + advanceFlushedFloorIfPossible() + return + } + // Move every inflight range to the retransmit queue, + // preserving offset order (inFlight is sorted by offset + // ascending, so addLast preserves sort within retransmit + // for these new entries — though retransmit is a FIFO + // and doesn't strictly require sorted order, takeChunk + // pops front-first regardless). + for (r in inFlight) { + if (r.fin && !_finAcked) _finSent = false + retransmit.addLast(r) + } + inFlight.clear() + } + } + /** * Disposition for a range that overlapped a [markAcked] / [markLost] * range. ACK drops the range and latches `_finAcked` if the FIN diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PtoCryptoRetransmitTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PtoCryptoRetransmitTest.kt new file mode 100644 index 000000000..60f39234a --- /dev/null +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/PtoCryptoRetransmitTest.kt @@ -0,0 +1,209 @@ +/* + * 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 com.vitorpamplona.quic.frame.CryptoFrame +import com.vitorpamplona.quic.frame.Frame +import com.vitorpamplona.quic.packet.LongHeaderPacket +import kotlinx.coroutines.runBlocking +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertTrue + +/** + * RFC 9002 §6.2.4 spec-correct PTO probe behavior: when the timer fires + * pre-handshake the client MUST send an ack-eliciting packet at the + * encryption level with unacknowledged data, and SHOULD retransmit the + * lost data rather than emit a bare PING. + * + * Regression scenario from the quic-interop-runner ns-3 transfer test + * against aioquic: the first ClientHello datagram is dropped because + * the simulated server isn't fully started at t≈0.5s. Our PTO at + * t≈1.5s previously sent only a PING with the same DCID; the server + * had no state for that DCID (never saw the original Initial), so + * the PING was silently ignored and the connection died. + * + * The fix has two cooperating pieces: + * 1. [com.vitorpamplona.quic.stream.SendBuffer.requeueAllInflight] + * moves every sent-but-not-ACK'd byte range from the inflight + * list back onto the retransmit queue. + * 2. [QuicConnection.requeueAllInflightCrypto] exposes that to the + * driver, which calls it from the PTO branch when 1-RTT keys + * aren't installed yet (i.e. handshake not done). + * + * After the call, the next [drainOutbound] naturally emits a CRYPTO + * frame at the original offset — the server actually sees a fresh + * ClientHello attempt and can advance the handshake. + */ +class PtoCryptoRetransmitTest { + @Test + fun ptoPreHandshake_retransmitsClientHelloCryptoBytes() = + runBlocking { + val client = newClientWithStartedHandshake() + + // First drain — emits the ClientHello inside an Initial datagram. + val firstDrain = drainOutbound(client, nowMillis = 1L) + assertNotNull(firstDrain, "first drain must produce the ClientHello datagram") + assertTrue( + firstDrain.size >= 1200, + "RFC 9000 §14.1 — Initial-bearing datagram pads to ≥ 1200 bytes (got ${firstDrain.size})", + ) + + // Capture the original ClientHello bytes and offset by + // pulling them out of the SentPacket bookkeeping. + val firstSent = + client.initial.sentPackets.entries.firstOrNull { entry -> + entry.value.tokens.any { it is RecoveryToken.Crypto } + } + assertNotNull(firstSent, "Initial-level SentPacket must carry a Crypto token after first drain") + val firstCrypto = + firstSent.value.tokens + .filterIsInstance() + .single() + assertEquals(0L, firstCrypto.offset, "ClientHello starts at offset 0") + assertTrue(firstCrypto.length > 0L, "ClientHello has non-zero length") + + // Sanity: the second drain (immediately, no PTO yet) must + // produce nothing — bytes are inflight, not unsent. + val emptyDrain = drainOutbound(client, nowMillis = 2L) + assertTrue( + emptyDrain == null || emptyDrain.isEmpty(), + "no PTO yet, no fresh data → second drain must be empty (got ${emptyDrain?.size} bytes)", + ) + + // Simulate the driver's PTO branch firing: requeue the + // unacknowledged CRYPTO and set the PING flag. + client.lock.lock() + try { + client.requeueAllInflightCrypto(EncryptionLevel.INITIAL) + client.pendingPing = true + } finally { + client.lock.unlock() + } + + // Next drain must emit a fresh Initial packet carrying + // the ClientHello CRYPTO at the original offset. + val ptoDrain = drainOutbound(client, nowMillis = 3L) + assertNotNull(ptoDrain, "PTO drain must produce a retransmit datagram") + assertTrue( + ptoDrain.size >= 1200, + "RFC 9000 §14.1 — Initial-bearing datagram still pads to ≥ 1200 bytes on PTO retransmit (got ${ptoDrain.size})", + ) + + // The retransmit packet must carry a Crypto token at the + // original offset and length — that's how we know the + // CRYPTO frame went out (not a PING-only probe). + val replayEntries = + client.initial.sentPackets.entries + .filter { it.key != firstSent.key } + val replaySent = + replayEntries.firstOrNull { entry -> + entry.value.tokens.any { it is RecoveryToken.Crypto } + } + assertNotNull( + replaySent, + "PTO drain must produce a fresh Initial SentPacket carrying Crypto " + + "(saw ${replayEntries.map { it.value.tokens.map { t -> t::class.simpleName } }})", + ) + val replayCrypto = + replaySent.value.tokens + .filterIsInstance() + .single() + assertEquals(EncryptionLevel.INITIAL, replayCrypto.level) + assertEquals(firstCrypto.offset, replayCrypto.offset, "PTO retransmit replays original offset") + assertEquals(firstCrypto.length, replayCrypto.length, "PTO retransmit replays original length") + + // Decode the retransmit packet and assert the CRYPTO + // frame's payload bytes match the original. + val firstFrames = decodeInitialFrames(firstDrain, client) + val firstHello = + firstFrames.filterIsInstance().firstOrNull { + it.offset == firstCrypto.offset && it.data.size.toLong() == firstCrypto.length + } + assertNotNull(firstHello, "first drain's Initial must contain the ClientHello CRYPTO") + + val ptoFrames = decodeInitialFrames(ptoDrain, client) + val replayHello = + ptoFrames.filterIsInstance().firstOrNull { + it.offset == firstCrypto.offset + } + assertNotNull( + replayHello, + "PTO drain's Initial must contain a CRYPTO frame at offset 0 — bare PING is not enough " + + "(saw frames ${ptoFrames.map { it::class.simpleName }})", + ) + assertTrue( + replayHello.data.contentEquals(firstHello.data), + "PTO retransmit must carry the same ClientHello bytes (size first=${firstHello.data.size} replay=${replayHello.data.size})", + ) + + // pendingPing should be cleared since the CRYPTO frame + // satisfied the ack-eliciting requirement at the level. + assertEquals( + false, + client.pendingPing, + "pendingPing must be consumed once the PTO emit went out (CRYPTO covered the probe)", + ) + } + + /** + * Decode an Initial-level long-header packet from a freshly-drained + * datagram and return its frames. We open with the client's own + * SEND protection (client-side keys) — symmetric AEAD opens with + * the same key/iv that sealed it. + */ + private fun decodeInitialFrames( + datagram: ByteArray, + client: QuicConnection, + ): List { + val send = client.initial.sendProtection!! + val parsed = + LongHeaderPacket.parseAndDecrypt( + bytes = datagram, + offset = 0, + aead = send.aead, + key = send.key, + iv = send.iv, + hp = send.hp, + hpKey = send.hpKey, + largestReceivedInSpace = -1L, + ) + assertNotNull(parsed, "must parse the Initial packet header (datagram size=${datagram.size})") + return com.vitorpamplona.quic.frame + .decodeFrames(parsed.packet.payload) + } + + private fun newClientWithStartedHandshake(): QuicConnection = + runBlocking { + val client = + QuicConnection( + serverName = "example.test", + config = QuicConnectionConfig(), + tlsCertificateValidator = + com.vitorpamplona.quic.tls + .PermissiveCertificateValidator(), + ) + client.start() + client + } +}