Merge pull request #2751 from vitorpamplona/claude/stream-priority-followup-5ugk0
Implement stream priority scheduling for QUIC writer
This commit is contained in:
@@ -1,151 +0,0 @@
|
||||
# Stream priority for moq-lite group uni streams (T11.3 follow-up)
|
||||
|
||||
**Status**: deferred — flagged in T11 (drop `bestEffort=true`), not landed.
|
||||
|
||||
## Why
|
||||
|
||||
`T11` removed `bestEffort=true` from `MoqLiteSession.openGroupStream`
|
||||
because the QUIC contract it relied on (drop lost ranges silently
|
||||
without `RESET_STREAM`) created undeliverable streams that wedge
|
||||
peer reassembly buffers — the actual user-visible bug was a 30 s
|
||||
silent dropout per lost packet on lossy networks.
|
||||
|
||||
`bestEffort=true` was incidentally providing a "newer-groups-skip-
|
||||
queued-retransmits" effect: the writer never queued retransmits in
|
||||
the first place, so under congestion the loss budget naturally
|
||||
biased toward dropping older lost ranges. With `bestEffort=true`
|
||||
gone, all retransmits are now queued, and under sustained
|
||||
congestion the writer drains streams in `streamRoundRobinStart`
|
||||
order — which can mean the listener catches up on a stale group
|
||||
when a fresh one would be more useful.
|
||||
|
||||
The kixelated reference (`moq-rs`'s `Publisher::serve_group` in
|
||||
`rs/moq-lite/src/lite/publisher.rs:347-406`) addresses this by
|
||||
calling `stream.set_priority(priority.current())` on every group
|
||||
stream and biasing the writer toward higher-priority (newer)
|
||||
streams. We should do the same.
|
||||
|
||||
## Target shape
|
||||
|
||||
### `:quic` API
|
||||
|
||||
Add a stable priority hook to `QuicStream`:
|
||||
|
||||
```kotlin
|
||||
class QuicStream(...) {
|
||||
@Volatile
|
||||
var priority: Int = 0
|
||||
// Higher drains first under contention. Default 0 = unchanged
|
||||
// round-robin behaviour.
|
||||
}
|
||||
```
|
||||
|
||||
Expose at the WebTransport layer:
|
||||
|
||||
```kotlin
|
||||
interface WebTransportWriteStream {
|
||||
fun setPriority(priority: Int)
|
||||
}
|
||||
```
|
||||
|
||||
### `QuicConnectionWriter` — the load-bearing change
|
||||
|
||||
Today's send-frame loop (`QuicConnectionWriter.kt:411-416`):
|
||||
|
||||
```kotlin
|
||||
val streamsView = conn.streamsListLocked()
|
||||
val start = conn.streamRoundRobinStart % streamsView.size
|
||||
for (i in streamsView.indices) {
|
||||
val stream = streamsView[(start + i) % streamsView.size]
|
||||
// ...
|
||||
}
|
||||
```
|
||||
|
||||
Replace with priority-then-round-robin:
|
||||
|
||||
```kotlin
|
||||
val streamsView = conn.streamsListLocked()
|
||||
val sorted = streamsView.sortedByDescending { it.priority }
|
||||
val start = conn.streamRoundRobinStart % sorted.size
|
||||
for (i in sorted.indices) {
|
||||
val stream = sorted[(start + i) % sorted.size]
|
||||
// ...
|
||||
}
|
||||
```
|
||||
|
||||
Sort is stable; same-priority streams retain insertion order, so the
|
||||
existing round-robin behaviour holds within a priority tier. Higher-
|
||||
priority streams always come first in the iteration.
|
||||
|
||||
Cost: O(N log N) per drain pass, where N is active local-initiated
|
||||
streams. N is small (1–10 in the moq-lite audio path); the
|
||||
allocation of `sorted` per pass is the only real cost. If that
|
||||
shows up in profiling, switch to `kotlin.collections.IntArray`-
|
||||
backed indirect sort or maintain a priority-sorted list incrementally
|
||||
on `setPriority` calls.
|
||||
|
||||
### moq-lite wiring
|
||||
|
||||
`MoqLiteSession.openGroupStream:1022-1037`:
|
||||
|
||||
```kotlin
|
||||
internal suspend fun openGroupStream(
|
||||
subscribeId: Long,
|
||||
sequence: Long,
|
||||
): WebTransportWriteStream {
|
||||
val uni = transport.openUniStream()
|
||||
uni.setPriority(sequence.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||
uni.write(Varint.encode(MoqLiteDataType.Group.code))
|
||||
uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
|
||||
return uni
|
||||
}
|
||||
```
|
||||
|
||||
Newer groups have higher sequence → higher priority → drain first
|
||||
under congestion. The saturation conversion handles broadcasts that
|
||||
run long enough for `sequence` to exceed `Int.MAX_VALUE` (≈ 71
|
||||
years at our 1 group/sec production cadence; defensive only).
|
||||
|
||||
## Test
|
||||
|
||||
Add a unit test in `:quic` that builds a `QuicConnection`, opens
|
||||
two streams, sets `.priority = 0` on the first and `.priority = 1`
|
||||
on the second, queues bytes on both, and verifies the higher-
|
||||
priority stream's bytes hit the wire first. Doesn't need to fake
|
||||
real flow-control backpressure — pinning the iteration order via
|
||||
the writer's emitted-frames tape is sufficient to catch the
|
||||
regression case ("a later refactor accidentally re-introduces
|
||||
round-robin order").
|
||||
|
||||
A second test in `:nestsClient` that opens 3 group streams and
|
||||
asserts later-sequence streams drain before earlier ones is
|
||||
nice-to-have but secondary; the `:quic`-level test pins the load-
|
||||
bearing invariant.
|
||||
|
||||
## Why deferred
|
||||
|
||||
The change is small in lines but touches the QUIC writer's hot
|
||||
path. Rebasing a future `:quic` retransmit / pacing change onto a
|
||||
priority-sorted iteration order is doable but the diff conflicts
|
||||
get noisy. Better landed as a focused PR after the current
|
||||
interop-work series stabilises.
|
||||
|
||||
Risk profile:
|
||||
- **Bugs**: subtle starvation risk if a high-priority stream
|
||||
always has `streamRemaining > 0` — round-robin tiebreaker within
|
||||
a priority tier mitigates this for same-priority case, but the
|
||||
cross-tier case needs deliberate thought (do we want strict
|
||||
priority or weighted?).
|
||||
- **Performance**: per-pass sort allocation. Negligible for N≤10
|
||||
but worth measuring if N grows.
|
||||
- **Compat**: streams without an explicit priority default to 0,
|
||||
matching today's behaviour. Existing tests should pass unchanged.
|
||||
|
||||
## When to land
|
||||
|
||||
After the catalog interop series is fully verified (production
|
||||
audio with no degradation under realistic conditions), and ideally
|
||||
alongside a `:quic` perf review pass that touches the same code
|
||||
path. Not blocking on any audio bug today — `T11`'s reliable-
|
||||
delivery fix is the actual correctness change; this is the
|
||||
spec-aligned hardening.
|
||||
+7
@@ -1037,6 +1037,13 @@ class MoqLiteSession internal constructor(
|
||||
// 150 ms late still falls inside hang's default ~200 ms
|
||||
// jitter buffer) and avoids the dropout entirely.
|
||||
val uni = transport.openUniStream()
|
||||
// Mirror moq-rs `Publisher::serve_group`
|
||||
// (`rs/moq-lite/src/lite/publisher.rs`): newer groups get higher
|
||||
// priority so the QUIC writer drains fresh audio first when
|
||||
// retransmits queue up under loss. Saturating cast guards a
|
||||
// theoretical broadcast that runs long enough for `sequence` to
|
||||
// exceed Int.MAX_VALUE (≈ 71 years at 1 group/sec); defensive only.
|
||||
uni.setPriority(sequence.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||
uni.write(Varint.encode(MoqLiteDataType.Group.code))
|
||||
uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
|
||||
return uni
|
||||
|
||||
+4
@@ -162,6 +162,8 @@ class FakeBidiStream internal constructor(
|
||||
override suspend fun finish() {
|
||||
write.close()
|
||||
}
|
||||
|
||||
override fun setPriority(priority: Int) = Unit
|
||||
}
|
||||
|
||||
class FakeReadStream internal constructor(
|
||||
@@ -186,4 +188,6 @@ private class ChannelWriteStream(
|
||||
override suspend fun finish() {
|
||||
channel.close()
|
||||
}
|
||||
|
||||
override fun setPriority(priority: Int) = Unit
|
||||
}
|
||||
|
||||
+16
@@ -112,6 +112,22 @@ interface WebTransportWriteStream {
|
||||
|
||||
/** Half-close the write side (FIN). No further writes after this call. */
|
||||
suspend fun finish()
|
||||
|
||||
/**
|
||||
* Hint to the transport about this stream's drain priority relative
|
||||
* to other streams on the same session. Higher value drains first
|
||||
* under congestion; same-priority streams keep round-robin order.
|
||||
* Default 0 = unchanged round-robin behaviour.
|
||||
*
|
||||
* moq-lite uses this to bias the writer toward newer group streams
|
||||
* (sequence-numbered, fresher audio) so a backlog of retransmits on
|
||||
* an older group doesn't starve the listener of fresh frames. See
|
||||
* `Publisher::serve_group` in `rs/moq-lite/src/lite/publisher.rs`.
|
||||
*
|
||||
* Implementations that don't model priority (e.g. the in-memory
|
||||
* fake) MAY treat this as a no-op.
|
||||
*/
|
||||
fun setPriority(priority: Int)
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+18
@@ -365,6 +365,10 @@ private class QuicBidiStreamAdapter(
|
||||
stream.send.finish()
|
||||
driver.wakeup()
|
||||
}
|
||||
|
||||
override fun setPriority(priority: Int) {
|
||||
stream.priority = priority
|
||||
}
|
||||
}
|
||||
|
||||
private class QuicReadStreamAdapter(
|
||||
@@ -392,6 +396,10 @@ private class QuicUniWriteStreamAdapter(
|
||||
stream.send.finish()
|
||||
driver.wakeup()
|
||||
}
|
||||
|
||||
override fun setPriority(priority: Int) {
|
||||
stream.priority = priority
|
||||
}
|
||||
}
|
||||
|
||||
/** Adapter for a WT peer-initiated uni stream whose prefix has been stripped. */
|
||||
@@ -431,4 +439,14 @@ private class StrippedWtBidiStreamAdapter(
|
||||
?: error("peer-initiated bidi stream has no finish — demux didn't wire one")
|
||||
finish()
|
||||
}
|
||||
|
||||
/**
|
||||
* No-op: peer-initiated bidi streams arrive through the demux as a
|
||||
* [com.vitorpamplona.quic.webtransport.StrippedWtStream] which exposes
|
||||
* only `send`/`finish` closures, not the underlying [QuicStream]. The
|
||||
* moq-lite priority use case targets locally-opened uni group streams
|
||||
* only, so this path doesn't need to model priority — see the
|
||||
* [WebTransportWriteStream.setPriority] contract.
|
||||
*/
|
||||
override fun setPriority(priority: Int) = Unit
|
||||
}
|
||||
|
||||
+66
-38
@@ -410,46 +410,74 @@ private fun buildApplicationPacket(
|
||||
// insertion-ordered and stays in sync with the streams map.
|
||||
val streamsView = conn.streamsListLocked()
|
||||
if (streamsView.isNotEmpty()) {
|
||||
val start = conn.streamRoundRobinStart % streamsView.size
|
||||
for (i in streamsView.indices) {
|
||||
if (packetBudget <= 64) break
|
||||
val stream = streamsView[(start + i) % streamsView.size]
|
||||
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
|
||||
// Skip if both stream and connection have no credit; FIN-only
|
||||
// (zero-byte) chunks may still go through because they don't
|
||||
// consume credit.
|
||||
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
|
||||
val effectiveCap = minOf(streamRemaining, connBudget)
|
||||
val maxBytes =
|
||||
minOf(packetBudget - 32, effectiveCap.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||
val chunk = stream.send.takeChunk(maxBytes = maxBytes) ?: continue
|
||||
if (chunk.data.isNotEmpty() || chunk.fin) {
|
||||
frames +=
|
||||
StreamFrame(
|
||||
streamId = stream.streamId,
|
||||
offset = chunk.offset,
|
||||
data = chunk.data,
|
||||
fin = chunk.fin,
|
||||
explicitLength = true,
|
||||
)
|
||||
// Step C of the deferred-follow-ups pass: track this
|
||||
// STREAM emission so RFC 9002 retransmit can re-queue
|
||||
// the byte range on loss. SendBuffer.markLost (commit B)
|
||||
// moves the range from in-flight back to the retransmit
|
||||
// queue, and the next takeChunk replays it.
|
||||
tokens +=
|
||||
RecoveryToken.Stream(
|
||||
streamId = stream.streamId,
|
||||
offset = chunk.offset,
|
||||
length = chunk.data.size.toLong(),
|
||||
fin = chunk.fin,
|
||||
)
|
||||
packetBudget -= chunk.data.size + 32
|
||||
connBudget -= chunk.data.size
|
||||
conn.sendConnectionFlowConsumed += chunk.data.size
|
||||
// Strict priority across tiers, round-robin within each tier.
|
||||
// Higher-priority streams (e.g. moq-lite newer-sequence group
|
||||
// streams) ALWAYS drain ahead of lower-priority ones; the
|
||||
// rotating start-index only rotates among same-priority peers.
|
||||
// This is the spec-aligned shape — applying the rotation
|
||||
// globally over the sorted list would flip cross-tier order on
|
||||
// alternating drains, defeating the priority hint entirely.
|
||||
//
|
||||
// Default priority is 0; if every stream is at the default, all
|
||||
// streams form a single tier and iteration order matches the
|
||||
// pre-priority round-robin behaviour exactly.
|
||||
//
|
||||
// Cost: O(N log N) per drain pass plus one transient sorted
|
||||
// list. N is small (1–10 in the moq-lite audio path); if it
|
||||
// ever grows enough to matter, switch to an indirect index sort
|
||||
// or maintain an incrementally-sorted view on setPriority.
|
||||
val sorted =
|
||||
if (streamsView.size > 1) streamsView.sortedByDescending { it.priority } else streamsView
|
||||
val rotation = conn.streamRoundRobinStart
|
||||
var tierStart = 0
|
||||
outer@ while (tierStart < sorted.size) {
|
||||
// Walk the contiguous run of same-priority streams.
|
||||
val tierPriority = sorted[tierStart].priority
|
||||
var tierEnd = tierStart + 1
|
||||
while (tierEnd < sorted.size && sorted[tierEnd].priority == tierPriority) tierEnd++
|
||||
val tierSize = tierEnd - tierStart
|
||||
val tierRotation = if (tierSize > 1) rotation % tierSize else 0
|
||||
for (k in 0 until tierSize) {
|
||||
if (packetBudget <= 64) break@outer
|
||||
val stream = sorted[tierStart + ((tierRotation + k) % tierSize)]
|
||||
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
|
||||
// Skip if both stream and connection have no credit; FIN-only
|
||||
// (zero-byte) chunks may still go through because they don't
|
||||
// consume credit.
|
||||
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
|
||||
val effectiveCap = minOf(streamRemaining, connBudget)
|
||||
val maxBytes =
|
||||
minOf(packetBudget - 32, effectiveCap.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||
val chunk = stream.send.takeChunk(maxBytes = maxBytes) ?: continue
|
||||
if (chunk.data.isNotEmpty() || chunk.fin) {
|
||||
frames +=
|
||||
StreamFrame(
|
||||
streamId = stream.streamId,
|
||||
offset = chunk.offset,
|
||||
data = chunk.data,
|
||||
fin = chunk.fin,
|
||||
explicitLength = true,
|
||||
)
|
||||
// Step C of the deferred-follow-ups pass: track this
|
||||
// STREAM emission so RFC 9002 retransmit can re-queue
|
||||
// the byte range on loss. SendBuffer.markLost (commit B)
|
||||
// moves the range from in-flight back to the retransmit
|
||||
// queue, and the next takeChunk replays it.
|
||||
tokens +=
|
||||
RecoveryToken.Stream(
|
||||
streamId = stream.streamId,
|
||||
offset = chunk.offset,
|
||||
length = chunk.data.size.toLong(),
|
||||
fin = chunk.fin,
|
||||
)
|
||||
packetBudget -= chunk.data.size + 32
|
||||
connBudget -= chunk.data.size
|
||||
conn.sendConnectionFlowConsumed += chunk.data.size
|
||||
}
|
||||
}
|
||||
tierStart = tierEnd
|
||||
}
|
||||
conn.streamRoundRobinStart = (start + 1) % streamsView.size
|
||||
conn.streamRoundRobinStart = (rotation + 1) % streamsView.size
|
||||
}
|
||||
|
||||
if (frames.isEmpty()) return null
|
||||
|
||||
@@ -47,6 +47,27 @@ class QuicStream(
|
||||
val send = SendBuffer(bestEffort = bestEffort)
|
||||
val receive = ReceiveBuffer()
|
||||
|
||||
/**
|
||||
* Send-side scheduling priority. The connection writer's drain loop
|
||||
* iterates streams by descending priority; same-priority streams keep
|
||||
* their existing round-robin order. Higher value = drains first under
|
||||
* congestion. Default 0 matches pre-priority round-robin behaviour for
|
||||
* every existing call site.
|
||||
*
|
||||
* Used by moq-lite group streams: the publisher assigns each new group
|
||||
* a priority equal to its sequence number so that newer groups
|
||||
* (fresher audio) drain ahead of older ones when retransmits queue up
|
||||
* on a lossy link. Mirrors `Publisher::serve_group` in
|
||||
* `rs/moq-lite/src/lite/publisher.rs` (`stream.set_priority`).
|
||||
*
|
||||
* `@Volatile` because callers (e.g. moq-lite's openGroupStream)
|
||||
* assign from arbitrary coroutines while the writer reads it under
|
||||
* [com.vitorpamplona.quic.connection.QuicConnection.lock] during a
|
||||
* drain pass.
|
||||
*/
|
||||
@Volatile
|
||||
var priority: Int = 0
|
||||
|
||||
/**
|
||||
* Bytes received and confirmed contiguous, exposed as a flow to the consumer.
|
||||
*
|
||||
|
||||
@@ -327,6 +327,46 @@ class InMemoryQuicPipe(
|
||||
*/
|
||||
fun buildServerApplicationPacket(frames: List<com.vitorpamplona.quic.frame.Frame>): ByteArray? = buildServerApplicationDatagram(frames)
|
||||
|
||||
/**
|
||||
* Decrypt the application-level (1-RTT) packet inside a client-emitted
|
||||
* datagram and return its frames in wire order. Test-only helper for
|
||||
* assertions that depend on per-frame ordering inside a packet (e.g.
|
||||
* stream priority scheduling). Walks past any coalesced long-header
|
||||
* packets (Initial / Handshake) at the front of the datagram, since the
|
||||
* client may still flush ACKs at those levels post-handshake. Returns
|
||||
* null if no short-header packet is present or decryption fails.
|
||||
*/
|
||||
fun decryptClientApplicationFrames(datagram: ByteArray): List<com.vitorpamplona.quic.frame.Frame>? {
|
||||
if (datagram.isEmpty()) return null
|
||||
var offset = 0
|
||||
while (offset < datagram.size) {
|
||||
val first = datagram[offset].toInt() and 0xFF
|
||||
if ((first and 0x80) == 0) {
|
||||
val proto = serverApplicationRx ?: return null
|
||||
val parsed =
|
||||
ShortHeaderPacket.parseAndDecrypt(
|
||||
bytes = datagram,
|
||||
offset = offset,
|
||||
dcidLen = serverScid.length,
|
||||
aead = proto.aead,
|
||||
key = proto.key,
|
||||
iv = proto.iv,
|
||||
hp = proto.hp,
|
||||
hpKey = proto.hpKey,
|
||||
largestReceivedInSpace = applicationPnSpace.largestReceived,
|
||||
) ?: return null
|
||||
applicationPnSpace.observeInbound(parsed.packet.packetNumber, 0L)
|
||||
return decodeFrames(parsed.packet.payload)
|
||||
}
|
||||
// Long header — skip past it using the encoded length field so
|
||||
// we can inspect any short-header packet that was coalesced
|
||||
// after it.
|
||||
val peeked = LongHeaderPacket.peekHeader(datagram, offset) ?: return null
|
||||
offset += peeked.totalLength
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
private fun buildServerInitialPacket(crypto: ByteArray): ByteArray {
|
||||
val proto =
|
||||
PacketProtection(
|
||||
|
||||
+80
@@ -110,6 +110,86 @@ class QuicConnectionWriterTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun writer_drains_higher_priority_streams_before_lower_priority() {
|
||||
// T11.3 follow-up: priority must dominate iteration order on
|
||||
// EVERY drain, not just the first. The naive "sort by priority
|
||||
// then apply the existing rotating start globally" shape looks
|
||||
// right at a glance but the rotating start advances on every
|
||||
// drain, so cross-tier ordering flips on alternate drains and
|
||||
// the priority hint is silently defeated. The correct shape is
|
||||
// strict priority across tiers + round-robin only within a
|
||||
// tier, which we pin here by draining TWICE and asserting the
|
||||
// higher-priority stream's StreamFrame lands first in BOTH
|
||||
// packets. Low-priority stream is opened FIRST so insertion-
|
||||
// order can't accidentally pass for priority ordering.
|
||||
runBlocking {
|
||||
val (client, pipe) = connectedClient()
|
||||
val low = client.openBidiStream()
|
||||
val high = client.openBidiStream()
|
||||
low.priority = 0
|
||||
high.priority = 10
|
||||
|
||||
repeat(2) { round ->
|
||||
low.send.enqueue(ByteArray(200) { 0xAA.toByte() })
|
||||
high.send.enqueue(ByteArray(200) { 0xBB.toByte() })
|
||||
|
||||
val datagram = drainOutbound(client, nowMillis = 0L)
|
||||
assertNotNull(datagram, "drain on round $round must emit a packet")
|
||||
val frames = pipe.decryptClientApplicationFrames(datagram)
|
||||
assertNotNull(frames, "decrypt must succeed on round $round")
|
||||
val streamFrames = frames.filterIsInstance<StreamFrame>()
|
||||
assertEquals(
|
||||
2,
|
||||
streamFrames.size,
|
||||
"expected one StreamFrame per stream on round $round, got $streamFrames",
|
||||
)
|
||||
assertEquals(
|
||||
high.streamId,
|
||||
streamFrames[0].streamId,
|
||||
"higher-priority stream must drain first on round $round; " +
|
||||
"saw ${streamFrames.map { it.streamId }}",
|
||||
)
|
||||
assertEquals(low.streamId, streamFrames[1].streamId, "low second on round $round")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun writer_round_robins_within_a_priority_tier() {
|
||||
// Regression guard for the tier-local round-robin: same-priority
|
||||
// streams must still rotate so an early-opened stream doesn't
|
||||
// monopolise a packet's stream-frame slot indefinitely. We open
|
||||
// three streams at the default (0) priority and verify the
|
||||
// rotating start advances by one per drain, matching the
|
||||
// pre-priority behaviour.
|
||||
runBlocking {
|
||||
val (client, pipe) = connectedClient()
|
||||
val a = client.openBidiStream()
|
||||
val b = client.openBidiStream()
|
||||
val c = client.openBidiStream()
|
||||
// All default priority — single tier, three streams.
|
||||
val expectedRotation =
|
||||
listOf(
|
||||
listOf(a.streamId, b.streamId, c.streamId),
|
||||
listOf(b.streamId, c.streamId, a.streamId),
|
||||
listOf(c.streamId, a.streamId, b.streamId),
|
||||
)
|
||||
for ((round, expected) in expectedRotation.withIndex()) {
|
||||
a.send.enqueue(ByteArray(64) { 0xA1.toByte() })
|
||||
b.send.enqueue(ByteArray(64) { 0xB2.toByte() })
|
||||
c.send.enqueue(ByteArray(64) { 0xC3.toByte() })
|
||||
|
||||
val datagram = drainOutbound(client, nowMillis = 0L)
|
||||
assertNotNull(datagram, "drain $round must emit a packet")
|
||||
val frames = pipe.decryptClientApplicationFrames(datagram)
|
||||
assertNotNull(frames, "decrypt must succeed on round $round")
|
||||
val ids = frames.filterIsInstance<StreamFrame>().map { it.streamId }
|
||||
assertEquals(expected, ids, "round $round round-robin order")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun writer_respects_connection_level_send_credit_cap() {
|
||||
// Audit-4 #9: pre-fix the writer ignored sendConnectionFlowCredit
|
||||
|
||||
Reference in New Issue
Block a user