From bd681a0a5fc3667a3a0955915407fd3d9c3815a0 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 5 May 2026 18:06:56 +0000 Subject: [PATCH] test(nests): close FRAME1 delivery race in subscribeSpeaker_survives_session_swap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The test emitted FRAME1 into first.frames and immediately called first.fail(...), trusting that FRAME1 would propagate through the pump to the wrapper's frames before the session swap. emit() is non-suspending — it only enqueues FRAME1 into the pump's slot. On slower hosts (observed on macOS CI) the orchestrator's reconnect path can flip activeListener to the second session before the pump's collect lambda runs, at which point collectLatest cancels pump-iteration-1 mid-resume and FRAME1 is dropped. The consumer's take(2) then only ever sees FRAME2 and the async's withTimeout fires after 5 s. Add a consumerProgress StateFlow that the consumer's collector bumps on each frame, and wait for it to reach 1 before failing the listener. Same shape as the existing consumerSubscribed gate that closed the "emit before consumer subscribes" race. --- .../ReconnectingNestsListenerTest.kt | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListenerTest.kt b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListenerTest.kt index dfccc847fa..3d4249c161 100644 --- a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListenerTest.kt +++ b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListenerTest.kt @@ -39,6 +39,7 @@ import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.first +import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.flow.onSubscription import kotlinx.coroutines.flow.take import kotlinx.coroutines.flow.toList @@ -189,6 +190,18 @@ class ReconnectingNestsListenerTest { // before the test thread races into the first emit. val consumerSubscribed = CompletableDeferred() + // Counter for frames the consumer has actually observed. + // Needed to break the FRAME1-delivery race: emit() into + // first.frames is non-suspending and just enqueues into + // the pump's slot. If we trigger a session swap before + // the pump's collect lambda runs (`frames.emit(it)` to + // the wrapper), collectLatest cancels iteration 1 mid- + // resume and FRAME1 is lost — consumer ends with 1/2 + // frames and the async's withTimeout fires. The + // collector is single-coroutine so a plain StateFlow + // update is safe; the test thread reads it via .first. + val consumerProgress = MutableStateFlow(0) + // The wrapper backs handle.objects with a SharedFlow // (frames.asSharedFlow); the cast lets us use // SharedFlow.onSubscription, which fires AFTER the @@ -201,6 +214,7 @@ class ReconnectingNestsListenerTest { objectsAsShared .onSubscription { consumerSubscribed.complete(Unit) } .take(2) + .onEach { consumerProgress.value += 1 } .toList() } } @@ -220,6 +234,19 @@ class ReconnectingNestsListenerTest { first.frames.emit(frame(byteArrayOf(0x01))) + // Wait for FRAME1 to traverse pump → wrapper.frames → + // consumer collector. Without this sync the next + // `first.fail(...)` can race ahead and trigger + // collectLatest cancellation of pump-iteration-1 while + // FRAME1 is still queued in first.frames; the cancel + // interrupts the pump's resume before its lambda runs + // and FRAME1 never reaches the wrapper. Consumer then + // observes only FRAME2 (1 frame ≠ take(2)) and the + // async's withTimeout fires after 5 s. + withTimeout(5_000L) { + consumerProgress.first { it >= 1 } + } + // Force a reconnect: fail the first listener, the // orchestrator opens the next one. first.fail("scripted-disconnect")