From 6e054583a977dfaf5c092cbcc41fd94c2f38c967 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 9 May 2026 04:44:04 +0000 Subject: [PATCH] =?UTF-8?q?fix(quic):=20RFC=209000=20=C2=A74.1=20connectio?= =?UTF-8?q?n-level=20inbound=20flow-control=20enforcement?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pre-fix the receiver enforced ONLY per-stream `receiveLimit`. The aggregate cap (`initial_max_data` / latest MAX_DATA) was advertised and respected on the SEND side but never checked on the RECEIVE side — a peer that opened many streams and pushed each up to its per-stream cap could collectively exceed `initial_max_data` without us closing. Implementation: * `QuicStream.receiveHighestOffset`: per-stream high-water mark of "largest stream offset received" — the spec quantity that gets summed across streams for the connection-level check. * `QuicConnection.connectionInboundOffsetSum`: running sum across all streams of `receiveHighestOffset`. Updated in the parser at STREAM and RESET_STREAM frame ingest time. * Parser closes with FLOW_CONTROL_ERROR whenever the running sum exceeds `advertisedMaxData`. * RESET_STREAM final-size accounting: per §4.5 the finalSize counts toward connection-level flow control even though no STREAM frame delivered those bytes — a peer that resets a stream at a finalSize larger than it had previously sent gets the extra bytes added to the running sum at RESET arrival. Different from the writer's MAX_DATA threshold logic (which uses the cheaper contiguous-end approximation): the receiver-side enforcement needs the strict spec quantity to close before bookkeeping diverges. Test: `ConnectionLevelFlowControlTest` (5 cases) covers below-cap pass, over-cap close, retransmit-doesn't-double-count, RESET_STREAM finalSize counts toward limit, fresh-handshake invariants. https://claude.ai/code/session_01XGmmaVuy2wsSZQx2AqN2ik --- .../quic/connection/QuicConnection.kt | 18 +++ .../quic/connection/QuicConnectionParser.kt | 38 +++++ .../vitorpamplona/quic/stream/QuicStream.kt | 14 ++ .../ConnectionLevelFlowControlTest.kt | 137 ++++++++++++++++++ 4 files changed, 207 insertions(+) create mode 100644 quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ConnectionLevelFlowControlTest.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 dee3f687d..7ae599cae 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnection.kt @@ -481,6 +481,24 @@ class QuicConnection( */ internal var advertisedMaxData: Long = config.initialMaxData + /** + * RFC 9000 §4.1 connection-level inbound flow-control accounting. + * Sum across all streams (live + retired) of the largest stream + * offset ever observed on an inbound STREAM / RESET_STREAM frame. + * The receiver MUST close the connection with FLOW_CONTROL_ERROR + * when this sum exceeds [advertisedMaxData]; that's enforced in + * the parser at frame-ingest time. + * + * Different from the writer's `totalRecvAdvanced` which uses + * `stream.receive.contiguousEnd()` (a slight under-count) for + * MAX_DATA threshold logic. Here we track the strict spec + * quantity — necessary to close before bookkeeping diverges + * if the peer sends data with gaps that pushes total received + * past the limit. + */ + @Volatile + internal var connectionInboundOffsetSum: Long = 0L + /** * Cumulative count of peer-initiated unidirectional streams we've * accepted (incremented in [getOrCreatePeerStreamLocked]). Compared 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 946f716d8..3778bd5fb 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/connection/QuicConnectionParser.kt @@ -675,6 +675,25 @@ private fun dispatchFrames( ) return } + // RFC 9000 §4.1 connection-level enforcement: the SUM + // across all streams of "largest stream offset received" + // MUST NOT exceed the most recently advertised MAX_DATA. + // We update [conn.connectionInboundOffsetSum] only when + // this frame's end-offset advances the per-stream + // high-water mark, so duplicate / reordered frames don't + // double-count. + if (frameEnd > stream.receiveHighestOffset) { + val delta = frameEnd - stream.receiveHighestOffset + stream.receiveHighestOffset = frameEnd + conn.connectionInboundOffsetSum += delta + if (conn.connectionInboundOffsetSum > conn.advertisedMaxData) { + conn.markClosedExternally( + "FLOW_CONTROL_ERROR: peer exceeded connection MAX_DATA " + + "(${conn.connectionInboundOffsetSum} > ${conn.advertisedMaxData})", + ) + return + } + } // RFC 9000 §4.5: enforce final-size invariants. The // [com.vitorpamplona.quic.stream.ReceiveBuffer.insert] // surface returns a typed result; map any non-OK result @@ -830,6 +849,25 @@ private fun dispatchFrames( return } } + // RFC 9000 §4.5: a RESET_STREAM's finalSize counts toward + // connection-level flow control even though no STREAM frame + // ever delivered those bytes — it represents what the peer + // committed to send before aborting. Top up + // [conn.connectionInboundOffsetSum] when finalSize advances + // beyond what STREAM frames already counted, then re-check + // against the connection cap. + if (target != null && frame.finalSize > target.receiveHighestOffset) { + val delta = frame.finalSize - target.receiveHighestOffset + target.receiveHighestOffset = frame.finalSize + conn.connectionInboundOffsetSum += delta + if (conn.connectionInboundOffsetSum > conn.advertisedMaxData) { + conn.markClosedExternally( + "FLOW_CONTROL_ERROR: peer RESET_STREAM finalSize pushed connection past MAX_DATA " + + "(${conn.connectionInboundOffsetSum} > ${conn.advertisedMaxData})", + ) + return + } + } // Mark the peer's stream aborted and close our read side; the // application sees a truncated incoming flow. target?.closeIncoming() diff --git a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/QuicStream.kt b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/QuicStream.kt index 97974fec1..e2c54d34d 100644 --- a/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/QuicStream.kt +++ b/quic/src/commonMain/kotlin/com/vitorpamplona/quic/stream/QuicStream.kt @@ -128,6 +128,20 @@ class QuicStream( var receiveLimit: Long = 0L internal set + /** + * RFC 9000 §4.1: highest stream offset (offset + length) ever seen + * on an inbound STREAM or RESET_STREAM frame. Used by + * [QuicConnection]'s connection-level flow-control accounting — + * the spec requires the receiver to compare the SUM across all + * streams of "largest received offset" against the advertised + * `initial_max_data` / latest MAX_DATA, NOT the contiguous read + * frontier. The writer's MAX_DATA threshold logic still uses the + * cheaper contiguous-end approximation; this field is purely the + * enforcement signal on receive. + */ + @Volatile + internal var receiveHighestOffset: Long = 0L + /** * Marker the parser sets whenever [receive.contiguousEnd] advances; the * writer's appendFlowControlUpdates consumes it to skip streams that diff --git a/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ConnectionLevelFlowControlTest.kt b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ConnectionLevelFlowControlTest.kt new file mode 100644 index 000000000..1a4f5549c --- /dev/null +++ b/quic/src/commonTest/kotlin/com/vitorpamplona/quic/connection/ConnectionLevelFlowControlTest.kt @@ -0,0 +1,137 @@ +/* + * 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.frame.ResetStreamFrame +import com.vitorpamplona.quic.frame.StreamFrame +import kotlinx.coroutines.runBlocking +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertTrue + +/** + * RFC 9000 §4.1: connection-level inbound flow-control enforcement. + * + * "All data sent in STREAM frames counts toward this limit. The sum of + * the largest received offsets on all streams — including streams in + * the Reset Recvd state — MUST NOT exceed the value advertised by a + * receiver." + * + * Pre-fix the receiver enforced ONLY per-stream `receiveLimit` (each + * stream individually couldn't push past `initial_max_stream_data_*`). + * The aggregate connection-level cap (`initial_max_data`) was advertised + * AND respected on the SEND side, but never checked on the RECEIVE + * side. A peer that opened many streams and pushed each one up to its + * per-stream cap could collectively exceed `initial_max_data` without + * us closing. + */ +class ConnectionLevelFlowControlTest { + private fun connectClient( + maxData: Long = 16 * 1024, + maxStreamData: Long = 8 * 1024, + maxStreamsUni: Long = 16, + ): Pair = + newConnectedClient( + maxData = maxData, + maxStreamData = maxStreamData, + maxStreamsUni = maxStreamsUni, + maxStreamsBidi = 16, + ) + + @Test + fun connection_recv_sum_below_max_data_accepted() { + val (client, pipe) = connectClient(maxData = 16 * 1024, maxStreamData = 8 * 1024) + // Two server-uni streams (id 3, 7) each carrying 4 KiB. Sum = 8 KiB + // < 16 KiB connection cap. Should sail through. + val frame1 = StreamFrame(streamId = 3L, offset = 0L, data = ByteArray(4096), fin = false) + val frame2 = StreamFrame(streamId = 7L, offset = 0L, data = ByteArray(4096), fin = false) + val datagram = pipe.buildServerApplicationDatagram(listOf(frame1, frame2))!! + feedDatagram(client, datagram, nowMillis = 0L) + assertEquals(QuicConnection.Status.CONNECTED, client.status) + assertEquals(8L * 1024, client.connectionInboundOffsetSum) + } + + @Test + fun connection_recv_sum_exceeding_max_data_closes() { + val (client, pipe) = connectClient(maxData = 16 * 1024, maxStreamData = 16 * 1024) + // Three server-uni streams each shipping 6 KiB → 18 KiB > 16 KiB cap. + // Per-stream limit (16 KiB) is fine; the aggregate breaches MAX_DATA. + val frame1 = StreamFrame(streamId = 3L, offset = 0L, data = ByteArray(6 * 1024), fin = false) + val frame2 = StreamFrame(streamId = 7L, offset = 0L, data = ByteArray(6 * 1024), fin = false) + val frame3 = StreamFrame(streamId = 11L, offset = 0L, data = ByteArray(6 * 1024), fin = false) + val datagram = pipe.buildServerApplicationDatagram(listOf(frame1, frame2, frame3))!! + feedDatagram(client, datagram, nowMillis = 0L) + assertEquals(QuicConnection.Status.CLOSED, client.status) + val reason = client.closeReason + assertNotNull(reason) + assertTrue(reason.contains("FLOW_CONTROL_ERROR"), "expected FLOW_CONTROL_ERROR, got: $reason") + assertTrue(reason.contains("MAX_DATA")) + } + + @Test + fun retransmitted_stream_frame_does_not_double_count() { + val (client, pipe) = connectClient(maxData = 16 * 1024, maxStreamData = 8 * 1024) + val frame = StreamFrame(streamId = 3L, offset = 0L, data = ByteArray(4 * 1024), fin = false) + // First arrival. + feedDatagram(client, pipe.buildServerApplicationDatagram(listOf(frame))!!, nowMillis = 0L) + val sumAfterFirst = client.connectionInboundOffsetSum + assertEquals(4L * 1024, sumAfterFirst) + // Identical retransmit — does NOT advance the per-stream high-water, + // so connection counter stays put. + feedDatagram(client, pipe.buildServerApplicationDatagram(listOf(frame))!!, nowMillis = 0L) + assertEquals(sumAfterFirst, client.connectionInboundOffsetSum) + assertEquals(QuicConnection.Status.CONNECTED, client.status) + } + + @Test + fun reset_stream_final_size_counts_toward_connection_limit() { + val (client, pipe) = connectClient(maxData = 16 * 1024, maxStreamData = 16 * 1024) + // Stream 3 ships 4 KiB then resets at finalSize = 20 KiB. The + // RESET_STREAM commits the peer to those 20 KiB even though we + // haven't seen them on the wire — they count toward MAX_DATA per + // §4.5. 20 KiB on a single stream alone breaches the 16 KiB + // connection cap, so we MUST close. + val streamData = StreamFrame(streamId = 3L, offset = 0L, data = ByteArray(4 * 1024), fin = false) + val reset = ResetStreamFrame(streamId = 3L, applicationErrorCode = 0L, finalSize = 20L * 1024) + feedDatagram(client, pipe.buildServerApplicationDatagram(listOf(streamData, reset))!!, nowMillis = 0L) + assertEquals(QuicConnection.Status.CLOSED, client.status) + val reason = client.closeReason + assertNotNull(reason) + assertTrue(reason.contains("FLOW_CONTROL_ERROR"), "expected FLOW_CONTROL_ERROR, got: $reason") + } + + @Test + fun handshake_fixture_starts_with_zero_inbound_offset_sum() { + // Sanity check on the bookkeeping: a fresh handshake-only + // connection has not received any STREAM bytes, so the sum is 0 + // even though [advertisedMaxData] is the configured initial. + val (client, pipe) = connectClient() + runBlocking { + // Make sure pipe-driven handshake completes to CONNECTED before + // we read the field — otherwise we'd be looking at the value + // mid-handshake with no peer streams in play. + pipe.drive(maxRounds = 4) + } + assertEquals(0L, client.connectionInboundOffsetSum) + assertTrue(client.advertisedMaxData > 0L) + } +}