test(nests): fix nextSubscribeBidi helper to terminate after match

The takeWhile/collect pattern only re-checks its predicate when the
NEXT upstream value emits, so once nextSubscribeBidi found its
target Subscribe bidi, the helper sat blocked waiting for a
follow-up bidi that may never come — turning groups_are_demuxed_by_
subscribeId into a 5-minute hang on tests that open exactly two
subscribes (the announce-watch bidi happened to land between, but
relying on it to nudge the flow forward is fragile).

Switched to transformWhile { … emit(…); false }.firstOrNull() so
the upstream collection terminates synchronously after the first
match. Warm-daemon test time drops from multi-minute timeout back
to the expected ~10 ms range.

https://claude.ai/code/session_01HXf3zG3F2ev2ASeQju7Y5S
This commit is contained in:
Claude
2026-04-28 16:22:01 +00:00
parent cd9279d23b
commit 5f71dc7c76
@@ -30,8 +30,8 @@ import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.flow.first import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.firstOrNull import kotlinx.coroutines.flow.firstOrNull
import kotlinx.coroutines.flow.take import kotlinx.coroutines.flow.take
import kotlinx.coroutines.flow.takeWhile
import kotlinx.coroutines.flow.toList import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.flow.transformWhile
import kotlinx.coroutines.runBlocking import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout import kotlinx.coroutines.withTimeout
import kotlin.test.AfterTest import kotlin.test.AfterTest
@@ -445,13 +445,11 @@ class MoqLiteSessionTest {
* are simply abandoned — their pump-side `bidi.incoming()` * are simply abandoned — their pump-side `bidi.incoming()`
* collect just sits idle until the test ends. * collect just sits idle until the test ends.
*/ */
private suspend fun nextSubscribeBidi(serverSide: FakeWebTransport): Pair<FakeBidiStream, MoqLiteSubscribe> { private suspend fun nextSubscribeBidi(serverSide: FakeWebTransport): Pair<FakeBidiStream, MoqLiteSubscribe> =
val bidiFlow = serverSide.peerOpenedBidiStreams() serverSide
var found: Pair<FakeBidiStream, MoqLiteSubscribe>? = null .peerOpenedBidiStreams()
bidiFlow .transformWhile { bidi ->
.takeWhile { found == null } val firstChunk = bidi.incoming().firstOrNull() ?: return@transformWhile true
.collect { bidi ->
val firstChunk = bidi.incoming().firstOrNull() ?: return@collect
val code = MoqLiteFrameBuffer().apply { push(firstChunk) }.readVarint() val code = MoqLiteFrameBuffer().apply { push(firstChunk) }.readVarint()
if (code != MoqLiteControlType.Subscribe.code) { if (code != MoqLiteControlType.Subscribe.code) {
// Housekeeping bidi (announce watch, etc.) — // Housekeeping bidi (announce watch, etc.) —
@@ -459,7 +457,7 @@ class MoqLiteSessionTest {
// collector will idle indefinitely, which is // collector will idle indefinitely, which is
// fine for unit tests under runBlocking + // fine for unit tests under runBlocking +
// pumpScope cleanup. // pumpScope cleanup.
return@collect return@transformWhile true
} }
val bodyChunk = val bodyChunk =
bidi.incoming().firstOrNull() bidi.incoming().firstOrNull()
@@ -467,10 +465,15 @@ class MoqLiteSessionTest {
val payload = val payload =
MoqLiteFrameBuffer().apply { push(bodyChunk) }.readSizePrefixed() MoqLiteFrameBuffer().apply { push(bodyChunk) }.readSizePrefixed()
?: error("subscribe body chunk did not contain a complete size-prefixed payload") ?: error("subscribe body chunk did not contain a complete size-prefixed payload")
found = bidi to MoqLiteCodec.decodeSubscribe(payload) emit(bidi to MoqLiteCodec.decodeSubscribe(payload))
} // Terminate upstream collection — without this the
return found ?: error("flow ended without a Subscribe bidi") // helper would block waiting for a NEXT bidi that
} // may never come, since `takeWhile` only re-checks
// its predicate when the next value emits. Tests
// that open exactly two subscribes (no third bidi
// to nudge the flow forward) used to hang here.
false
}.firstOrNull() ?: error("flow ended without a Subscribe bidi")
private fun framePayload(bytes: ByteArray): ByteArray = Varint.encode(bytes.size.toLong()) + bytes private fun framePayload(bytes: ByteArray): ByteArray = Varint.encode(bytes.size.toLong()) + bytes