refactor(nests): catalog emit-on-subscribe hook replaces 2s republish loop

The catalog publisher's previous shape was "send once at startup
(silently fails if no subs are attached yet) + a 2 s republish loop".
That worked but had two problems:

1. The startup send is a no-op (PublisherStateImpl.send returns false
   when inboundSubs is empty), so the catalog isn't actually on the
   wire until the first loop tick — up to 2 s of catalog dead time
   for a watcher that subscribes immediately.
2. Re-emitting a static blob every 2 s for the broadcast's lifetime
   wastes a fresh group on the relay's per-track cache; not a hot path
   today but unnecessary.

Replace with an emit-on-subscribe hook. New API on MoqLitePublisherHandle:

  fun setOnNewSubscriber(hook: (suspend () -> Unit)?)

PublisherStateImpl fires the hook once per accepted (track-matching)
inbound SUBSCRIBE, OUTSIDE its serialisation lock so the hook can
safely call send/endGroup without deadlock. Track-mismatched subs
(SubscribeDrop reply) do NOT fire the hook — covered by a new
"hook does not fire on track mismatch" test.

MoqLiteNestsSpeaker.startBroadcasting and ReconnectingNestsSpeaker's
hot-swap iteration both now register a hook that writes the catalog
JSON on every fresh subscribe instead of looping. The
catalogRepublishJob field on MoqLiteBroadcastHandle and the
hotSwapCatalogJob field on ReissuingBroadcastHandle (plus their close
paths) are dropped.

Net effect for hang.js / NestsUI-v2 watchers:
  - Catalog reaches the watcher on the SAME group as the subscribe ack
    (no 0–2 s startup gap).
  - Late joiners hit the relay's track-latest cache for the catalog
    blob — a fresh hook fire only happens when the relay opens a new
    SUBSCRIBE bidi to us (cliff-detector recycle on listener side, or
    JWT-refresh on our side).
  - One catalog group per relay-side subscribe instead of one every
    2 s for the broadcast lifetime.

Two tests pin the new behaviour: hook fires on inbound subscribe;
hook does NOT fire on a track-mismatched subscribe (which gets a
SubscribeDrop reply and never enters inboundSubs).

