diff --git a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt index c719a8cac..00a8622af 100644 --- a/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt +++ b/nestsClient/src/commonMain/kotlin/com/vitorpamplona/nestsclient/MoqLiteNestsSpeaker.kt @@ -25,10 +25,15 @@ import com.vitorpamplona.nestsclient.audio.NestMoqLiteBroadcaster import com.vitorpamplona.nestsclient.audio.OpusEncoder import com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.launch import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock @@ -96,6 +101,26 @@ class MoqLiteNestsSpeaker internal constructor( // session's `activePublisher` slot — otherwise a subsequent // startBroadcasting fails with "already publishing" and the // publisher's announce state never tears down. + // + // Companion catalog publisher: the kixelated/moq browser + // watcher (and any standards-aligned moq-lite consumer) + // discovers a broadcast's tracks by subscribing to the + // `catalog.json` track and parsing the latest group's + // payload as a JSON manifest. Without this track our + // broadcasts are *invisible* to the canonical web watcher + // even though our audio frames are streaming fine — the + // watcher has nothing to subscribe to. Open it alongside + // audio so a single broadcast advertises both tracks. + val catalogPublisher = + try { + session.publish( + broadcastSuffix = speakerPubkeyHex, + track = MoqLiteNestsListener.CATALOG_TRACK, + ) + } catch (t: Throwable) { + runCatching { publisher.close() } + throw t + } val broadcaster = try { NestMoqLiteBroadcaster( @@ -120,6 +145,39 @@ class MoqLiteNestsSpeaker internal constructor( ) } } catch (t: Throwable) { + runCatching { catalogPublisher.close() } + runCatching { publisher.close() } + throw t + } + // Push catalog once before launching the periodic republish + // pump so the manifest is on the wire before any subscriber + // attaches; the republisher then re-emits a fresh group + // every CATALOG_REPUBLISH_INTERVAL_MS so late-joining + // watchers don't wait an arbitrary amount of time for the + // next refresh. + val catalogJson = SPEAKER_CATALOG_JSON.encodeToByteArray() + val republishJob = + try { + catalogPublisher.send(catalogJson) + catalogPublisher.endGroup() + scope.launch { + try { + while (true) { + delay(CATALOG_REPUBLISH_INTERVAL_MS) + catalogPublisher.send(catalogJson) + catalogPublisher.endGroup() + } + } catch (ce: CancellationException) { + throw ce + } catch (_: Throwable) { + // Best-effort: a transient catalog send + // failure is non-fatal — the audio path + // owns terminal-failure detection. + } + } + } catch (t: Throwable) { + runCatching { catalogPublisher.close() } + runCatching { broadcaster.stop() } runCatching { publisher.close() } throw t } @@ -133,6 +191,8 @@ class MoqLiteNestsSpeaker internal constructor( MoqLiteBroadcastHandle( broadcaster = broadcaster, publisher = publisher, + catalogPublisher = catalogPublisher, + catalogRepublishJob = republishJob, parent = this, ) activeHandle = handle @@ -140,6 +200,31 @@ class MoqLiteNestsSpeaker internal constructor( } } + companion object { + /** + * moq-lite catalog manifest broadcast on the + * [MoqLiteNestsListener.CATALOG_TRACK] track. The shape mirrors + * what the kixelated/moq browser watcher parses + * ([com.vitorpamplona.amethyst.commons.viewmodels.RoomSpeakerCatalog]): + * a versioned envelope with an `audio[]` of track descriptors. + * Hard-coded because the speaker's encoder is fixed at + * 48 kHz mono Opus today; if [OpusEncoder] becomes parameterised + * this should be derived from the encoder config instead. + */ + const val SPEAKER_CATALOG_JSON: String = + "{\"version\":1,\"audio\":[{\"track\":\"audio/data\"," + + "\"codec\":\"opus\",\"sample_rate\":48000,\"channel_count\":1}]}" + + /** + * How often the catalog group is re-emitted. The relay drops + * the catalog group when its last subscriber drops, so a + * watcher that attaches mid-broadcast needs a freshly-published + * group to receive the manifest. 2 s keeps late-attach worst + * case bounded without spamming the relay. + */ + const val CATALOG_REPUBLISH_INTERVAL_MS: Long = 2_000L + } + /** * [HotSwappablePublisherSource] implementation. See the interface * kdoc — this method mints a fresh publisher on the session @@ -228,6 +313,8 @@ class MoqLiteNestsSpeaker internal constructor( internal class MoqLiteBroadcastHandle( private val broadcaster: NestMoqLiteBroadcaster, private val publisher: MoqLitePublisherHandle, + private val catalogPublisher: MoqLitePublisherHandle, + private val catalogRepublishJob: Job, private val parent: MoqLiteNestsSpeaker, ) : BroadcastHandle { @Volatile private var muted: Boolean = false @@ -246,12 +333,27 @@ internal class MoqLiteBroadcastHandle( override suspend fun close() { if (closed) return closed = true + // Stop the catalog republisher first so it doesn't race a + // concurrent send against the publisher.close below. + try { + catalogRepublishJob.cancelAndJoin() + } catch (ce: kotlinx.coroutines.CancellationException) { + // The parent scope is being cancelled. Continue cleanup of + // the audio + publisher resources, then rethrow. + runCatching { catalogPublisher.close() } + runCatching { publisher.close() } + parent.broadcastClosed(this) + throw ce + } catch (_: Throwable) { + // Best-effort. + } try { broadcaster.stop() } catch (ce: kotlinx.coroutines.CancellationException) { // Even on cancel, run the rest of cleanup before rethrowing // — broadcaster.stop already cancels its own job, so the // mic + encoder + publisher are owed their close paths. + runCatching { catalogPublisher.close() } runCatching { publisher.close() } parent.broadcastClosed(this) throw ce @@ -263,6 +365,15 @@ internal class MoqLiteBroadcastHandle( // failures on the broadcaster.stop path. try { publisher.close() + } catch (ce: kotlinx.coroutines.CancellationException) { + runCatching { catalogPublisher.close() } + parent.broadcastClosed(this) + throw ce + } catch (_: Throwable) { + // Best-effort. + } + try { + catalogPublisher.close() } catch (ce: kotlinx.coroutines.CancellationException) { parent.broadcastClosed(this) 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 72afdb4ce..55f27fa12 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 @@ -98,8 +98,22 @@ class MoqLiteSession internal constructor( */ private var announceWatchJob: Job? = null - /** Single active publisher per session (moq-lite doesn't model multi-broadcast publishers). */ - private var activePublisher: PublisherStateImpl? = null + /** + * Active publishers on this session. moq-lite models a broadcast as + * `(suffix, set-of-tracks)`; the relay multiplexes Subscribe + * messages to whichever publisher claims the requested track. We + * allow N concurrent publishers per session as long as they share + * the same suffix (= same broadcast) and have distinct tracks. + * + * The nests audio room use case needs at least two: `audio/data` + * for Opus frames and `catalog.json` for the broadcast metadata + * the canonical kixelated/moq watcher subscribes to first to + * discover what's available. Without the catalog the broadcast is + * invisible to the JS reference watcher (the browser-side nests + * web client) — Amethyst-to-Amethyst kept working only because + * both sides hardcoded the audio track name and skipped discovery. + */ + private val activePublishers: MutableList = mutableListOf() @Volatile private var closed: Boolean = false @@ -683,8 +697,12 @@ class MoqLiteSession internal constructor( * - we open uni streams ourselves to push group data * (`session.open_uni()`) * - * Only one [publish] is supported per session for now. Calling - * [publish] twice on the same session is rejected with [IllegalStateException]. + * Multiple [publish] calls are supported on the same session as + * long as every call shares the same `broadcastSuffix` (you publish + * one broadcast per session, with multiple tracks) and uses a + * distinct `track`. Re-publishing the same `(suffix, track)` pair + * or mixing different suffixes is rejected with + * [IllegalStateException]. */ suspend fun publish( broadcastSuffix: String, @@ -695,11 +713,16 @@ class MoqLiteSession internal constructor( val publisher: PublisherStateImpl state.withLock { check(!closed) { "session is closed" } - check(activePublisher == null) { - "MoqLiteSession.publish called twice — only one broadcast per session is supported" + val existingSuffix = activePublishers.firstOrNull()?.suffix + check(existingSuffix == null || existingSuffix == normalised) { + "MoqLiteSession.publish suffix mismatch: existing='$existingSuffix', new='$normalised'. " + + "moq-lite models one broadcast per session — open a new session for a different broadcast." + } + check(activePublishers.none { it.track == track }) { + "MoqLiteSession.publish called twice for the same track '$track' on suffix '$normalised'." } publisher = PublisherStateImpl(suffix = normalised, track = track) - activePublisher = publisher + activePublishers += publisher // Lazy launch — the inbound-bidi pump needs to keep running // for the lifetime of any active publisher. if (bidiPump == null) bidiPump = scope.launch { pumpInboundBidis() } @@ -734,7 +757,21 @@ class MoqLiteSession internal constructor( } private suspend fun handleInboundBidi(bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream) { - val publisher = state.withLock { activePublisher } ?: return + // Snapshot of publishers at bidi-arrival time. All publishers + // on a single session share a suffix (enforced in [publish]), + // so the Announce response is the same regardless of which + // publisher we route the bidi to. Subscribe responses pick the + // publisher whose `track` matches `sub.track`. + val publishersSnapshot = state.withLock { activePublishers.toList() } + if (publishersSnapshot.isEmpty()) return + // Designated publisher for Announce-bidi ownership: the first + // one. The Active/Ended pair is keyed off the broadcast + // SUFFIX (not per track), so emitting it once via one + // publisher matches the single-broadcast model the relay + // expects. All publishers close together in [close], so + // there's no risk of one closing while the announce bidi + // stays alive on another. + val announcePublisher = publishersSnapshot.first() // Single long-running collector for the bidi's full lifetime. // Pre-fix this dispatch was split into a `firstOrNull()` to @@ -759,6 +796,16 @@ class MoqLiteSession internal constructor( var typeCode: Long? = null var dispatched = false var inboundSub: MoqLiteSubscribe? = null + // Track which publisher this bidi was routed to so the peer-FIN + // cleanup at the bottom calls removeInboundSubscription on the + // right one. `inboundAnnouncePublisher` is currently unused for + // cleanup (announce bidis live until publisher.close fires + // Ended), but we keep the assignment for symmetry / future + // explicit teardown. + var inboundSubPublisher: PublisherStateImpl? = null + + @Suppress("UNUSED_VARIABLE") + var inboundAnnouncePublisher: PublisherStateImpl? = null try { bidi.incoming().collect { chunk -> buffer.push(chunk) @@ -776,7 +823,7 @@ class MoqLiteSession internal constructor( val pleasePayload = buffer.readSizePrefixed() ?: return@collect val please = MoqLiteCodec.decodeAnnouncePlease(pleasePayload) val emittedSuffix = - MoqLitePath.stripPrefix(please.prefix, publisher.suffix) ?: publisher.suffix + MoqLitePath.stripPrefix(please.prefix, announcePublisher.suffix) ?: announcePublisher.suffix bidi.write( MoqLiteCodec.encodeAnnounce( MoqLiteAnnounce( @@ -786,15 +833,31 @@ class MoqLiteSession internal constructor( ), ), ) - publisher.registerAnnounceBidi(bidi, emittedSuffix) - Log.d("NestTx") { "ANNOUNCE inbound prefix='${please.prefix}' → emitted Active suffix='$emittedSuffix' (publisher.suffix='${publisher.suffix}')" } + announcePublisher.registerAnnounceBidi(bidi, emittedSuffix) + // Routed to publisher that owns the announce bidi for the lifecycle + // — for the multi-track use case (audio + catalog), all share one + // suffix and one announce bidi pointing at the first publisher. + inboundAnnouncePublisher = announcePublisher + Log.d("NestTx") { "ANNOUNCE inbound prefix='${please.prefix}' → emitted Active suffix='$emittedSuffix' (publisher.suffix='${announcePublisher.suffix}')" } dispatched = true } MoqLiteControlType.Subscribe -> { val subPayload = buffer.readSizePrefixed() ?: return@collect val sub = MoqLiteCodec.decodeSubscribe(subPayload) - Log.d("NestTx") { "SUBSCRIBE inbound id=${sub.id} broadcast='${sub.broadcast}' track='${sub.track}' (publisher track='${publisher.track}')" } + // Find the publisher that claims this track. With the + // single-track-per-publisher model, only one match is possible. + val targetPublisher = publishersSnapshot.firstOrNull { it.track == sub.track } + if (targetPublisher == null) { + Log.w("NestTx") { + "SUBSCRIBE inbound id=${sub.id} track='${sub.track}' has no matching publisher " + + "on this session (have ${publishersSnapshot.map { it.track }}) — finishing bidi" + } + runCatching { bidi.finish() } + dispatched = true + return@collect + } + Log.d("NestTx") { "SUBSCRIBE inbound id=${sub.id} broadcast='${sub.broadcast}' track='${sub.track}' (publisher track='${targetPublisher.track}')" } // Register the subscription BEFORE sending Ok so the // peer's observation of Ok is a happens-after of // `inboundSubs += sub`. Otherwise on dispatchers that @@ -803,8 +866,9 @@ class MoqLiteSession internal constructor( // (notably Windows under Dispatchers.Default), the // peer's first `publisher.send` after Ok races the // registration and observes an empty subscriber set. - publisher.registerInboundSubscription(sub) + targetPublisher.registerInboundSubscription(sub) inboundSub = sub + inboundSubPublisher = targetPublisher bidi.write( MoqLiteCodec.encodeSubscribeOk( MoqLiteSubscribeOk( @@ -843,7 +907,11 @@ class MoqLiteSession internal constructor( // 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) } + val sub = inboundSub + val pub = inboundSubPublisher + if (sub != null && pub != null) { + pub.removeInboundSubscription(sub) + } } /** @@ -872,18 +940,20 @@ class MoqLiteSession internal constructor( if (closed) return closed = true val toClose: List - val publisherToClose: PublisherStateImpl? + val publishersToClose: List state.withLock { toClose = subscriptionsBySubscribeId.values.toList() subscriptionsBySubscribeId.clear() - publisherToClose = activePublisher - activePublisher = null + publishersToClose = activePublishers.toList() + activePublishers.clear() } for (sub in toClose) { runCatching { sub.bidi.finish() } sub.frames.close() } - runCatching { publisherToClose?.close() } + for (p in publishersToClose) { + runCatching { p.close() } + } groupPump?.cancelAndJoin() bidiPump?.cancelAndJoin() runCatching { transport.close() } @@ -1081,7 +1151,7 @@ class MoqLiteSession internal constructor( runCatching { groupToFinish?.uni?.finish() } // Detach from the session so a subsequent `publish` can run. state.withLock { - if (activePublisher === this) activePublisher = null + activePublishers.remove(this@PublisherStateImpl) } }