Merge pull request #2750 from vitorpamplona/claude/debug-audio-dropout-n0g6Z

Add audio format negotiation and cliff detection for Nests
This commit is contained in:
Vitor Pamplona
2026-05-06 16:07:59 -04:00
committed by GitHub
40 changed files with 4578 additions and 578 deletions
@@ -0,0 +1,91 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.viewmodels
import com.vitorpamplona.nestsclient.audio.AudioPlayer
import com.vitorpamplona.nestsclient.audio.NestPlayer
import com.vitorpamplona.nestsclient.moq.SubscribeHandle
/**
* Per-pubkey state slot held in [NestViewModel.activeSubscriptions].
*
* Three lifecycle phases:
* - **Pending**: just-reserved by `reconcileSubscriptions` to dedupe
* concurrent reconciles. No handle / player attached yet. Constructor
* via [pending].
* - **Active**: [attach] wires the moq subscribe handle, the
* [NestPlayer] decode loop, and the device player. [isPlaying]
* flips true.
* - **Detached**: [detach] returns the handle + roomPlayer pair so the
* caller can run the suspending teardown (`NestPlayer.stop()` +
* `SubscribeHandle.unsubscribe()`) in its own coroutine scope —
* keeps the native MediaCodec / AudioTrack release ordered after
* the decode loop has unwound (audit MoQ #11/#12).
*
* `internal` because only [NestViewModel]'s subscription-lifecycle
* paths (open / close / reconcile / mute / hush) construct or mutate
* these. Lifted to a top-level type rather than a nested class so the
* file split tracks concerns: this class is "subscription state
* machine"; the rest of NestViewModel is "ViewModel public surface +
* orchestration."
*/
internal class ActiveSubscription private constructor(
val pubkey: String,
) {
private var handle: SubscribeHandle? = null
private var roomPlayer: NestPlayer? = null
var player: AudioPlayer? = null
private set
var isPlaying: Boolean = false
private set
fun attach(
handle: SubscribeHandle,
roomPlayer: NestPlayer,
player: AudioPlayer,
) {
this.handle = handle
this.roomPlayer = roomPlayer
this.player = player
this.isPlaying = true
}
/**
* Hand the player + handle back to the caller's coroutine scope —
* `NestPlayer.stop()` and `SubscribeHandle.unsubscribe()` are
* both suspend, and the caller has the right scope to await them
* (so native MediaCodec/AudioTrack release runs after the decode
* loop has unwound, per audit MoQ #11/#12).
*/
fun detach(): Pair<NestPlayer?, SubscribeHandle?> {
isPlaying = false
val p = roomPlayer
val h = handle
roomPlayer = null
handle = null
player = null
return p to h
}
companion object {
fun pending(pubkey: String) = ActiveSubscription(pubkey)
}
}
@@ -35,16 +35,17 @@ import com.vitorpamplona.nestsclient.NestsSpeaker
import com.vitorpamplona.nestsclient.NestsSpeakerState
import com.vitorpamplona.nestsclient.audio.AudioCapture
import com.vitorpamplona.nestsclient.audio.AudioException
import com.vitorpamplona.nestsclient.audio.AudioFormat
import com.vitorpamplona.nestsclient.audio.AudioPlayer
import com.vitorpamplona.nestsclient.audio.NestPlayer
import com.vitorpamplona.nestsclient.audio.OpusDecoder
import com.vitorpamplona.nestsclient.audio.OpusEncoder
import com.vitorpamplona.nestsclient.connectReconnectingNestsListener
import com.vitorpamplona.nestsclient.connectReconnectingNestsSpeaker
import com.vitorpamplona.nestsclient.moq.SubscribeHandle
import com.vitorpamplona.nestsclient.transport.WebTransportFactory
import com.vitorpamplona.quartz.nip01Core.signers.NostrSigner
import com.vitorpamplona.quartz.nip53LiveActivities.chat.LiveActivitiesChatMessageEvent
import com.vitorpamplona.quartz.utils.Log
import kotlinx.collections.immutable.ImmutableSet
import kotlinx.collections.immutable.persistentSetOf
import kotlinx.collections.immutable.toPersistentSet
@@ -55,9 +56,13 @@ import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
import kotlinx.coroutines.withTimeoutOrNull
import kotlin.coroutines.cancellation.CancellationException
/**
@@ -80,7 +85,12 @@ import kotlin.coroutines.cancellation.CancellationException
* Audio-pipeline construction is injected via [decoderFactory] /
* [playerFactory] so commonMain doesn't have to know which platform's
* MediaCodec / AudioTrack is in play. M1 wires Android-only — desktop
* passes nothing here yet.
* passes nothing here yet. Both factories take a `channelCount` (1 for
* mono, 2 for stereo) AND a `sampleRate` (Hz) discovered from the
* publisher's `catalog.json` audio rendition; subscriptions await a
* brief catalog-arrival window before constructing the decoder + player
* so a stereo or non-48 kHz web publisher doesn't get its frames
* decoded with the wrong channel layout / clock.
*
* **Threading contract:** all public methods (`connect`, `disconnect`,
* `updateSpeakers`, `setMuted`, `setMicMuted`, `startBroadcast`,
@@ -97,8 +107,8 @@ import kotlin.coroutines.cancellation.CancellationException
class NestViewModel(
private val httpClient: NestsClient,
private val transport: WebTransportFactory,
private val decoderFactory: () -> OpusDecoder,
private val playerFactory: () -> AudioPlayer,
private val decoderFactory: (channelCount: Int, sampleRate: Int) -> OpusDecoder,
private val playerFactory: (channelCount: Int, sampleRate: Int) -> AudioPlayer,
private val signer: NostrSigner,
private val room: NestsRoomConfig,
// Speaker-side audio capture/encode actuals. Optional — desktop and
@@ -265,6 +275,41 @@ class NestViewModel(
private val activeSubscriptions = mutableMapOf<String, ActiveSubscription>()
private val speakingExpiryJobs = mutableMapOf<String, Job>()
/**
* Monotonic [kotlin.time.TimeMark] of the last MoQ object delivered
* to the consumer flow per speaker pubkey. Updated in
* [onSpeakerActivity] alongside the speaking-now flag. Read by
* [cliffDetectorJob] to spot the relay-side forward-queue cliff
* documented in
* `nestsClient/plans/2026-05-01-quic-stream-cliff-investigation.md`:
* the relay stops opening new uni streams to the listener while
* the announce stream still says the broadcast is Active and the
* subscribe bidi is still open from QUIC's POV — meaning no other
* layer notices.
*/
private val lastFrameAt = mutableMapOf<String, kotlin.time.TimeMark>()
/**
* Periodic scan of [activeSubscriptions] vs [_announcedSpeakers]
* vs [lastFrameAt]. When a speaker has been announced Active
* AND we have an active subscription AND no frames have arrived
* for [ROOM_AUDIO_CLIFF_TIMEOUT_MS], force a `recycleSession()`
* so the orchestrator opens a fresh QUIC transport — a clean
* recovery from the relay-side per-subscriber forward-queue
* starvation that has no QUIC- or moq-lite-level signal.
*/
private var cliffDetectorJob: Job? = null
/**
* [kotlin.time.TimeMark] of the most recent recycle the cliff
* detector triggered, or `null` if it has never fired. Acts as a
* cooldown so a single cliff event doesn't kick off multiple
* recycles back-to-back while the wrapper is mid-handshake on
* the new session — the new session has no incoming frames yet,
* which would otherwise trip the detector again immediately.
*/
private var lastCliffRecycleAt: kotlin.time.TimeMark? = null
/**
* Per-speaker catalog-fetch coroutines. Each entry is the
* background `subscribeCatalog` collector launched in
@@ -561,7 +606,10 @@ class NestViewModel(
* session over a separate WebTransport — nests' protocol uses one
* session per direction.
*/
fun startBroadcast(speakerPubkeyHex: String) {
fun startBroadcast(
speakerPubkeyHex: String,
initialMuted: Boolean = false,
) {
if (closed || !canBroadcast) return
if (_uiState.value.connection !is ConnectionUiState.Connected) return
if (_uiState.value.broadcast is BroadcastUiState.Broadcasting ||
@@ -608,7 +656,26 @@ class NestViewModel(
},
)
broadcastHandle = handle
_uiState.update { it.copy(broadcast = BroadcastUiState.Broadcasting(isMuted = false)) }
// Auto-start path (room creation): the host wants
// listeners to be able to subscribe BEFORE they hit
// unmute, so we open the publisher session pre-muted.
// The mic stays open + the broadcaster keeps
// capturing + encoding, but `setMuted(true)` makes
// every send drop on the floor inside
// `NestMoqLiteBroadcaster.start`'s `if (muted) continue`
// gate. Listener-side announces fire as soon as the
// publisher registers with the relay, so the audience
// can subscribe immediately and switch from
// "Connecting…" to "Live (silent)" while they wait
// for the host's first word. Composes with `focusMuted`
// exactly the same way [setMicMuted] does so an
// inbound phone call during the auto-start window
// doesn't get a silent unmute.
if (initialMuted) {
val effective = initialMuted || focusMuted
runCatching { handle.setMuted(effective) }
}
_uiState.update { it.copy(broadcast = BroadcastUiState.Broadcasting(isMuted = initialMuted)) }
} catch (ce: CancellationException) {
throw ce
} catch (t: Throwable) {
@@ -763,15 +830,31 @@ class NestViewModel(
announcesJob?.cancel()
announcesJob =
viewModelScope.launch {
runCatching {
l.announces().collect { ann ->
if (closed) return@collect
if (ann.active) {
_announcedSpeakers.update { it + ann.pubkey }
} else {
_announcedSpeakers.update { it - ann.pubkey }
com.vitorpamplona.quartz.utils.Log
.d("NestRx") { "observeAnnounces launching collect on l.announces()" }
var emissionCount = 0
val outcome =
runCatching {
l.announces().collect { ann ->
if (closed) return@collect
emissionCount += 1
com.vitorpamplona.quartz.utils.Log.d("NestRx") {
"observeAnnounces #$emissionCount active=${ann.active} pubkey='${ann.pubkey.take(12)}' → ${if (ann.active) "ADD" else "REMOVE"}"
}
com.vitorpamplona.nestsclient.trace.NestsTrace.emit("vm_observe_announce") {
"\"emission\":$emissionCount,\"active\":${ann.active}," +
"\"pubkey\":${com.vitorpamplona.nestsclient.trace.jsonStr(ann.pubkey)}"
}
if (ann.active) {
_announcedSpeakers.update { it + ann.pubkey }
} else {
_announcedSpeakers.update { it - ann.pubkey }
}
}
}
com.vitorpamplona.quartz.utils.Log.w("NestRx") {
val why = outcome.exceptionOrNull()?.let { "${it::class.simpleName}: ${it.message}" } ?: "naturally"
"observeAnnounces collect ENDED $why (emissions=$emissionCount, _announcedSpeakers.size=${_announcedSpeakers.value.size})"
}
}
}
@@ -794,6 +877,18 @@ class NestViewModel(
// this is typically a no-op in toAdd / toRemove
// — runs anyway for the first-Connected path.
reconcileSubscriptions()
// Defensive restart: if the wrapper went
// through a Closed transient (e.g. the cliff
// detector itself triggered `recycleSession()`,
// which closes the inner listener and waits
// for the orchestrator to reopen), our
// [Closed] branch below cancelled the
// detector via [teardown]. We need it back
// up on the fresh session — otherwise a
// second cliff event during the same room
// sits undetected. Idempotent: returns
// immediately if a job is already active.
startCliffDetector()
}
is NestsListenerState.Failed -> {
@@ -802,11 +897,34 @@ class NestViewModel(
// wait for a manual reconnect tap.
}
// Server-initiated Closed: tear down stale
// local state so any later user-driven reconnect
// starts fresh.
NestsListenerState.Closed -> {
teardown(targetState = ConnectionUiState.Closed)
// Closed from the wrapper is *almost always*
// transient — the orchestrator emits Closed
// → Reconnecting → Connecting → Connected
// around every cliff-detector recycle and
// every JWT refresh. We must NOT teardown
// here: teardown cancels `cliffDetectorJob`,
// `announcesJob`, etc., AND calls
// `wrapper.close()`, which cancels the
// wrapper's orchestrator before it can
// reopen the inner listener. End result
// pre-fix: the very first cliff-detector
// recycle permanently kills the room
// instead of recovering it (visible in the
// 15:56:25 receiver log: cliff fires →
// teardown fires 521 ms later → wrapper
// never reopens → `cliff-detector EXITED
// closed=false`).
//
// User-initiated close goes through
// `disconnect()` / `onCleared()` which call
// `teardown` directly; the wrapper's
// subsequent Closed emission here is a
// redundant signal — no-op'ing it doesn't
// change anything for those paths. The UI
// still picks up the Closed via
// `state.toUiState(ui.connection)` above
// for visual feedback.
}
else -> { /* no extra side effect */ }
@@ -855,6 +973,7 @@ class NestViewModel(
listener = l
observeListenerState(l)
observeAnnounces(l)
startCliffDetector()
} catch (ce: CancellationException) {
throw ce
} catch (t: Throwable) {
@@ -915,6 +1034,11 @@ class NestViewModel(
if (rawAudioLevels.remove(slot.pubkey) != null) {
_audioLevels.value = rawAudioLevels.toMap()
}
// Drop the cliff-detector heartbeat for this pubkey so the next
// tick doesn't see a stale "last frame was 30s ago" entry for a
// speaker we already stopped subscribing to and trigger a
// gratuitous recycle.
lastFrameAt.remove(slot.pubkey)
if (roomPlayer != null || handle != null) {
viewModelScope.launch {
roomPlayer?.runCatching { stop() }
@@ -923,6 +1047,92 @@ class NestViewModel(
}
}
/**
* Wait briefly for [pubkey]'s catalog to land in [_speakerCatalogs]
* and pick the audio config (channel count + sample rate) for the
* decoder + AudioTrack. Falls back to [AudioFormat.CHANNELS] /
* [AudioFormat.SAMPLE_RATE_HZ] on timeout (the catalog never
* arrived within [timeoutMs]) or when the catalog declares
* unsupported values (channelCount outside `1..2`, or non-positive
* sampleRate).
*
* Caller is responsible for having started [fetchSpeakerCatalog]
* first; otherwise this always times out.
*
* Why a timeout: the catalog handshake is best-effort. If the
* publisher doesn't publish a catalog (legacy publishers) or the
* relay never forwards the catalog group, audio playback should
* still proceed in the default config rather than block forever.
* 500 ms is generous — with the speaker-side
* `setOnNewSubscriber` hook the catalog group is on the wire
* within one round-trip of the SUBSCRIBE_OK, so this typically
* returns within tens of ms.
*/
private suspend fun awaitAudioPipelineConfig(
pubkey: String,
timeoutMs: Long,
): AudioPipelineConfig {
val catalog =
withTimeoutOrNull(timeoutMs) {
_speakerCatalogs
.map { it[pubkey] }
.filterNotNull()
.first()
}
val rendition = catalog?.primaryAudio()
val declaredChannels = rendition?.numberOfChannels
val channels =
when {
declaredChannels == null -> {
AudioFormat.CHANNELS
}
declaredChannels !in 1..2 -> {
Log.w("NestRx") {
"publisher catalog for pubkey='${pubkey.take(8)}' declares numberOfChannels=$declaredChannels " +
"(only 1 / 2 supported); falling back to mono"
}
AudioFormat.CHANNELS
}
else -> {
declaredChannels
}
}
val declaredRate = rendition?.sampleRate
val rate =
when {
declaredRate == null -> {
AudioFormat.SAMPLE_RATE_HZ
}
declaredRate <= 0 -> {
Log.w("NestRx") {
"publisher catalog for pubkey='${pubkey.take(8)}' declares sampleRate=$declaredRate " +
"(must be positive); falling back to ${AudioFormat.SAMPLE_RATE_HZ} Hz"
}
AudioFormat.SAMPLE_RATE_HZ
}
else -> {
declaredRate
}
}
return AudioPipelineConfig(channelCount = channels, sampleRate = rate)
}
/**
* Output of [awaitAudioPipelineConfig] — the validated audio
* pipeline configuration for one subscription. Internal because the
* factories' two-arg shape (`(channelCount, sampleRate)`) is what
* the public surface deals in; this struct just bundles them inside
* `openSubscription`.
*/
private data class AudioPipelineConfig(
val channelCount: Int,
val sampleRate: Int,
)
/**
* Open the speaker's `catalog.json` track in the background, parse
* the first frame, and stash it in [speakerCatalogs]. Best-effort —
@@ -961,7 +1171,21 @@ class NestViewModel(
) {
if (closed || activeSubscriptions[pubkey] !== slot) return
try {
val handle = l.subscribeSpeaker(pubkey)
// [maxLatencyMs] tells the relay our staleness tolerance:
// groups older than this should be evicted, not queued.
// With `0L` (the wire default the JS reference uses) the
// relay's per-subscriber forward queue grows unbounded
// until it cliffs — see
// `nestsClient/plans/2026-05-01-quic-stream-cliff-investigation.md`.
// 500 ms is long enough to ride out a typical relay-side
// hiccup without pruning a healthy group, short enough
// that a wedged subscriber drains its buffer before the
// queue overflows. Combined with the cliff-detector below
// that triggers `recycleSession()` when frames stop
// arriving despite the broadcast still being announced,
// this is a defense-in-depth pair against the moq-rs
// forward-queue starvation behaviour.
val handle = l.subscribeSpeaker(pubkey, maxLatencyMs = ROOM_AUDIO_MAX_LATENCY_MS)
// Re-check after the suspending subscribeSpeaker — the user
// may have removed this speaker via updateSpeakers / disconnected
// while the SUBSCRIBE was in flight. If so, abandon the handle
@@ -971,15 +1195,41 @@ class NestViewModel(
viewModelScope.launch { runCatching { handle.unsubscribe() } }
return
}
// Start the catalog fetch BEFORE we wait for it — runs in
// parallel with the main subscribe so the round-trip overlaps
// with the audio-side handshake. Cancelled by closeSubscription.
fetchSpeakerCatalog(l, pubkey)
// Wait briefly for the catalog so we can configure the
// decoder + AudioTrack to match the publisher's actual
// channel layout. With the speaker-side emit-on-subscribe
// hook the catalog frame is on the wire within one round-trip
// of the SUBSCRIBE_OK, so this typically returns within tens
// of ms. Falls back to mono if no catalog arrives within
// [CATALOG_AWAIT_TIMEOUT_MS] — covers legacy publishers and
// the failure-mode where the relay never forwards the catalog
// group; audio still plays, just in the default config.
val audioCfg = awaitAudioPipelineConfig(pubkey, CATALOG_AWAIT_TIMEOUT_MS)
// Allocate native resources (MediaCodec decoder + AudioTrack
// player on Android). Both are heavy and leaky if dropped on
// the floor — wrap them in a nested try so any cancellation
// or throw between here and slot.attach releases them
// (audit round-2 VM #7).
val decoder = decoderFactory()
// Per-subscription decoder factory closure: captures the
// catalog-derived [audioCfg] so a publisher-boundary
// decoder rebuild (see [NestPlayer]'s `decoderFactory`
// kdoc) reuses the SAME channel layout AND sample rate —
// without it, a rebuild after a cliff-recycle would
// default to mono / 48 kHz and a stereo or non-48 kHz
// publisher would silently downmix / clock-mismatch on
// the new decoder. The factory is also called for the
// initial decoder so the construction path is uniform.
val perSubscriptionDecoderFactory: () -> OpusDecoder = {
decoderFactory(audioCfg.channelCount, audioCfg.sampleRate)
}
val decoder = perSubscriptionDecoderFactory()
val player =
try {
playerFactory()
playerFactory(audioCfg.channelCount, audioCfg.sampleRate)
} catch (t: Throwable) {
runCatching { decoder.release() }
throw t
@@ -994,7 +1244,7 @@ class NestViewModel(
val isHushed = pubkey in _uiState.value.locallyHushed
val roomPlayer =
NestPlayer(
decoder = decoder,
initialDecoder = decoder,
player = player,
scope = viewModelScope,
// ~100 ms of audio buffered before the AudioTrack
@@ -1004,6 +1254,14 @@ class NestViewModel(
// tuned in the audio-rooms audit; see NestPlayer
// kdoc for details.
prerollFrames = ROOM_PLAYER_PREROLL_FRAMES,
// Trigger a decoder rebuild on every publisher
// boundary (re-issuing wrapper spliced in a new
// SUBSCRIBE → trackAlias changes). Without this,
// Opus's predictor state from the prior
// publisher's last frame is fed into the new
// publisher's first frame and produces audible
// warble at every cliff-recycle / hot-swap.
decoderFactory = perSubscriptionDecoderFactory,
)
// Apply current mute + per-speaker hush state before play()
// opens the device so the first frame respects them.
@@ -1014,6 +1272,19 @@ class NestViewModel(
// Tap the object flow to drive the speaking-now indicator before
// the decoder consumes it.
val instrumented = handle.objects.onEach { onSpeakerActivity(pubkey) }
// Per-subscription consecutive-decoder-error counter. Single-frame
// Opus errors are normal noise (PLC handles them), but a structural
// mismatch — wrong wire format, missing legacy-container varint
// strip, codec mismatch — fails *every* frame in a row. The
// previous behaviour was to silently swallow per-packet
// [AudioException.Kind.DecoderError] indefinitely, which made
// exactly that class of bug invisible (no audio, no log, no UI
// signal). Track the streak and surface it once per crossing so
// the next protocol regression shows up in `adb logcat | grep
// NestRx` rather than in a user bug report.
val consecutiveDecodeErrors =
java.util.concurrent.atomic
.AtomicLong(0L)
roomPlayer.play(
instrumented,
onError = { err ->
@@ -1030,7 +1301,28 @@ class NestViewModel(
AudioException.Kind.DecoderError,
AudioException.Kind.EncoderError,
-> {
Unit
val streak = consecutiveDecodeErrors.incrementAndGet()
// 50 frames = 1 s of solid decode failures.
// Anything past that is a structural bug, not
// jitter. Log once per crossing (use exact
// equality, not modulus) so we get a single
// attention-grabbing line per stuck stream
// rather than a 50 Hz flood.
if (streak == DECODE_ERROR_STREAK_LOG_THRESHOLD) {
// Non-lambda overload here (rather than the
// allocation-free lambda one) so we preserve
// the throwable stacktrace — this fires
// exactly once per stuck-stream streak, so
// the unconditional string build is fine.
Log.w(
"NestRx",
"decoder failed $DECODE_ERROR_STREAK_LOG_THRESHOLD consecutive frames " +
"for pubkey='${pubkey.take(8)}' — likely wire-format mismatch " +
"(missing legacy-container varint strip, codec/channel-count mismatch, " +
"or corrupted Opus). last error: ${err.message}",
err.cause,
)
}
}
AudioException.Kind.PlaybackFailed,
@@ -1042,7 +1334,13 @@ class NestViewModel(
}
}
},
onLevel = { onAudioLevel(pubkey, it) },
onLevel = { lvl ->
// A successful decode resets the error streak.
// Sample-rate-fast (50 Hz) but cheap on the
// success path: Atomic.set with no allocation.
consecutiveDecodeErrors.set(0L)
onAudioLevel(pubkey, lvl)
},
)
slot.attach(handle, roomPlayer, player)
publishActiveSpeakers()
@@ -1052,12 +1350,10 @@ class NestViewModel(
_uiState.update {
it.copy(connectingSpeakers = (it.connectingSpeakers + pubkey).toPersistentSet())
}
// Parallel catalog fetch — best-effort, doesn't gate
// audio playback. Tracked in catalogJobs; cancelled
// by closeSubscription so a removed speaker doesn't
// leave the catalog collector running on the
// wrapper's still-live re-issuing handle.
fetchSpeakerCatalog(l, pubkey)
// Catalog fetch was started earlier (before the decoder
// wait) so it could overlap with the audio-side handshake;
// its job is already in [catalogJobs] and gets cancelled
// by [closeSubscription] alongside the audio path.
} catch (t: Throwable) {
// Either CancellationException (scope cancelled mid-construction)
// or an unexpected throw — release the half-built pipeline and
@@ -1089,6 +1385,10 @@ class NestViewModel(
announcesJob = null
levelEmitterJob?.cancel()
levelEmitterJob = null
cliffDetectorJob?.cancel()
cliffDetectorJob = null
lastFrameAt.clear()
lastCliffRecycleAt = null
// Audio-focus observation only ends on the final teardown —
// a transient disconnect+reconnect (user retry, room swap)
// keeps the bus subscription alive so a focus loss that
@@ -1175,6 +1475,14 @@ class NestViewModel(
*/
private fun onSpeakerActivity(pubkey: String) {
if (closed) return
// Cliff-detector heartbeat: each MoQ object proves the
// listener-side relay-forward queue is still flowing for this
// speaker. The detector only acts when an announced speaker's
// last-frame timestamp ages past `ROOM_AUDIO_CLIFF_TIMEOUT_MS`,
// so an active stream resets it on every frame.
lastFrameAt[pubkey] =
kotlin.time.TimeSource.Monotonic
.markNow()
speakingExpiryJobs[pubkey]?.cancel()
// First frame for this subscription — clear the buffering
// overlay. Subsequent frames are no-ops here.
@@ -1191,6 +1499,122 @@ class NestViewModel(
}
}
/**
* Periodic scan that detects the relay-side forward-queue cliff:
* an announced speaker whose subscription is open but hasn't
* delivered any object in [ROOM_AUDIO_CLIFF_TIMEOUT_MS]. Triggers
* `recycleSession()` on the listener so the wrapper opens a fresh
* QUIC transport — the new session has a fresh per-subscriber
* queue on the relay's side.
*
* Idempotent on start. Stops itself when the listener is gone
* (explicit teardown signal) — no-ops on transient empty
* `activeSubscriptions` so an audience-only listener doesn't
* needlessly recycle.
*/
private fun startCliffDetector() {
if (cliffDetectorJob?.isActive == true) return
com.vitorpamplona.quartz.utils.Log
.d("NestRx") { "cliff-detector launched" }
cliffDetectorJob =
viewModelScope.launch {
var ticks = 0L
try {
while (true) {
delay(ROOM_AUDIO_CLIFF_CHECK_INTERVAL_MS)
ticks += 1
if (closed) return@launch
val l = listener
val connState = _uiState.value.connection
if (l == null || connState !is ConnectionUiState.Connected) {
// Periodic visibility into "why isn't the cliff
// detector acting" — the receiver-side cliff
// we're chasing has been silent in production
// logs even when the wire conditions look right,
// so dump our gating state every 5 s.
if (ticks % CLIFF_DIAG_LOG_EVERY == 0L) {
com.vitorpamplona.quartz.utils.Log.d("NestRx") {
"cliff-detector tick=$ticks gated: listener=${l != null} conn=${connState::class.simpleName}"
}
}
continue
}
val activeSpeakers =
activeSubscriptions
.asSequence()
.filter { (_, slot) -> slot.isPlaying }
.map { it.key }
.toSet()
val announced = _announcedSpeakers.value
if (ticks % CLIFF_DIAG_LOG_EVERY == 0L) {
val ages =
lastFrameAt.entries.joinToString { (k, mark) ->
"${k.take(8)}=${mark.elapsedNow().inWholeMilliseconds}ms"
}
com.vitorpamplona.quartz.utils.Log.d("NestRx") {
"cliff-detector tick=$ticks active=${activeSpeakers.size} announced=${announced.size} lastFrameAges=[$ages]"
}
com.vitorpamplona.nestsclient.trace.NestsTrace.emit("cliff_tick") {
"\"tick\":$ticks," +
"\"active\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(activeSpeakers)}," +
"\"announced\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(announced)}," +
"\"last_frame_age_ms\":{" +
lastFrameAt.entries.joinToString {
"${com.vitorpamplona.nestsclient.trace.jsonStr(it.key)}:${it.value.elapsedNow().inWholeMilliseconds}"
} + "}"
}
}
val stalled =
computeStalledSpeakers(
activeSpeakers = activeSpeakers,
announcedSpeakers = announced,
lastFrameAt = lastFrameAt,
lastRecycleAt = lastCliffRecycleAt,
cliffTimeoutMs = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
cooldownMs = ROOM_AUDIO_CLIFF_COOLDOWN_MS,
)
if (stalled.isEmpty()) continue
com.vitorpamplona.quartz.utils.Log.w("NestRx") {
"cliff-detector: announced+subscribed but silent for ≥${ROOM_AUDIO_CLIFF_TIMEOUT_MS}ms — recycling session. stalled=$stalled"
}
com.vitorpamplona.nestsclient.trace.NestsTrace.emit("cliff_recycle") {
"\"timeout_ms\":$ROOM_AUDIO_CLIFF_TIMEOUT_MS," +
"\"stalled\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(stalled)}"
}
val recycleMark =
kotlin.time.TimeSource.Monotonic
.markNow()
lastCliffRecycleAt = recycleMark
// Reset `lastFrameAt` for every stalled pubkey so
// the cliff timer starts counting from the recycle
// moment, not from the old (pre-recycle) last
// frame timestamp. Without this reset the next
// tick after the cooldown immediately re-trips —
// because `lastFrameAt[pubkey].elapsedNow()` is
// still the giant pre-recycle value plus the
// cooldown window. Production logs at commit
// ea08c43 showed exactly this: 4 recycles in
// ~30 s eventually drove the relay into a
// "subscribe stream FIN before reply" loop where
// it refused fresh subscribes entirely. Using
// the recycle moment as a synthetic frame event
// gives the new session [ROOM_AUDIO_CLIFF_TIMEOUT_MS]
// wall-clock to deliver before the next cliff
// check, matching the expectation a freshly-
// attached subscription should clear inside
// that window if the relay is healthy.
for (pubkey in stalled) {
lastFrameAt[pubkey] = recycleMark
}
runCatching { l.recycleSession() }
}
} finally {
com.vitorpamplona.quartz.utils.Log
.w("NestRx") { "cliff-detector EXITED after $ticks ticks (closed=$closed)" }
}
}
}
private fun clearSpeaking(pubkey: String) {
speakingExpiryJobs.remove(pubkey)
if (_uiState.value.speakingNow.contains(pubkey)) {
@@ -1247,48 +1671,8 @@ class NestViewModel(
}
}
private class ActiveSubscription private constructor(
val pubkey: String,
) {
private var handle: SubscribeHandle? = null
private var roomPlayer: NestPlayer? = null
var player: AudioPlayer? = null
private set
var isPlaying: Boolean = false
private set
fun attach(
handle: SubscribeHandle,
roomPlayer: NestPlayer,
player: AudioPlayer,
) {
this.handle = handle
this.roomPlayer = roomPlayer
this.player = player
this.isPlaying = true
}
/**
* Hand the player + handle back to the caller's coroutine scope —
* `NestPlayer.stop()` and `SubscribeHandle.unsubscribe()` are
* both suspend, and the caller has the right scope to await them
* (so native MediaCodec/AudioTrack release runs after the decode
* loop has unwound, per audit MoQ #11/#12).
*/
fun detach(): Pair<NestPlayer?, SubscribeHandle?> {
isPlaying = false
val p = roomPlayer
val h = handle
roomPlayer = null
handle = null
player = null
return p to h
}
companion object {
fun pending(pubkey: String) = ActiveSubscription(pubkey)
}
}
// ActiveSubscription was extracted to its own file in the same
// package — see [ActiveSubscription.kt].
// Platform-specific Factory lives in `amethyst/.../nests/room/`,
// not commonMain — the lifecycle KMP `ViewModelProvider.Factory`
@@ -1421,16 +1805,72 @@ sealed class BroadcastUiState {
*/
const val SPEAKING_TIMEOUT_MS: Long = 250L
/**
* How long [NestViewModel.openSubscription] waits for the publisher's
* `catalog.json` to land before constructing the decoder + AudioTrack.
* The catalog declares the audio rendition's `numberOfChannels` and
* `sampleRate`, which the decoder needs at construction time so a
* stereo or non-48 kHz web publisher doesn't get its frames decoded
* with the wrong layout / clock.
*
* Sized to be generous enough that a typical publisher's catalog
* (which arrives within a single round-trip of the SUBSCRIBE_OK with
* the speaker-side emit-on-subscribe hook in place) lands well before
* the timeout, and short enough that a publisher that never emits a
* catalog (legacy publishers, relay-side bug) doesn't visibly stall
* the listener — audio playback proceeds in the default config.
*
* **Frame-buffer budget**: while we wait, audio frames stream from
* the relay into the wrapper's
* [com.vitorpamplona.nestsclient.ReconnectingNestsListener]
* `SUBSCRIBE_BUFFER = 64` SharedFlow with `DROP_OLDEST` overflow.
* At 50 fps Opus that's 1.28 s of headroom; the 500 ms wait
* consumes at most 25 frames of buffer, leaving ≥ 39 frames of
* margin (≈780 ms) before the oldest frame would be dropped. Even
* at the production `framesPerGroup = 50` (= 1 group / sec) cadence
* this never trips the eviction path during normal startup.
*/
const val CATALOG_AWAIT_TIMEOUT_MS: Long = 500L
/**
* Per-subscription consecutive Opus decode-error count that triggers a
* single conspicuous warning log. Single-frame decode errors are normal
* — Opus PLC papers over them — and are silently swallowed; a streak of
* this size is structural (wire-format mismatch, missing
* legacy-container varint strip, codec/channel-count mismatch, corrupt
* publisher) and demands attention because the user-visible symptom is
* "speaker tile shows up but no audio plays" with no other signal.
*
* 50 × 20 ms ≈ 1 s. Long enough that a transient flurry of bad frames
* doesn't cry wolf; short enough that an "every frame fails" structural
* bug surfaces in a logcat line within the first second of audio.
*
* Logged once per crossing (exact equality, not modulus) so a stuck
* stream produces ONE attention-grabbing line rather than 50 Hz of
* noise.
*/
const val DECODE_ERROR_STREAK_LOG_THRESHOLD: Long = 50L
/**
* Per-speaker pre-roll: number of decoded PCM frames buffered before the
* underlying [com.vitorpamplona.nestsclient.audio.AudioPlayer] starts
* consuming. 5 × 20 ms ≈ 100 ms of audio — long enough to mask a typical
* Main-thread stall (Compose recomposition / GC) without adding perceptible
* join latency. Combines with [com.vitorpamplona.nestsclient.audio.AudioTrackPlayer]'s
* ~250 ms AudioTrack buffer for ~350 ms of total slack between the
* consuming. 10 × 20 ms ≈ 200 ms of audio.
*
* Sized to absorb the slowest-cadence publisher we expect: kixelated/moq's
* web reference encoder uses `groupDuration: 100ms` (5 frames per group)
* by default but adds another 100200 ms of jitter on the relay's
* outbound forward queue under residential network conditions. 100 ms of
* preroll (the previous value) reliably underran on web speakers — the
* AudioTrack ran out of samples between adjacent groups and the user
* heard short clicks. 200 ms covers the observed jitter envelope while
* keeping perceived join latency well under the typical "speaker becomes
* audible" delay on the wire.
*
* Combines with [com.vitorpamplona.nestsclient.audio.AudioTrackPlayer]'s
* ~250 ms AudioTrack buffer for ~450 ms of total slack between the
* arrival-of-frame and the underrun horizon.
*/
const val ROOM_PLAYER_PREROLL_FRAMES: Int = 5
const val ROOM_PLAYER_PREROLL_FRAMES: Int = 10
/**
* Coalescing interval for [NestViewModel.audioLevels]. The decode loop
@@ -1442,6 +1882,125 @@ const val ROOM_PLAYER_PREROLL_FRAMES: Int = 5
*/
const val LEVEL_TICK_MS: Long = 100L
/**
* `max_latency` we send on every speaker SUBSCRIBE. Tells the relay
* to evict groups older than this from the per-subscriber forward
* queue rather than holding them up to the relay's `MAX_GROUP_AGE`
* default (30 s). Defends against the moq-rs forward-queue cliff
* documented in
* `nestsClient/plans/2026-05-01-quic-stream-cliff-investigation.md`.
* 500 ms = 25 Opus frames at 20 ms cadence — long enough to ride
* out a typical hiccup without dropping a healthy frame, short
* enough to keep the relay's per-subscriber queue bounded.
*/
const val ROOM_AUDIO_MAX_LATENCY_MS: Long = 500L
/**
* Inactivity threshold for the listener-side cliff detector. If
* a speaker is `_announcedSpeakers` (= relay says they're broadcasting)
* AND the wrapper's MoqObject flow doesn't deliver a frame for this
* many milliseconds, the listener forces a `recycleSession()` so the
* orchestrator opens a fresh QUIC transport. Resets the relay's
* per-subscriber forward queue.
*
* 2 500 ms is past every legitimate group boundary (`framesPerGroup =
* 50` packs 1 s of audio per uni stream, so the receiver sees one
* `drainOneGroup` FIN every ~1 s on a healthy stream — 2.5 s of
* silence means more than two group cycles missed, very likely a
* cliff). Earlier than 2 s would false-positive on the natural 1 s
* cadence; later than 3 s wastes audible audio at every cliff edge.
*
* Production logs (`6e4df4a` run 18:37:43..18:38:08) showed the
* cliff-after-streaming case losing ~6 s of audio at the 4 s
* detector + ~2 s reconnect handshake. Tightening detection to
* 2.5 s shaves ~1.5 s off the audible gap.
*/
const val ROOM_AUDIO_CLIFF_TIMEOUT_MS: Long = 2_500L
/**
* How often the cliff detector wakes up to scan
* `activeSubscriptions` against `lastFrameAt`. 1 s is a coarse
* enough cadence not to cost anything (one map walk per second)
* while keeping detection latency well under the 4 s threshold.
*/
const val ROOM_AUDIO_CLIFF_CHECK_INTERVAL_MS: Long = 1_000L
/**
* Cooldown after a cliff-detector-driven recycle. Suppresses
* back-to-back recycles while the wrapper is mid-handshake on the
* fresh session AND while the relay is recovering its
* per-subscriber forward queue.
*
* Production logs at commit ea08c43 showed an 8 s cooldown was
* not enough — moq-rs needed longer to recover between aggressive
* recycles, and 4 recycles in ~30 s drove the relay into a
* "subscribe stream FIN before reply" loop that refused all
* subsequent subscribes. 30 s gives the relay's per-subscriber
* forward task pool time to drain stale pending writes from the
* previous subscription before the new subscribe lands. The
* audible-gap cost rises from ~5 s to ~30 s in the worst case
* (a fully-stalled relay), but the trade is correct: 30 s of
* silence + recovered audio beats endless silence with the relay
* locked out.
*/
const val ROOM_AUDIO_CLIFF_COOLDOWN_MS: Long = 30_000L
/**
* Diagnostic log frequency for the cliff detector. Emit a state-dump
* line every Nth 1-second tick — i.e. every 5 s while connected and
* idle, so logcat captures the "active / announced / lastFrameAges"
* snapshot we use to debug silent-cliff cases without flooding when
* everything is fine.
*/
private const val CLIFF_DIAG_LOG_EVERY: Long = 5L
/**
* Pure logic of the cliff detector — extracted so headless tests can
* exercise it with a [kotlin.time.TestTimeSource] without standing up
* the full [NestViewModel] coroutine machinery.
*
* Returns the list of pubkeys that should trigger a `recycleSession()`
* RIGHT NOW. Empty list means "no cliff observed; do nothing on this
* tick."
*
* Recycle conditions, all-of:
* - the pubkey has an active subscription on our side
* ([activeSpeakers]) — i.e. we're actively pulling its audio,
* - the relay has told us this pubkey is broadcasting
* ([announcedSpeakers]),
* - we've received at least one MoQ object in the past
* ([lastFrameAt] has an entry — we don't preemptively recycle
* a brand-new subscription that simply hasn't ramped up yet),
* - the elapsed time since that last object is at or past
* [cliffTimeoutMs].
*
* Suppression:
* - if [lastRecycleAt] is non-null and less than [cooldownMs] has
* elapsed since it, returns empty (no cascading recycles while
* the wrapper is mid-handshake on a fresh session).
*/
internal fun computeStalledSpeakers(
activeSpeakers: Set<String>,
announcedSpeakers: Set<String>,
lastFrameAt: Map<String, kotlin.time.TimeMark>,
lastRecycleAt: kotlin.time.TimeMark?,
cliffTimeoutMs: Long = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
cooldownMs: Long = ROOM_AUDIO_CLIFF_COOLDOWN_MS,
): List<String> {
if (lastRecycleAt != null &&
lastRecycleAt.elapsedNow().inWholeMilliseconds < cooldownMs
) {
return emptyList()
}
return activeSpeakers
.asSequence()
.filter { it in announcedSpeakers }
.filter { pubkey ->
val mark = lastFrameAt[pubkey] ?: return@filter false
mark.elapsedNow().inWholeMilliseconds >= cliffTimeoutMs
}.toList()
}
/**
* How long a kind-7 reaction stays in
* [NestViewModel.recentReactions] before the eviction sweep
@@ -22,35 +22,67 @@ package com.vitorpamplona.amethyst.commons.viewmodels
import androidx.compose.runtime.Immutable
import com.vitorpamplona.quartz.nip01Core.core.JsonMapper
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
/**
* Decoded shape of a moq-lite `catalog.json` track for a single
* audio-room speaker. Each broadcast advertises one or more tracks
* (one per codec / quality tier) — for nests audio rooms there's
* exactly one Opus track today, but the model carries a list for
* forward compatibility.
* audio-room speaker. Mirrors the kixelated/moq `hang` reference
* catalog format — the canonical shape that `@kixelated/hang`'s
* browser watcher and the `moq-rs` Rust hang crate both produce and
* consume — so Amethyst publishers are visible to the moq-lite browser
* reference, and Amethyst listeners parse standards-aligned catalogs
* from non-Amethyst publishers.
*
* Best-effort schema: nostrnests' upstream moq-lite catalog spec
* (kixelated/moq) has evolved across revisions, and not every
* publisher fills in every field. The parser tolerates missing
* keys via `ignoreUnknownKeys = true` on [JsonMapper] and the
* `?`-marked properties.
* Shape (verbatim from `kixelated/moq/rs/hang/src/catalog/`):
*
* {
* "audio": {
* "renditions": {
* "<trackName>": {
* "codec": "opus",
* "container": { "kind": "legacy" },
* "sampleRate": 48000,
* "numberOfChannels": 1,
* "bitrate": 32000 // optional
* }
* }
* }
* }
*
* Keys are camelCase per the upstream serde `rename_all = "camelCase"`.
* Every field except the rendition-key string is optional in the parser
* — older / partial publishers are tolerated. Field semantics:
*
* - rendition key: the moq-lite `track` name a subscriber should
* subscribe to for that audio rendition's frames (commonly the same
* string for single-rendition broadcasts). Subscribers pick a
* rendition (e.g. by codec / bitrate) and use this key as the
* Subscribe.track string.
* - `codec`: codec mimetype string (`"opus"`, `"mp4a.40.2"` for AAC).
* - `container.kind`: frame wrapper. `"legacy"` = each frame is
* `varint(timestamp_us) + raw_codec_payload` inside the moq-lite
* group. `"cmaf"` = MOOF/MDAT-fragmented MP4. We only emit/parse
* `legacy` today; unknown kinds are ignored at parse time.
* - `sampleRate` / `numberOfChannels`: PCM source params.
*/
@Immutable
@Serializable
data class RoomSpeakerCatalog(
val version: Int? = null,
val audio: List<AudioTrack> = emptyList(),
val audio: Audio? = null,
) {
@Immutable
@Serializable
data class AudioTrack(
val track: String? = null,
data class Audio(
val renditions: Map<String, AudioConfig> = emptyMap(),
)
@Immutable
@Serializable
data class AudioConfig(
val codec: String? = null,
@SerialName("sample_rate") val sampleRate: Int? = null,
@SerialName("channel_count") val channelCount: Int? = null,
val container: Container? = null,
val sampleRate: Int? = null,
val numberOfChannels: Int? = null,
val bitrate: Int? = null,
) {
/**
@@ -64,14 +96,48 @@ data class RoomSpeakerCatalog(
buildList {
codec?.takeIf { it.isNotBlank() }?.let { add(it.uppercase()) }
sampleRate?.let { add("${it / 1000}kHz") }
channelCount?.let { add(if (it == 1) "mono" else "${it}ch") }
numberOfChannels?.let { add(if (it == 1) "mono" else "${it}ch") }
}
return parts.takeIf { it.isNotEmpty() }?.joinToString(" · ")
}
}
/** First audio track, if any. The current single-Opus reality. */
fun primaryAudio(): AudioTrack? = audio.firstOrNull()
@Immutable
@Serializable
data class Container(
val kind: String? = null,
) {
companion object {
/**
* The only `container.kind` value the Amethyst decoder
* pipeline currently understands: each moq-lite frame is
* `varint(timestamp_us) + raw_codec_payload`. Source:
* `kixelated/moq/rs/hang/src/container/legacy.rs`.
*
* Other documented kinds — `"cmaf"` (MOOF/MDAT-fragmented
* MP4) and a handful of in-flight experimental shapes —
* would require a different decoder path and are
* intentionally rejected by [primaryAudio].
*/
const val LEGACY_KIND: String = "legacy"
}
}
/**
* First audio rendition with a [Container.LEGACY_KIND] container,
* which is the only container layout the Amethyst decoder pipeline
* understands today (`varint(timestamp_us) + opus_packet` per
* frame). Returns null when no legacy rendition exists — the UI
* surfaces "no audio info" rather than a CMAF rendition we'd try
* to feed to our legacy decoder.
*
* Picks the first match in JSON iteration order
* (kotlinx.serialization uses LinkedHashMap), so the publisher's
* preferred legacy rendition wins. A future publisher that emits
* CMAF-first then legacy still surfaces the legacy entry; CMAF-
* only publishers correctly return null.
*/
fun primaryAudio(): AudioConfig? = audio?.renditions?.values?.firstOrNull { it.container?.kind == Container.LEGACY_KIND }
companion object {
/**
@@ -0,0 +1,313 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.viewmodels
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.ExperimentalTime
import kotlin.time.TestTimeSource
/**
* Unit tests for the [computeStalledSpeakers] predicate that drives the
* listener-side cliff detector in [NestViewModel]. The predicate is
* extracted as a pure function specifically so we can drive it with a
* [TestTimeSource] — neither `runTest`'s virtual scheduler nor a real
* monotonic clock would let us deterministically place events relative
* to the 4 s `ROOM_AUDIO_CLIFF_TIMEOUT_MS` threshold.
*
* The detector defends against the moq-rs production-relay forward-
* queue cliff documented in
* `nestsClient/plans/2026-05-01-quic-stream-cliff-investigation.md`:
* the relay opens fewer (or zero) new uni streams to a slow
* subscriber while the announce stream still says the broadcast is
* Active and the subscribe bidi is still alive at QUIC level — no
* other layer notices. These tests pin the recycle gate so a future
* refactor doesn't accidentally widen or narrow when we recycle.
*/
@OptIn(ExperimentalTime::class)
class CliffDetectorTest {
@Test
fun returnsEmptyWhenNoSpeakersActive() {
val ts = TestTimeSource()
val result =
computeStalledSpeakers(
activeSpeakers = emptySet(),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to ts.markNow()),
lastRecycleAt = null,
)
assertTrue(result.isEmpty())
}
@Test
fun ignoresSpeakerNotInAnnounced() {
// A speaker we're subscribed to but who isn't in the relay's
// announce-stream output is most likely either (a) already
// off stage and the relay just hasn't propagated the Ended
// yet, or (b) in the brief gap between an Ended and the
// wrapper re-issuing. Recycling the whole transport here
// would be wasteful — we only act on confirmed cliffs.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 10_000.milliseconds // way past threshold
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = emptySet(),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertTrue(result.isEmpty())
}
@Test
fun ignoresSpeakerWithNoFrameYet() {
// Brand-new subscription that hasn't ramped up: lastFrameAt
// has no entry. Don't recycle — the wrapper's per-speaker
// re-issue + opener-throws backoff already retries on its
// own; the cliff detector only acts once we've proven the
// subscription was working.
val ts = TestTimeSource()
ts += 10_000.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = emptyMap(),
lastRecycleAt = null,
)
assertTrue(result.isEmpty())
}
@Test
fun returnsEmptyJustBeforeThreshold() {
// Boundary case: 1 ms under the 2.5 s threshold. The detector
// should still consider this healthy. Relevant because audio
// groups arrive at ~1 s cadence (framesPerGroup=50), so a
// false positive at ≈ threshold would recycle on every
// group-rollover hiccup.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 2_499.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertTrue(result.isEmpty())
}
@Test
fun returnsStalledSpeakerAtThresholdInclusive() {
// Exactly at threshold: include. The check is `>=`, not `>`.
// Tested explicitly because boundary semantics here are
// user-visible (this is the "audio went silent" detection).
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 2_500.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertEquals(listOf(ALICE), result)
}
@Test
fun returnsStalledSpeakerWellPastThreshold() {
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 30_000.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertEquals(listOf(ALICE), result)
}
@Test
fun returnsAllStalledSpeakersAtOnceMixedWithFreshOnes() {
// Multi-speaker stage: alice and bob have stalled, charlie's
// subscription is fresh (just received a frame). Recycle
// should fire for the stalled set and leave charlie alone —
// the `recycleSession` itself takes the whole listener down,
// but the test here is the predicate's classification.
// Threshold = 2.5 s default. Spacings: alice/bob 4 s old (past),
// charlie 1.5 s old (under).
val ts = TestTimeSource()
val aliceFrame = ts.markNow()
val bobFrame = ts.markNow()
ts += 2_500.milliseconds
val charlieFrame = ts.markNow()
ts += 1_500.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE, BOB, CHARLIE),
announcedSpeakers = setOf(ALICE, BOB, CHARLIE),
lastFrameAt =
mapOf(
ALICE to aliceFrame,
BOB to bobFrame,
CHARLIE to charlieFrame,
),
lastRecycleAt = null,
)
assertEquals(setOf(ALICE, BOB), result.toSet())
}
@Test
fun cooldownSuppressesRecycleEvenWhenStalled() {
// After a recycle, the wrapper opens a fresh QUIC transport.
// The new session has no `lastFrameAt` entries yet for any
// pubkey; without a cooldown we would re-trigger on the
// very next 1 s tick because the prior tick's stalled
// pubkeys still age past the threshold (their lastFrameAt
// hasn't been updated by the new session yet). 30 s cooldown
// covers the typical reconnect handshake AND gives moq-rs
// time to drain its per-subscriber forward queue from the
// prior subscription before the new subscribe lands.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 4_500.milliseconds // past threshold
val recycleMark = ts.markNow() // recycle just fired
ts += 5_000.milliseconds // 5 s into cooldown — well within 30 s window
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = recycleMark,
)
assertTrue(result.isEmpty())
}
@Test
fun cooldownReleasesAfterTimeoutPasses() {
// Once cooldown elapses, a still-stalled subscription
// becomes eligible to recycle again. Important for the
// case where the recycle didn't actually fix the cliff —
// we want a second attempt rather than getting wedged.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 4_500.milliseconds
val recycleMark = ts.markNow()
ts += 30_001.milliseconds // 1 ms past 30 s cooldown
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = recycleMark,
)
assertEquals(listOf(ALICE), result)
}
@Test
fun customTimeoutsAreHonored() {
// Defaults are wired into the production VM, but the function
// accepts overrides so a test for tighter / looser tolerances
// can drive it without mutating the constants.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 1_000.milliseconds
val withDefault =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertTrue(withDefault.isEmpty(), "1 s elapsed shouldn't trip 2.5 s default")
val withTightTimeout =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
cliffTimeoutMs = 500L,
)
assertEquals(listOf(ALICE), withTightTimeout, "1 s elapsed should trip a 500 ms threshold")
}
@Test
fun activeSpeakerNotInLastFrameAtIsIgnoredEvenWithStalledPeers() {
// Mixed: alice is announced + active + stalled, bob is
// announced + active but has no frame yet. Detector should
// return alice only — bob hasn't proven the subscription
// works, so it's not "stalled" yet, just slow to start.
val ts = TestTimeSource()
val aliceFrame = ts.markNow()
ts += 5_000.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE, BOB),
announcedSpeakers = setOf(ALICE, BOB),
lastFrameAt = mapOf(ALICE to aliceFrame),
lastRecycleAt = null,
)
assertEquals(listOf(ALICE), result)
}
@Test
fun nullLastRecycleNeverSuppresses() {
// Null `lastRecycleAt` means "we've never recycled" — the
// very first cliff event of a session must fire. The
// cooldown check is gated on non-null.
val ts = TestTimeSource()
val frameMark = ts.markNow()
ts += 5_000.milliseconds
val result =
computeStalledSpeakers(
activeSpeakers = setOf(ALICE),
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = null,
)
assertEquals(listOf(ALICE), result)
}
companion object {
private val ALICE = "a".repeat(64)
private val BOB = "b".repeat(64)
private val CHARLIE = "c".repeat(64)
}
}
@@ -455,8 +455,8 @@ class NestViewModelTest {
NestViewModel(
httpClient = NoopNestsClient,
transport = NoopWebTransportFactory,
decoderFactory = { NoopOpusDecoder },
playerFactory = { NoopAudioPlayer() },
decoderFactory = { _, _ -> NoopOpusDecoder },
playerFactory = { _, _ -> NoopAudioPlayer() },
signer = NoopSigner,
room = ROOM_CONFIG,
connector =
@@ -27,69 +27,80 @@ import kotlin.test.assertNull
class RoomSpeakerCatalogTest {
@Test
fun parsesNostrnestsShape() {
fun parsesKixelatedHangShape() {
// Verbatim shape produced by the canonical kixelated/moq `hang`
// crate (`rs/hang/src/catalog/`). Verifies that an Amethyst
// listener picks up a standards-aligned publisher's catalog.
val json =
"""
{
"version": 1,
"audio": [
{
"track": "audio/data",
"codec": "opus",
"sample_rate": 48000,
"channel_count": 2,
"bitrate": 64000
"audio": {
"renditions": {
"audio/data": {
"codec": "opus",
"container": { "kind": "legacy" },
"sampleRate": 48000,
"numberOfChannels": 2,
"bitrate": 64000
}
}
]
}
}
""".trimIndent()
val catalog = RoomSpeakerCatalog.parseOrNull(json.encodeToByteArray())
assertNotNull(catalog)
assertEquals(1, catalog.version)
assertEquals(1, catalog.audio.size)
val track = catalog.primaryAudio()
assertNotNull(track)
assertEquals("audio/data", track.track)
assertEquals("opus", track.codec)
assertEquals(48_000, track.sampleRate)
assertEquals(2, track.channelCount)
assertEquals(64_000, track.bitrate)
assertEquals(1, catalog.audio?.renditions?.size)
val rendition = catalog.primaryAudio()
assertNotNull(rendition)
assertEquals("opus", rendition.codec)
assertEquals("legacy", rendition.container?.kind)
assertEquals(48_000, rendition.sampleRate)
assertEquals(2, rendition.numberOfChannels)
assertEquals(64_000, rendition.bitrate)
}
@Test
fun describeFormatsHumanReadable() {
val track =
RoomSpeakerCatalog.AudioTrack(
val rendition =
RoomSpeakerCatalog.AudioConfig(
codec = "opus",
sampleRate = 48_000,
channelCount = 2,
numberOfChannels = 2,
)
// Codec uppercased, kHz / channel count short-formed —
// intended for a single-line tooltip in the participant
// sheet, not a verbose codec dump.
assertEquals("OPUS · 48kHz · 2ch", track.describe())
assertEquals("OPUS · 48kHz · 2ch", rendition.describe())
}
@Test
fun describeMonoIsLabelled() {
val track = RoomSpeakerCatalog.AudioTrack(codec = "opus", channelCount = 1)
assertEquals("OPUS · mono", track.describe())
val rendition = RoomSpeakerCatalog.AudioConfig(codec = "opus", numberOfChannels = 1)
assertEquals("OPUS · mono", rendition.describe())
}
@Test
fun describeAllNullReturnsNull() {
val track = RoomSpeakerCatalog.AudioTrack()
val rendition = RoomSpeakerCatalog.AudioConfig()
// Empty catalog → caller doesn't render a "unknown · unknown"
// line; describe() returns null so the UI can omit cleanly.
assertNull(track.describe())
assertNull(rendition.describe())
}
@Test
fun toleratesUnknownKeys() {
// Forward-compat: future moq-lite catalog revisions can add
// fields without breaking older clients.
// Forward-compat: future hang catalog revisions can add
// fields without breaking older clients. Container.kind is
// declared as `legacy` because `primaryAudio()` filters to
// that kind — the unknown-key tolerance we're testing here
// is about unrecognized siblings, not about the container
// requirement itself.
val json =
"""{"version":2,"audio":[{"track":"audio/data","codec":"opus","extra":"future-only"}],"new_top_level":true}"""
"""{"audio":{"renditions":{"audio/data":{
| "codec":"opus","container":{"kind":"legacy"},
| "extra":"future-only"
|}}},"video":{},"newTopLevel":true}
""".trimMargin()
val catalog = RoomSpeakerCatalog.parseOrNull(json.encodeToByteArray())
assertNotNull(catalog)
assertEquals("opus", catalog.primaryAudio()?.codec)
@@ -104,12 +115,103 @@ class RoomSpeakerCatalogTest {
}
@Test
fun emptyAudioListIsAllowed() {
fun emptyAudioRenditionsIsAllowed() {
// A publisher might emit an "advertising" catalog with no
// tracks yet (e.g. between codec switches). Accept the shape
// and let primaryAudio() return null.
val catalog = RoomSpeakerCatalog.parseOrNull("""{"version":1,"audio":[]}""".encodeToByteArray())
// renditions yet (e.g. between codec switches). Accept the
// shape and let primaryAudio() return null.
val catalog = RoomSpeakerCatalog.parseOrNull("""{"audio":{"renditions":{}}}""".encodeToByteArray())
assertNotNull(catalog)
assertNull(catalog.primaryAudio())
}
@Test
fun primaryAudioPicksLegacyEvenWhenCmafComesFirst() {
// A future publisher may emit CMAF-first then legacy in the
// renditions map. We only know how to decode the legacy
// container today (varint(ts)+opus); CMAF (MOOF/MDAT) would
// be fed bytes-as-Opus and decode to garbage. The filter
// skips non-legacy renditions so the chosen entry is one we
// can actually play.
val mixed =
"""{"audio":{"renditions":{
| "video/cmaf":{"codec":"opus","container":{"kind":"cmaf"},"sampleRate":48000,"numberOfChannels":1},
| "audio/data":{"codec":"opus","container":{"kind":"legacy"},"sampleRate":48000,"numberOfChannels":1}
|}}}
""".trimMargin()
val catalog = RoomSpeakerCatalog.parseOrNull(mixed.encodeToByteArray())
assertNotNull(catalog)
val rendition = catalog.primaryAudio()
assertNotNull(rendition, "expected the legacy rendition, not CMAF")
assertEquals("legacy", rendition.container?.kind)
}
@Test
fun primaryAudioReturnsNullForCmafOnlyPublisher() {
// No legacy rendition → no decoder path. Returning null is
// the right contract: the caller falls back to "unknown
// codec / no audio" rather than feeding CMAF bytes to the
// legacy decoder.
val cmafOnly =
"""{"audio":{"renditions":{"audio/cmaf":{
| "codec":"opus","container":{"kind":"cmaf"},
| "sampleRate":48000,"numberOfChannels":1
|}}}}
""".trimMargin()
val catalog = RoomSpeakerCatalog.parseOrNull(cmafOnly.encodeToByteArray())
assertNotNull(catalog)
assertNull(catalog.primaryAudio())
}
@Test
fun primaryAudioReturnsNullWhenContainerKindIsMissing() {
// Defensive: a malformed publisher that omits `container`
// entirely must NOT be treated as legacy by default — we'd
// be guessing the wire shape. Same null fallback as the
// CMAF-only case.
val noContainer =
"""{"audio":{"renditions":{"audio/data":{
| "codec":"opus","sampleRate":48000,"numberOfChannels":1
|}}}}
""".trimMargin()
val catalog = RoomSpeakerCatalog.parseOrNull(noContainer.encodeToByteArray())
assertNotNull(catalog)
assertNull(catalog.primaryAudio())
}
@Test
fun missingAudioReturnsNullPrimary() {
// hang catalogs declaring only video MUST still parse without
// throwing — primaryAudio() returns null.
val catalog = RoomSpeakerCatalog.parseOrNull("""{}""".encodeToByteArray())
assertNotNull(catalog)
assertNull(catalog.primaryAudio())
}
@Test
fun stripPrefixRoundTripsCanonicalCatalog() {
// The catalog payload `MoqLiteHangCatalog.opusMono48k(...)` emits
// (in `:nestsClient`) MUST round-trip through this parser — the
// two classes target the same wire shape independently because
// `:nestsClient` does not depend on `:commons` and vice versa.
// `MoqLiteHangCatalog` isn't reachable from this module; keep an
// inline literal that mirrors what the publisher emits verbatim
// so a desync on either side trips this assertion.
val emitted =
"{\"audio\":{\"renditions\":{\"audio/data\":{" +
"\"codec\":\"opus\",\"container\":{\"kind\":\"legacy\"}," +
"\"sampleRate\":48000,\"numberOfChannels\":1,\"jitter\":20}}}}"
val catalog = RoomSpeakerCatalog.parseOrNull(emitted.encodeToByteArray())
assertNotNull(catalog)
val rendition = catalog.primaryAudio()
assertNotNull(rendition)
assertEquals("opus", rendition.codec)
assertEquals("legacy", rendition.container?.kind)
assertEquals(48_000, rendition.sampleRate)
assertEquals(1, rendition.numberOfChannels)
// `jitter` is intentionally not asserted on `rendition` — the
// commons-side parser drops unknown fields via the JsonMapper's
// `ignoreUnknownKeys = true`. Adding it to `RoomSpeakerCatalog`
// is forward work; the byte-exact match in the literal above
// already pins the wire shape.
}
}