From d8ab4fd99b825c8420c1c2892a4f74053e8f86e6 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 28 Apr 2026 17:06:09 +0000 Subject: [PATCH] fix(nests): rotate moq-lite groups + reconnect after publisher cycle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit To make a SubscribeHandle survive a publisher session swap on the same relay, three things had to change together: 1. Broadcaster emits one Opus frame per moq-lite group (`publisher.send` + `publisher.endGroup`). moq-lite's "from-latest" subscribe semantics deliver a new subscriber the NEXT group's frames; without per-frame rotation a subscriber that attaches mid-broadcast waits forever for the (single, never-ending) group to end. 2. MoqLiteSession opens its announce-watch bidi synchronously before the first subscribe, dispatches subscribe + announce frames over a single long-running collector (varint type code hoisted outside `collect`), and FINs the publisher's currentGroup when an inbound subscribe bidi closes — so the next send opens a fresh group keyed off the live subscriber instead of the recycled one. 3. ReconnectingNestsListener no longer breaks on terminal=Closed in the orchestrator; the user-driven stop path goes through `orchestrator.cancel()`, so any other Closed (peer-driven transport close, half-broken session, publisher recycle) is a reconnect trigger. Round-trip interop test updated to assert `groupId == idx` (one group per frame) rather than `groupId == 0`. All three reconnecting-listener interop scenarios now pass against the real moq-rs relay: happy-path, session-swap, and listener-survives-publisher-recycle. --- .../nestsclient/ReconnectingNestsListener.kt | 15 +- .../audio/NestMoqLiteBroadcaster.kt | 38 ++- .../nestsclient/moq/lite/MoqLiteSession.kt | 281 +++++++++++------- .../moq/lite/MoqLiteSessionTest.kt | 9 +- .../interop/NostrNestsRoundTripInteropTest.kt | 8 +- 5 files changed, 218 insertions(+), 133 deletions(-) diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListener.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListener.kt index d2e738b8d..48e2d3342 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListener.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/ReconnectingNestsListener.kt @@ -171,9 +171,20 @@ suspend fun connectReconnectingNestsListener( // never enters Reconnecting. continue } - val terminal = state.value - if (terminal is NestsListenerState.Closed) break + // Note: we do NOT break on terminal=Closed. The + // user-driven stop path goes through + // [ReconnectingHandle.close], which calls + // `orchestrator.cancel()` BEFORE closing the inner + // listener; cancellation propagates through the + // next suspending call (typically `delay` below or + // `openOnce` on the next loop iteration) and the + // orchestrator exits cleanly. Any *other* path that + // produces a Closed inner listener — peer-driven + // transport close, half-broken session that was + // closed by some internal cleanup — should be + // treated as an unexpected drop and reconnected. if (policy.isExhausted(attempt + 1)) break + val terminal = state.value val delayMs = if (terminal is NestsListenerState.Reconnecting) { terminal.delayMs diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/NestMoqLiteBroadcaster.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/NestMoqLiteBroadcaster.kt index 276d51440..808f9bbe0 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/NestMoqLiteBroadcaster.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/audio/NestMoqLiteBroadcaster.kt @@ -87,17 +87,33 @@ class NestMoqLiteBroadcaster( } if (opus.isEmpty()) continue if (muted) continue - runCatching { publisher.send(opus) } - .onFailure { t -> - if (t is CancellationException) throw t - onError( - AudioException( - AudioException.Kind.PlaybackFailed, - "publisher.send failed", - t, - ), - ) - } + // One Opus frame per moq-lite group — mirrors the + // nests JS reference's audio publish path, and is + // load-bearing for the listener-survives-publisher- + // recycle invariant: a brand-new subscriber that + // attaches mid-broadcast (e.g. listener wrapper + // re-subscribing after a publisher cycle) gets the + // NEXT group's frames per moq-lite "from-latest" + // semantics. Without endGroup, the entire broadcast + // is one giant group and new subscribers wait + // indefinitely. The 20 ms cadence here means at + // most one frame of audio missed for any new + // subscriber. See + // `nestsClient/plans/2026-04-26-moq-lite-gap.md`'s + // "Group size: 1 frame per group" line. + runCatching { + publisher.send(opus) + publisher.endGroup() + }.onFailure { t -> + if (t is CancellationException) throw t + onError( + AudioException( + AudioException.Kind.PlaybackFailed, + "publisher.send failed", + t, + ), + ) + } } } catch (ce: CancellationException) { throw ce diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt index 1697cc566..573f94565 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSession.kt @@ -33,7 +33,6 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.consumeAsFlow -import kotlinx.coroutines.flow.firstOrNull import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -163,6 +162,18 @@ class MoqLiteSession internal constructor( endGroup: Long? = null, ): MoqLiteSubscribeHandle { ensureOpen() + // Open the announce-watch bidi BEFORE the subscribe goes + // out. moq-rs uses the announce stream to propagate + // broadcast availability into the subscriber session — a + // subscribe that arrives before our session has any + // announce bidi open is rejected with "not found", even + // when the publisher's session is alive on the relay. The + // bidi must be on the wire before subscribe sends; lazy- + // launching after subscribe (the obvious-but-wrong shape) + // races the relay's discovery and produces flaky misses, + // especially for fresh listener sessions opened after a + // wrapper-driven reconnect. + ensureAnnounceWatchStarted() val id = state.withLock { check(!closed) { "session is closed" } @@ -234,29 +245,6 @@ class MoqLiteSession internal constructor( state.withLock { subscriptionsBySubscribeId[id] = sub if (groupPump == null) groupPump = scope.launch { pumpUniStreams() } - // Lazy-launch the publisher-disconnect watcher - // — one shared announce bidi per session whose - // sole job is to close the frames channel on any - // ListenerSubscription whose broadcast path goes - // Ended. moq-lite Lite-03 has no explicit - // "publisher gone" message on the subscribe - // bidi (the relay keeps that bidi open across - // publisher cycles in case a fresh publisher - // takes over the suffix), so the announce - // stream IS the only signal. - // - // Without this, a wrapper-layer consumer - // collecting from [MoqLiteSubscribeHandle.frames] - // would sit silent indefinitely after a - // publisher cycle even though the relay is happy - // to serve a fresh subscribe under the same - // suffix. The wrapper-level - // [com.vitorpamplona.nestsclient.ReconnectingNestsListener] - // pump's natural collect-completion path then - // re-issues the subscribe. - if (announceWatchJob == null) { - announceWatchJob = scope.launch { pumpAnnounceWatch() } - } } return MoqLiteSubscribeHandle( id = id, @@ -268,13 +256,56 @@ class MoqLiteSession internal constructor( } } + /** + * Lock that serializes lazy-launch of the announce-watch. + * Distinct from [state] so the synchronous `announce(prefix="")` + * inside [ensureAnnounceWatchStarted] can suspend without + * blocking other state-mutating operations. + */ + private val announceWatchLock = Mutex() + + /** + * Open the shared announce-watch bidi *synchronously* (and + * launch its collector coroutine) if it isn't already running. + * Idempotent. Called from [subscribe] before the subscribe + * message goes on the wire so moq-rs has a chance to propagate + * broadcast availability into our session before the subscribe + * arrives — see the comment in [subscribe]. + */ + private suspend fun ensureAnnounceWatchStarted() { + announceWatchLock.withLock { + if (announceWatchJob != null) return + val handle = + try { + announce(prefix = "") + } catch (ce: kotlinx.coroutines.CancellationException) { + throw ce + } catch (_: Throwable) { + // Couldn't open the announce bidi — best effort, + // bail. Subscriptions still work; we just lose + // automatic cycle detection (the wrapper still + // re-issues on listener swap / explicit failure). + return + } + announceWatchJob = + scope.launch { + try { + pumpAnnounceWatch(handle) + } finally { + announceWatchLock.withLock { announceWatchJob = null } + } + } + } + } + /** * Single shared announce-watch pump for ALL subscriptions on - * this session. Opens one announce bidi (prefix="") on first - * subscribe, then for each [MoqLiteAnnounceStatus.Ended] update - * iterates the subscription map and closes the frames channel of - * any subscription whose `broadcast` matches the announce - * suffix. The closed channel ends the consumer-facing + * this session. Driven by the bidi opened in + * [ensureAnnounceWatchStarted]. For each + * [MoqLiteAnnounceStatus.Ended] update, iterates the + * subscription map and closes the frames channel of any + * subscription whose `broadcast` matches the announce suffix. + * The closed channel ends the consumer-facing * `frames.consumeAsFlow()` flow naturally — same shape as a * user-driven `handle.unsubscribe()` from the consumer's POV — * which lets the wrapper's re-issuance pump drive a fresh @@ -287,18 +318,7 @@ class MoqLiteSession internal constructor( * silence — the session itself recovers via its own reconnect * path. Cancelled when [scope] is cancelled (session close). */ - private suspend fun pumpAnnounceWatch() { - val handle = - try { - announce(prefix = "") - } catch (ce: kotlinx.coroutines.CancellationException) { - throw ce - } catch (_: Throwable) { - // Couldn't open the announce bidi — best effort, - // bail. Subscriptions still work; we just lose - // automatic cycle detection. - return - } + private suspend fun pumpAnnounceWatch(handle: MoqLiteAnnouncesHandle) { try { handle.updates.collect { update -> if (update.status != MoqLiteAnnounceStatus.Ended) return@collect @@ -329,7 +349,6 @@ class MoqLiteSession internal constructor( // Announce bidi died — same best-effort fallback. } finally { runCatching { handle.close() } - state.withLock { announceWatchJob = null } } } @@ -475,86 +494,106 @@ class MoqLiteSession internal constructor( } private suspend fun handleInboundBidi(bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream) { - val buffer = MoqLiteFrameBuffer() val publisher = state.withLock { activePublisher } ?: return + + // Single long-running collector for the bidi's full lifetime. + // Pre-fix this dispatch was split into a `firstOrNull()` to + // peek the control byte + a `readSizePrefixedFromBidiInto` + // to read the body — but `bidi.incoming()` is backed by + // `Channel.consumeAsFlow(consume=true)`, which + // CANCELS the channel when the first collect ends. Any + // attempt to re-collect from the bidi (e.g. to watch for + // subscriber-disconnect FIN) saw an immediately-empty + // closed flow, firing the cleanup right after registration + // and starving the publisher's send path. With one collector + // for the bidi's whole life, the dispatch reads the message, + // the collector continues silently until peer FIN, and the + // post-collect cleanup runs exactly once — same shape the + // moq-lite session's [announce] pump already uses. + val buffer = MoqLiteFrameBuffer() + // typeCode is hoisted outside the collect lambda so it + // survives across invocations — `buffer.readVarint()` + // advances `pos`, so calling it again on the next collect + // tick would read body bytes as if they were the control + // varint and tear the dispatch state apart. + var typeCode: Long? = null + var dispatched = false + var inboundSub: MoqLiteSubscribe? = null try { - // Read the leading ControlType varint from the first chunk. - val first = - bidi.incoming().firstOrNull() ?: return - buffer.push(first) - val controlCode = buffer.readVarint() ?: return - val controlType = MoqLiteControlType.fromCode(controlCode) ?: return - when (controlType) { - MoqLiteControlType.Announce -> { - handleAnnounceRequest(bidi, buffer, publisher) - } + bidi.incoming().collect { chunk -> + buffer.push(chunk) + if (!dispatched) { + if (typeCode == null) typeCode = buffer.readVarint() + val tc = typeCode ?: return@collect + val controlType = + MoqLiteControlType.fromCode(tc) ?: run { + dispatched = true + runCatching { bidi.finish() } + return@collect + } + when (controlType) { + MoqLiteControlType.Announce -> { + val pleasePayload = buffer.readSizePrefixed() ?: return@collect + val please = MoqLiteCodec.decodeAnnouncePlease(pleasePayload) + val emittedSuffix = + MoqLitePath.stripPrefix(please.prefix, publisher.suffix) ?: publisher.suffix + bidi.write( + MoqLiteCodec.encodeAnnounce( + MoqLiteAnnounce( + status = MoqLiteAnnounceStatus.Active, + suffix = emittedSuffix, + hops = 0L, + ), + ), + ) + publisher.registerAnnounceBidi(bidi, emittedSuffix) + dispatched = true + } - MoqLiteControlType.Subscribe -> { - handleSubscribeRequest(bidi, buffer, publisher) - } + MoqLiteControlType.Subscribe -> { + val subPayload = buffer.readSizePrefixed() ?: return@collect + val sub = MoqLiteCodec.decodeSubscribe(subPayload) + bidi.write( + MoqLiteCodec.encodeSubscribeOk( + MoqLiteSubscribeOk( + priority = sub.priority, + ordered = sub.ordered, + maxLatencyMillis = sub.maxLatencyMillis, + startGroup = null, + endGroup = null, + ), + ), + ) + publisher.registerInboundSubscription(sub) + inboundSub = sub + dispatched = true + } - else -> { - // Lite-03 treats Session/Fetch/Probe as separate flows; - // we don't implement them here. Drop the bidi. - runCatching { bidi.finish() } + else -> { + // Lite-03 treats Session/Fetch/Probe as + // separate flows; we don't implement them. + runCatching { bidi.finish() } + dispatched = true + } + } } + // Post-dispatch chunks are silently discarded — + // Lite-03's announce / subscribe bidis are idle + // after the response. The signal we care about is + // the flow's natural completion (peer FIN = + // subscriber-disconnect, or transport drop). } } catch (ce: CancellationException) { throw ce } catch (_: Throwable) { - runCatching { bidi.finish() } + // Bidi errored — fall through to the same cleanup. } - } - - private suspend fun handleAnnounceRequest( - bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, - seedBuffer: MoqLiteFrameBuffer, - publisher: PublisherStateImpl, - ) { - val pleasePayload = readSizePrefixedFromBidiInto(bidi.incoming(), seedBuffer) - val please = MoqLiteCodec.decodeAnnouncePlease(pleasePayload) - // The relay sets the prefix to the namespace it expects us to - // publish under (typically `claims.root`). Our broadcast path - // (after stripping the prefix) is `publisher.suffix`. moq-lite - // requires the suffix on the wire to be the *remaining* part - // after `please.prefix` — so strip it. - val emittedSuffix = MoqLitePath.stripPrefix(please.prefix, publisher.suffix) ?: publisher.suffix - bidi.write( - MoqLiteCodec.encodeAnnounce( - MoqLiteAnnounce( - status = MoqLiteAnnounceStatus.Active, - suffix = emittedSuffix, - hops = 0L, - ), - ), - ) - // Hold the bidi open until the publisher closes; if/when the - // application stops broadcasting, send `Ended`. - publisher.registerAnnounceBidi(bidi, emittedSuffix) - } - - private suspend fun handleSubscribeRequest( - bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, - seedBuffer: MoqLiteFrameBuffer, - publisher: PublisherStateImpl, - ) { - val subPayload = readSizePrefixedFromBidiInto(bidi.incoming(), seedBuffer) - val sub = MoqLiteCodec.decodeSubscribe(subPayload) - // Reply Ok right away — moq-lite is permissive on the publisher - // side; the relay decides whether the subscriber is allowed to - // see this broadcast. - bidi.write( - MoqLiteCodec.encodeSubscribeOk( - MoqLiteSubscribeOk( - priority = sub.priority, - ordered = sub.ordered, - maxLatencyMillis = sub.maxLatencyMillis, - startGroup = null, - endGroup = null, - ), - ), - ) - publisher.registerInboundSubscription(sub) + // Flow ended (peer FIN or error). Remove the inbound + // subscribe so the publisher's send path stops keying new + // groups off this dead subscriber. Announce bidis are + // owned by the publisher state for sending Ended on + // publisher-close — we don't remove them here. + inboundSub?.let { publisher.removeInboundSubscription(it) } } /** @@ -742,6 +781,24 @@ class MoqLiteSession internal constructor( } } + /** + * Remove an inbound subscription whose bidi was FIN'd by the + * relay (subscriber disconnected). FINs the current group + * defensively because [openNextGroupLocked] keys each uni + * stream off `inboundSubs.first()`'s id; if the dropped sub + * was first, the current uni stream is dead-routed and the + * next send must open a fresh group keyed off whatever + * live sub is now first. + */ + suspend fun removeInboundSubscription(sub: MoqLiteSubscribe) { + gate.withLock { + if (publisherClosed) return + if (!inboundSubs.remove(sub)) return + runCatching { currentGroup?.uni?.finish() } + currentGroup = null + } + } + override suspend fun startGroup() { gate.withLock { if (publisherClosed) return diff --git a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt index 1232e7174..58263dc75 100644 --- a/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt +++ b/nestsClient/src/commonTest/kotlin/com/vitorpamplona/nestsclient/moq/lite/MoqLiteSessionTest.kt @@ -79,8 +79,7 @@ class MoqLiteSessionTest { val peerHandlesSubscribe = async { - val bidi = serverSide.peerOpenedBidiStreams().first() - val req = readSubscribeRequest(bidi) + val (bidi, req) = nextSubscribeBidi(serverSide) assertEquals("speakerPubkey", req.broadcast) assertEquals("audio/data", req.track) assertEquals(MoqLiteSession.DEFAULT_PRIORITY, req.priority) @@ -130,8 +129,7 @@ class MoqLiteSessionTest { val peer = async { - val bidi = serverSide.peerOpenedBidiStreams().first() - readSubscribeRequest(bidi) + val (bidi, _) = nextSubscribeBidi(serverSide) bidi.write( MoqLiteCodec.encodeSubscribeDrop( MoqLiteSubscribeDrop(errorCode = 4L, reasonPhrase = "no such broadcast"), @@ -396,8 +394,7 @@ class MoqLiteSessionTest { var peerBidi: FakeBidiStream? = null val peer = async { - val bidi = serverSide.peerOpenedBidiStreams().first() - readSubscribeRequest(bidi) + val (bidi, _) = nextSubscribeBidi(serverSide) bidi.write(MoqLiteCodec.encodeSubscribeOk(okFor(0L))) peerBidi = bidi // Drain whatever the listener writes after Ok — moq-lite diff --git a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsRoundTripInteropTest.kt b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsRoundTripInteropTest.kt index 15a68f876..2dd29ed29 100644 --- a/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsRoundTripInteropTest.kt +++ b/nestsClient/src/jvmTest/kotlin/com/vitorpamplona/nestsclient/interop/NostrNestsRoundTripInteropTest.kt @@ -200,9 +200,13 @@ class NostrNestsRoundTripInteropTest { assertEquals( idx.toLong(), obj.objectId, - "object id at index $idx — MoQ requires monotonic ids per group", + "object id at index $idx — MoqLiteNestsListener uses a session-level counter, monotonic across all groups", + ) + assertEquals( + idx.toLong(), + obj.groupId, + "group id at index $idx — broadcaster emits one moq-lite group per Opus frame (audio-rooms NIP draft) so groupId increments 1:1 with the frame index", ) - assertEquals(0L, obj.groupId, "single-group track per audio-rooms NIP draft") assertContentEquals( "FRAME-".encodeToByteArray() + byteArrayOf(idx.toByte()), obj.payload,