nestsClient: expose moq-lite max_latency on subscribeSpeaker + correct plan framing

The moq-lite Lite-03 spec puts subscriber-side group-staleness control in
max_latency on the SUBSCRIBE frame: 0 = unlimited (relay falls back to its
MAX_GROUP_AGE = 30 s default); a positive value tells the relay to evict
groups older than that in favour of newer ones during transient backpressure.

Plumb it through:

  NestsListener.subscribeSpeaker(pubkey, maxLatencyMs = 0L)
    └─ MoqLiteNestsListener → MoqLiteSession.subscribe(.., maxLatencyMillis)
    └─ DefaultNestsListener (IETF) ignores — no equivalent on that wire
    └─ ReconnectingHandle re-issues across reconnects

Default stays at 0L for back-compat with the JS reference watcher; callers
opt in for low-latency live audio.

Also rewrites nestsClient/plans/2026-05-01-quic-stream-cliff-investigation.md
with corrected framing: of the four "upstream issues" the prior draft listed,
only our own missing MAX_STREAMS_UNI emission was a real bug (RFC 9000 §4.6
violation, fixed in d391ae1d). The relay's per-subscriber queue policy,
unbounded FuturesUnordered, no-timeout open_uni().await, and lack of
lagging-consumer detection are spec-permitted architectural decisions in a
gap moq-lite explicitly leaves to implementations — the spec puts staleness
control on the subscriber via max_latency, which we were sending as 0 the
whole time. Reframed as "feature request worth filing upstream", not "bugs".
This commit is contained in:
Claude
2026-05-01 18:30:37 +00:00
parent 9841e3ca2c
commit 897287be5f
7 changed files with 293 additions and 422 deletions
@@ -75,20 +75,24 @@ class MoqLiteNestsListener internal constructor(
internal val transport
get() = session.transport
override suspend fun subscribeSpeaker(speakerPubkeyHex: String): SubscribeHandle = wrapSubscription(broadcast = speakerPubkeyHex, track = AUDIO_TRACK)
override suspend fun subscribeSpeaker(
speakerPubkeyHex: String,
maxLatencyMs: Long,
): SubscribeHandle = wrapSubscription(broadcast = speakerPubkeyHex, track = AUDIO_TRACK, maxLatencyMs = maxLatencyMs)
override suspend fun subscribeCatalog(speakerPubkeyHex: String): SubscribeHandle = wrapSubscription(broadcast = speakerPubkeyHex, track = CATALOG_TRACK)
override suspend fun subscribeCatalog(speakerPubkeyHex: String): SubscribeHandle = wrapSubscription(broadcast = speakerPubkeyHex, track = CATALOG_TRACK, maxLatencyMs = 0L)
private suspend fun wrapSubscription(
broadcast: String,
track: String,
maxLatencyMs: Long,
): SubscribeHandle {
check(state.value is NestsListenerState.Connected) {
"NestsListener.subscribe requires Connected state, was ${state.value}"
}
val handle =
try {
session.subscribe(broadcast = broadcast, track = track)
session.subscribe(broadcast = broadcast, track = track, maxLatencyMillis = maxLatencyMs)
} catch (e: MoqLiteSubscribeException) {
// Surface protocol-level rejections (audio OR catalog)
// through MoqProtocolException, matching the IETF
@@ -142,7 +142,10 @@ private fun failedListener(state: MutableStateFlow<NestsListenerState>): NestsLi
object : NestsListener {
override val state = state
override suspend fun subscribeSpeaker(speakerPubkeyHex: String): SubscribeHandle = error("listener never connected: ${state.value}")
override suspend fun subscribeSpeaker(
speakerPubkeyHex: String,
maxLatencyMs: Long,
): SubscribeHandle = error("listener never connected: ${state.value}")
override suspend fun close() {
if (state.value !is NestsListenerState.Closed) {
@@ -48,10 +48,24 @@ interface NestsListener {
* Opus stream under the namespace `["nests", <roomId>]` with the
* speaker's pubkey hex as the track name.
*
* @param maxLatencyMs subscriber-side group-staleness tolerance (milliseconds)
* sent in the moq-lite SUBSCRIBE message. Per
* `kixelated/moq/doc/concept/layer/moq-lite.md` "Congestion": when the
* relay's per-track group queue can't drain fast enough, groups older
* than `maxLatencyMs` are evicted in favour of newer ones. `0`
* (default) means *unlimited* — the relay falls back to its own
* `MAX_GROUP_AGE = 30 s` default. For low-latency live audio set this
* to e.g. `500L` so the listener prefers fresh frames over stale
* buffered ones during transient relay backpressure. Ignored on the
* IETF MoQ-transport listener path (no equivalent field on that wire).
*
* @throws MoqProtocolException if the publisher rejects the subscription.
* @throws IllegalStateException if the listener is not in [NestsListenerState.Connected].
*/
suspend fun subscribeSpeaker(speakerPubkeyHex: String): SubscribeHandle
suspend fun subscribeSpeaker(
speakerPubkeyHex: String,
maxLatencyMs: Long = 0L,
): SubscribeHandle
/**
* Subscribe to a speaker's `catalog.json` track — moq-lite's
@@ -198,10 +212,18 @@ class DefaultNestsListener internal constructor(
) : NestsListener {
override val state: StateFlow<NestsListenerState> = mutableState.asStateFlow()
override suspend fun subscribeSpeaker(speakerPubkeyHex: String): SubscribeHandle {
override suspend fun subscribeSpeaker(
speakerPubkeyHex: String,
maxLatencyMs: Long,
): SubscribeHandle {
check(state.value is NestsListenerState.Connected) {
"NestsListener.subscribeSpeaker requires Connected state, was ${state.value}"
}
// IETF MoQ-transport-17 has no per-subscribe staleness tolerance
// field on the wire — `maxLatencyMs` is ignored on this path.
// Closest analogue would be the SUBSCRIBE filter, which is
// already pinned to LatestGroup so a slow listener wakes up at
// the freshest group rather than from history.
return session.subscribe(
namespace = roomNamespace,
trackName = speakerPubkeyHex.encodeToByteArray(),
@@ -217,7 +217,10 @@ private class ReconnectingHandle(
) : NestsListener {
override val state: StateFlow<NestsListenerState> = mutableState.asStateFlow()
override suspend fun subscribeSpeaker(speakerPubkeyHex: String): SubscribeHandle = reissuingSubscribe { listener -> listener.subscribeSpeaker(speakerPubkeyHex) }
override suspend fun subscribeSpeaker(
speakerPubkeyHex: String,
maxLatencyMs: Long,
): SubscribeHandle = reissuingSubscribe { listener -> listener.subscribeSpeaker(speakerPubkeyHex, maxLatencyMs) }
override suspend fun subscribeCatalog(speakerPubkeyHex: String): SubscribeHandle = reissuingSubscribe { listener -> listener.subscribeCatalog(speakerPubkeyHex) }