https://claude.ai/code/session_014JfZJHSTvyYYWJbC9VbB47
This commit is contained in:
Claude
2026-05-06 16:10:12 +00:00
parent 79c76d5831
commit 3417e44a6a
4 changed files with 200 additions and 103 deletions
@@ -28,13 +28,9 @@ import com.vitorpamplona.nestsclient.moq.lite.MoqLitePublisherHandle
import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession import com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession
import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope 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.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock import kotlinx.coroutines.sync.withLock
@@ -150,39 +146,25 @@ class MoqLiteNestsSpeaker internal constructor(
runCatching { publisher.close() } runCatching { publisher.close() }
throw t throw t
} }
// Push catalog once before launching the periodic republish // Catalog emit-on-subscribe: every time the relay opens a
// pump so the manifest is on the wire before any subscriber // SUBSCRIBE bidi for catalog.json, fire the hook to write
// attaches; the republisher then re-emits a fresh group // one group + FIN. moq-lite serves new listeners from the
// every CATALOG_REPUBLISH_INTERVAL_MS so late-joining // relay's per-track latest-group cache, so emitting once
// watchers don't wait an arbitrary amount of time for the // per relay-side subscribe is enough — late-joining
// next refresh. // watchers behind the same relay get the cached blob
// without us having to maintain a periodic re-emit loop.
// Set BEFORE the relay can race a SUBSCRIBE in; in
// practice the relay's SUBSCRIBE bidi takes a network
// round-trip after our ANNOUNCE Active, so this is safe
// even though the setter is non-suspending.
val catalogJson = val catalogJson =
MoqLiteHangCatalog.opusMono48k(MoqLiteNestsListener.AUDIO_TRACK).encodeJsonBytes() MoqLiteHangCatalog.opusMono48k(MoqLiteNestsListener.AUDIO_TRACK).encodeJsonBytes()
val republishJob = catalogPublisher.setOnNewSubscriber {
try { runCatching {
catalogPublisher.send(catalogJson) catalogPublisher.send(catalogJson)
catalogPublisher.endGroup() 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
} }
}
mutableState.value = mutableState.value =
NestsSpeakerState.Broadcasting( NestsSpeakerState.Broadcasting(
room = current.room, room = current.room,
@@ -194,7 +176,6 @@ class MoqLiteNestsSpeaker internal constructor(
broadcaster = broadcaster, broadcaster = broadcaster,
publisher = publisher, publisher = publisher,
catalogPublisher = catalogPublisher, catalogPublisher = catalogPublisher,
catalogRepublishJob = republishJob,
parent = this, parent = this,
) )
activeHandle = handle activeHandle = handle
@@ -202,17 +183,6 @@ class MoqLiteNestsSpeaker internal constructor(
} }
} }
companion object {
/**
* 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 * [HotSwappablePublisherSource] implementation. See the interface
* kdoc — this method mints a fresh publisher on the session * kdoc — this method mints a fresh publisher on the session
@@ -302,7 +272,6 @@ internal class MoqLiteBroadcastHandle(
private val broadcaster: NestMoqLiteBroadcaster, private val broadcaster: NestMoqLiteBroadcaster,
private val publisher: MoqLitePublisherHandle, private val publisher: MoqLitePublisherHandle,
private val catalogPublisher: MoqLitePublisherHandle, private val catalogPublisher: MoqLitePublisherHandle,
private val catalogRepublishJob: Job,
private val parent: MoqLiteNestsSpeaker, private val parent: MoqLiteNestsSpeaker,
) : BroadcastHandle { ) : BroadcastHandle {
@Volatile private var muted: Boolean = false @Volatile private var muted: Boolean = false
@@ -321,20 +290,8 @@ internal class MoqLiteBroadcastHandle(
override suspend fun close() { override suspend fun close() {
if (closed) return if (closed) return
closed = true closed = true
// Stop the catalog republisher first so it doesn't race a // Stop the broadcaster first so the audio capture + encoder
// concurrent send against the publisher.close below. // don't keep producing into a closing publisher.
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 { try {
broadcaster.stop() broadcaster.stop()
} catch (ce: kotlinx.coroutines.CancellationException) { } catch (ce: kotlinx.coroutines.CancellationException) {
@@ -30,7 +30,6 @@ import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Job import kotlinx.coroutines.Job
import kotlinx.coroutines.awaitCancellation import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.cancelAndJoin
import kotlinx.coroutines.delay import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.StateFlow
@@ -442,8 +441,6 @@ private class ReissuingBroadcastHandle(
*/ */
@Volatile private var hotSwapCatalogPublisher: MoqLitePublisherHandle? = null @Volatile private var hotSwapCatalogPublisher: MoqLitePublisherHandle? = null
@Volatile private var hotSwapCatalogJob: Job? = null
/** Legacy path's per-session handle. Cleared when the session swaps. */ /** Legacy path's per-session handle. Cleared when the session swaps. */
private val liveHandle = AtomicReference<BroadcastHandle?>(null) private val liveHandle = AtomicReference<BroadcastHandle?>(null)
private var pumpJob: Job? = null private var pumpJob: Job? = null
@@ -572,11 +569,10 @@ private class ReissuingBroadcastHandle(
// this the catalog goes silent the moment the session recycles // this the catalog goes silent the moment the session recycles
// and any watcher that attaches AFTER the recycle sees nothing // and any watcher that attaches AFTER the recycle sees nothing
// to subscribe to. Mirror of [MoqLiteNestsSpeaker.startBroadcasting]'s // to subscribe to. Mirror of [MoqLiteNestsSpeaker.startBroadcasting]'s
// catalog setup; same JSON, same republish cadence. // catalog setup; same JSON, same emit-on-subscribe pattern.
val catalogPayload = val catalogPayload =
MoqLiteHangCatalog.opusMono48k(MoqLiteNestsListener.AUDIO_TRACK).encodeJsonBytes() MoqLiteHangCatalog.opusMono48k(MoqLiteNestsListener.AUDIO_TRACK).encodeJsonBytes()
val priorCatalogPublisher = hotSwapCatalogPublisher val priorCatalogPublisher = hotSwapCatalogPublisher
val priorCatalogJob = hotSwapCatalogJob
val newCatalogPublisher = val newCatalogPublisher =
try { try {
hotSwap.openPublisherForHotSwap(MoqLiteNestsListener.CATALOG_TRACK) hotSwap.openPublisherForHotSwap(MoqLiteNestsListener.CATALOG_TRACK)
@@ -590,45 +586,23 @@ private class ReissuingBroadcastHandle(
null null
} }
if (newCatalogPublisher != null) { if (newCatalogPublisher != null) {
val newCatalogJob = // Set the emit-on-subscribe hook BEFORE installing the
try { // publisher reference, so the relay can't race a SUBSCRIBE
// in between.
newCatalogPublisher.setOnNewSubscriber {
runCatching {
newCatalogPublisher.send(catalogPayload) newCatalogPublisher.send(catalogPayload)
newCatalogPublisher.endGroup() newCatalogPublisher.endGroup()
scope.launch {
try {
while (true) {
delay(MoqLiteNestsSpeaker.CATALOG_REPUBLISH_INTERVAL_MS)
newCatalogPublisher.send(catalogPayload)
newCatalogPublisher.endGroup()
}
} catch (ce: kotlinx.coroutines.CancellationException) {
throw ce
} catch (_: Throwable) {
// Best-effort: a transient catalog send
// failure on this session is non-fatal —
// the audio path owns terminal-failure
// detection.
}
}
} catch (ce: kotlinx.coroutines.CancellationException) {
runCatching { newCatalogPublisher.close() }
throw ce
} catch (_: Throwable) {
runCatching { newCatalogPublisher.close() }
null
} }
if (newCatalogJob != null) {
hotSwapCatalogPublisher = newCatalogPublisher
hotSwapCatalogJob = newCatalogJob
// Tear down the prior session's catalog state AFTER the
// new one is installed so the watcher only sees a brief
// overlap rather than a gap. The prior publisher's
// session is about to be torn down by the orchestrator
// anyway, but graceful Announce(Ended) keeps the relay
// book-keeping clean.
if (priorCatalogJob != null) runCatching { priorCatalogJob.cancelAndJoin() }
if (priorCatalogPublisher != null) runCatching { priorCatalogPublisher.close() }
} }
hotSwapCatalogPublisher = newCatalogPublisher
// Tear down the prior session's catalog publisher AFTER the
// new one is installed so the watcher only sees a brief
// overlap rather than a gap. The prior publisher's session
// is about to be torn down by the orchestrator anyway, but
// graceful Announce(Ended) keeps the relay book-keeping
// clean.
if (priorCatalogPublisher != null) runCatching { priorCatalogPublisher.close() }
} }
try { try {
@@ -703,10 +677,10 @@ private class ReissuingBroadcastHandle(
// //
// Catalog gets explicit teardown because — unlike the audio // Catalog gets explicit teardown because — unlike the audio
// publisher which broadcaster.stop() closes for us — the // publisher which broadcaster.stop() closes for us — the
// catalog publisher + republisher have no broadcaster wrapper. // catalog publisher has no broadcaster wrapper. The
// Cancel the republish loop FIRST so it can't race the close. // emit-on-subscribe hook itself doesn't need shutdown: it
hotSwapCatalogJob?.let { runCatching { it.cancelAndJoin() } } // only fires inside [PublisherStateImpl.registerInboundSubscription],
hotSwapCatalogJob = null // and closing the publisher prevents any further hook fires.
hotSwapCatalogPublisher?.let { runCatching { it.close() } } hotSwapCatalogPublisher?.let { runCatching { it.close() } }
hotSwapCatalogPublisher = null hotSwapCatalogPublisher = null
hotSwapBroadcaster?.let { runCatching { it.stop() } } hotSwapBroadcaster?.let { runCatching { it.stop() } }
@@ -1079,6 +1079,18 @@ class MoqLiteSession internal constructor(
@Volatile private var publisherClosed = false @Volatile private var publisherClosed = false
/**
* Caller-installed hook fired once per accepted inbound
* SUBSCRIBE — see [MoqLitePublisherHandle.setOnNewSubscriber].
* Read-and-fire happens OUTSIDE [gate] to avoid deadlocking on
* the hook's own calls to [send] / [endGroup].
*/
@Volatile private var onNewSubscriberHook: (suspend () -> Unit)? = null
override fun setOnNewSubscriber(hook: (suspend () -> Unit)?) {
onNewSubscriberHook = hook
}
suspend fun registerAnnounceBidi( suspend fun registerAnnounceBidi(
bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream, bidi: com.vitorpamplona.nestsclient.transport.WebTransportBidiStream,
emittedSuffix: String, emittedSuffix: String,
@@ -1093,6 +1105,11 @@ class MoqLiteSession internal constructor(
} }
suspend fun registerInboundSubscription(sub: MoqLiteSubscribe) { suspend fun registerInboundSubscription(sub: MoqLiteSubscribe) {
// Capture the hook INSIDE the lock — guarantees the hook
// observes a fully-registered subscriber when it fires —
// but invoke it OUTSIDE so a hook that calls [send] /
// [endGroup] doesn't deadlock on the same gate.
var hookToFire: (suspend () -> Unit)? = null
gate.withLock { gate.withLock {
if (publisherClosed) { if (publisherClosed) {
Log.w("NestTx") { "SUBSCRIBE inbound rejected (publisher closed) id=${sub.id} track='${sub.track}'" } Log.w("NestTx") { "SUBSCRIBE inbound rejected (publisher closed) id=${sub.id} track='${sub.track}'" }
@@ -1103,8 +1120,15 @@ class MoqLiteSession internal constructor(
return return
} }
inboundSubs += sub inboundSubs += sub
hookToFire = onNewSubscriberHook
Log.d("NestTx") { "SUBSCRIBE registered id=${sub.id} broadcast='${sub.broadcast}' track='${sub.track}' inboundSubs.size=${inboundSubs.size}" } Log.d("NestTx") { "SUBSCRIBE registered id=${sub.id} broadcast='${sub.broadcast}' track='${sub.track}' inboundSubs.size=${inboundSubs.size}" }
} }
// Launch on the session's scope so the hook outlives the
// bidi pump's per-bidi coroutine if it's slow (e.g. a
// catalog send blocked on transport backpressure).
hookToFire?.let { hook ->
scope.launch { runCatching { hook.invoke() } }
}
} }
/** /**
@@ -1388,6 +1412,31 @@ interface MoqLitePublisherHandle {
/** FIN the current group's uni stream. The next [send] starts a fresh group. */ /** FIN the current group's uni stream. The next [send] starts a fresh group. */
suspend fun endGroup() suspend fun endGroup()
/**
* Register a callback that fires once each time a new inbound
* subscriber is registered against this publisher's track (i.e.
* each track-matching SUBSCRIBE bidi the relay opens to us). Used
* to push a "track-latest" payload — the canonical example is the
* broadcast catalog manifest, which a watcher needs to receive on
* subscribe but doesn't change between subscribers — without
* forcing the publisher to maintain a periodic re-emit loop.
*
* Called once per accepted SUBSCRIBE (track filter passed). Fires
* OUTSIDE the publisher's serialisation lock, so the hook can
* safely call [send] / [endGroup] without deadlocking.
*
* Caller MUST set the hook before any subscriber attaches (typically
* immediately after [com.vitorpamplona.nestsclient.moq.lite.MoqLiteSession.publish]
* returns) — there's no "fire-on-set for existing subscribers"
* replay. For the typical catalog use case the publisher is fresh
* when the hook is set, and the relay's SUBSCRIBE bidi takes a
* round-trip to arrive, so this is safe in practice.
*
* Pass `null` to clear the hook. Calling twice with non-null
* replaces the previous hook (no de-duplication).
*/
fun setOnNewSubscriber(hook: (suspend () -> Unit)?)
/** /**
* Stop publishing. Sends `Announce(Ended)` on every active announce * Stop publishing. Sends `Announce(Ended)` on every active announce
* bidi, FINs the current group, and releases all per-publisher * bidi, FINs the current group, and releases all per-publisher
@@ -357,6 +357,123 @@ class MoqLiteSessionTest {
session.close() session.close()
} }
@Test
fun publisher_setOnNewSubscriber_hook_fires_per_inbound_subscribe() =
runBlocking {
val (clientSide, serverSide) = FakeWebTransport.pair()
val session = MoqLiteSession.client(clientSide, pumpScope)
val publisher = session.publish(broadcastSuffix = "speakerPubkey", track = "audio/data")
val hookFireCount =
java.util.concurrent.atomic
.AtomicInteger(0)
publisher.setOnNewSubscriber {
hookFireCount.incrementAndGet()
}
// First inbound SUBSCRIBE → hook fires once.
val subBidi1 = serverSide.openBidiStream()
subBidi1.write(Varint.encode(MoqLiteControlType.Subscribe.code))
subBidi1.write(
MoqLiteCodec.encodeSubscribe(
MoqLiteSubscribe(
id = 0L,
broadcast = "speakerPubkey",
track = "audio/data",
priority = 0x80,
ordered = true,
maxLatencyMillis = 0L,
startGroup = null,
endGroup = null,
),
),
)
// Wait for SubscribeOk to drain so we know registerInboundSubscription
// ran (and therefore the hook had its chance to launch).
withTimeout(2_000) { subBidi1.incoming().first() }
// Hook fires asynchronously on the session scope; give it a
// moment to land. Use a bounded retry rather than a flat
// delay so the test is fast on the happy path.
withTimeout(2_000) {
while (hookFireCount.get() < 1) kotlinx.coroutines.yield()
}
assertEquals(1, hookFireCount.get())
// Second inbound SUBSCRIBE → hook fires again.
val subBidi2 = serverSide.openBidiStream()
subBidi2.write(Varint.encode(MoqLiteControlType.Subscribe.code))
subBidi2.write(
MoqLiteCodec.encodeSubscribe(
MoqLiteSubscribe(
id = 1L,
broadcast = "speakerPubkey",
track = "audio/data",
priority = 0x80,
ordered = true,
maxLatencyMillis = 0L,
startGroup = null,
endGroup = null,
),
),
)
withTimeout(2_000) { subBidi2.incoming().first() }
withTimeout(2_000) {
while (hookFireCount.get() < 2) kotlinx.coroutines.yield()
}
assertEquals(2, hookFireCount.get())
publisher.close()
session.close()
}
@Test
fun publisher_setOnNewSubscriber_hook_does_not_fire_on_track_mismatch() =
runBlocking {
val (clientSide, serverSide) = FakeWebTransport.pair()
val session = MoqLiteSession.client(clientSide, pumpScope)
// Publisher serves audio/data only.
val publisher = session.publish(broadcastSuffix = "speakerPubkey", track = "audio/data")
val hookFireCount =
java.util.concurrent.atomic
.AtomicInteger(0)
publisher.setOnNewSubscriber {
hookFireCount.incrementAndGet()
}
// Inbound SUBSCRIBE for a different track. Replies
// SubscribeDrop (covered in another test); hook MUST NOT
// fire because no subscriber was actually registered on
// this publisher.
val subBidi = serverSide.openBidiStream()
subBidi.write(Varint.encode(MoqLiteControlType.Subscribe.code))
subBidi.write(
MoqLiteCodec.encodeSubscribe(
MoqLiteSubscribe(
id = 7L,
broadcast = "speakerPubkey",
track = "video/data",
priority = 0x80,
ordered = true,
maxLatencyMillis = 0L,
startGroup = null,
endGroup = null,
),
),
)
// Drain the Drop reply.
withTimeout(2_000) { subBidi.incoming().first() }
// No way to wait deterministically for "the hook didn't
// fire"; sleep briefly to let any racing launch surface,
// then assert. Short delay because the hook would launch
// on the same scope as registerInboundSubscription's caller.
kotlinx.coroutines.delay(100)
assertEquals(0, hookFireCount.get())
publisher.close()
session.close()
}
@Test @Test
fun publisher_replies_subscribeDrop_when_track_is_not_published() = fun publisher_replies_subscribeDrop_when_track_is_not_published() =
runBlocking { runBlocking {