From f17e7adfa7acdb2bae9f666254f2aef2bf164756 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 5 May 2026 19:39:23 +0000 Subject: [PATCH] fix(nests): announces flow MUST use channelFlow, not flow{} MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The diagnostic logging in 1fc8dbc surfaced the real root cause — not the SharedFlow replay race I fixed in d6517cf, but a strict flow{}-builder concurrency-confinement violation. The receiver-side log (15:34:40, repro after the diagnostic patch) showed: session.announce(prefix='') bidi pump emit #1 status=Active suffix='fe52579aa30e' (chunks=1) MoqLiteNestsListener.announces bidi#2 #1 status=Active suffix='fe52579aa30e' MoqLiteNestsListener.announces bidi#2 finally (emissions=1) wrapper.announces() iter=2 inner collect ended IllegalStateException: Flow invariant is violated: Emission from another coroutine is detected. Child of StandaloneCoroutine{Active}@e9b5981, expected child of StandaloneCoroutine{Active}@8309a26. FlowCollector is not thread-safe and concurrent emissions are prohibited. To mitigate this restriction please use 'channelFlow' builder instead of 'flow' The session's bidi pump (`scope.launch { bidi.incoming().collect … updates.emit(decoded) }`) emits to the announce SharedFlow from the PUMP's coroutine. SharedFlow's `collect { … }` resumes the lambda inline on the emitter's coroutine when there's no dispatcher hop, so `MoqLiteNestsListener.announces`' `handle.updates.collect { emit(RoomAnnouncement(…)) }` lambda fires on the pump's coroutine, not the flow{} body's coroutine. `flow {}`'s SafeFlow guard throws on the cross-coroutine emit. The wrapper's `runCatching { listener.announces().collect { emit(it) } }` swallows the exception, the wrapper's collect ends with `fwd=1`, and `_announcedSpeakers` stays empty for the lifetime of the session. That, in turn, kept the cliff detector permanently gated: cliff-detector tick=N active=1 announced=0 lastFrameAges=[fe52579a=…ms] even though the receiver was clearly subscribed and clearly seeing frames. The cliff detector requires the speaker to be in `announcedSpeakers` to recycle, so the relay-forward cliff at ~135 streams went undetected — the original two-phone receiver-side silence symptom. Fix: switch `MoqLiteNestsListener.announces()` from `flow { … }` to `channelFlow { … }`. `channelFlow` is the documented escape hatch for cross-coroutine emission (the error message itself recommends it). The new shape: channelFlow { val handle = session.announce(prefix = "") val pump = launch { try { handle.updates.collect { send(RoomAnnouncement(…)) } } finally { withContext(NonCancellable) { runCatching { handle.close() } } } } awaitClose { pump.cancel() } } Notes: - `send` (channel) replaces `emit` (flow) — channelFlow buffers via a Channel, so cross-coroutine producers are explicitly supported. - The inner `collect` is wrapped in a launched child of the channelFlow scope so we can use `awaitClose` for the producer- side teardown contract channelFlow requires. - `handle.close` is suspend (FINs the bidi + joins the pump); putting it in the launched pump's `finally` under `NonCancellable` keeps the FIN reliable when the consumer cancels mid-stream. Diagnostic logging from 1fc8dbc preserved — the next two-phone repro should now show: observeAnnounces #1 active=true pubkey='fe52579aa30e' → ADD cliff-detector tick=N active=1 announced=1 lastFrameAges=[…] …and on a real cliff event, the recycle should fire. Tests: full suite still passes - CliffDetectorTest: 12/12 - NestPlayerTest: 12/12 - NestBroadcasterTest: 6/6 - MoqLiteSessionTest: 11/11 - All other nestsClient jvmTest classes (~240 tests across 37 classes) green. --- .../nestsclient/MoqLiteNestsListener.kt | 98 ++++++++++++------- 1 file changed, 65 insertions(+), 33 deletions(-) diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsListener.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsListener.kt index 3ee7fcf97..6b6c308bb 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsListener.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsListener.kt @@ -26,12 +26,16 @@ import com.vitorpamplona.nestsclient.moq.SubscribeOk import com.vitorpamplona.nestsclient.moq.lite.MoqLiteAnnounceStatus import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSubscribeException +import kotlinx.coroutines.NonCancellable +import kotlinx.coroutines.channels.awaitClose import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow -import kotlinx.coroutines.flow.flow +import kotlinx.coroutines.flow.channelFlow import kotlinx.coroutines.flow.map +import kotlinx.coroutines.launch +import kotlinx.coroutines.withContext import java.util.concurrent.atomic.AtomicLong /** @@ -127,9 +131,26 @@ class MoqLiteNestsListener internal constructor( } override fun announces(): Flow = - flow { - com.vitorpamplona.quartz.utils.Log - .d("NestRx") { "MoqLiteNestsListener.announces flow STARTING (state=${state.value::class.simpleName})" } + // `channelFlow` (NOT `flow`) is mandatory here. The downstream + // emission source is `handle.updates`, a MutableSharedFlow + // populated by a launched bidi pump on the session's scope. + // SharedFlow's `collect { lambda }` resumes the lambda inline + // when emit happens — so when the pump emits a chunk, the + // collect lambda runs on the pump's coroutine, not the + // collector's. `flow {}` enforces the "emit only from the + // builder coroutine" invariant and throws + // `IllegalStateException: Flow invariant is violated` the + // moment we call `emit()` from a different coroutine. Two- + // phone production logs (commit 1fc8dbc, run 15:34:40) + // showed exactly this — the first Active arrived, the + // collect lambda fired, emit() threw, the wrapper's + // `runCatching { listener.announces().collect { emit(it) } }` + // swallowed it, `_announcedSpeakers` stayed empty forever, + // and the cliff detector reported `active=1 announced=0` + // indefinitely. `channelFlow` is the documented fix + // (the error message itself recommends it): emissions go + // through a buffered channel that's safe across coroutines. + channelFlow { check(state.value is NestsListenerState.Connected) { "NestsListener.announces requires Connected state, was ${state.value}" } @@ -138,37 +159,48 @@ class MoqLiteNestsListener internal constructor( // suffix is the broadcast's path component within the // room — for nests this is the speaker pubkey hex. val handle = session.announce(prefix = "") - com.vitorpamplona.quartz.utils.Log - .d("NestRx") { "MoqLiteNestsListener.announces opened bidi#2 (VM-level)" } - var emissionCount = 0 - try { - handle.updates.collect { announce -> - emissionCount += 1 - com.vitorpamplona.quartz.utils.Log.d("NestRx") { - "MoqLiteNestsListener.announces bidi#2 #$emissionCount status=${announce.status} suffix='${announce.suffix.take(12)}'" + // Run the inner collect on a child of channelFlow's scope + // so we can `awaitClose` on the producer side and let + // close handle the unsub side-effect. Without the launch, + // we'd block here in `collect` forever and never reach + // `awaitClose` — but channelFlow's contract requires + // awaitClose for clean shutdown. + val pump = + launch { + try { + handle.updates.collect { announce -> + send( + RoomAnnouncement( + pubkey = announce.suffix, + active = announce.status == MoqLiteAnnounceStatus.Active, + ), + ) + } + } catch (ce: kotlinx.coroutines.CancellationException) { + throw ce + } catch (_: Throwable) { + // updates flow died — close the channel so the + // consumer sees end-of-flow and the wrapper + // can decide what to do (typically wait for + // the next activeListener emission). + } finally { + // Close the announce bidi on the same scope + // that owned the pump. `handle.close` is + // suspend (it FINs the bidi and joins the + // session-side pump); run it under + // NonCancellable so cancellation of this + // coroutine doesn't skip the FIN. + withContext(NonCancellable) { + runCatching { handle.close() } + } } - emit( - RoomAnnouncement( - pubkey = announce.suffix, - active = announce.status == MoqLiteAnnounceStatus.Active, - ), - ) - } - com.vitorpamplona.quartz.utils.Log - .w("NestRx") { "MoqLiteNestsListener.announces bidi#2 collect ENDED naturally after $emissionCount emissions" } - } finally { - com.vitorpamplona.quartz.utils.Log - .w("NestRx") { "MoqLiteNestsListener.announces bidi#2 finally (emissions=$emissionCount)" } - // The flow's `finally` runs on both normal completion - // and cancellation — let cancel propagate after closing - // the announce handle. - try { - handle.close() - } catch (ce: kotlinx.coroutines.CancellationException) { - throw ce - } catch (_: Throwable) { - // Best-effort. } + awaitClose { + // Caller cancelled, OR the wrapping `collectLatest` + // restarted (listener swap). Cancel the pump; its + // finally above runs the suspend close on a + // NonCancellable context. + pump.cancel() } }