Merge remote-tracking branch 'origin/main' into claude/research-quic-libraries-hH1Dc
This commit is contained in:
+126
-70
@@ -57,6 +57,7 @@ import androidx.compose.ui.Alignment
|
|||||||
import androidx.compose.ui.Modifier
|
import androidx.compose.ui.Modifier
|
||||||
import androidx.compose.ui.draw.alpha
|
import androidx.compose.ui.draw.alpha
|
||||||
import androidx.compose.ui.draw.drawBehind
|
import androidx.compose.ui.draw.drawBehind
|
||||||
|
import androidx.compose.ui.geometry.Offset
|
||||||
import androidx.compose.ui.graphics.Color
|
import androidx.compose.ui.graphics.Color
|
||||||
import androidx.compose.ui.graphics.drawscope.Stroke
|
import androidx.compose.ui.graphics.drawscope.Stroke
|
||||||
import androidx.compose.ui.res.pluralStringResource
|
import androidx.compose.ui.res.pluralStringResource
|
||||||
@@ -420,28 +421,18 @@ private fun MemberCell(
|
|||||||
},
|
},
|
||||||
label = "speaker-outer-ring-width",
|
label = "speaker-outer-ring-width",
|
||||||
)
|
)
|
||||||
|
// Reserve enough space around the avatar Box to fit the outer ring
|
||||||
|
// and glow halo. Without this padding the rings clip against the
|
||||||
|
// surrounding Surface / LazyVerticalGrid bounds (most visibly at
|
||||||
|
// the top edge for the first row, where the stage card's rounded
|
||||||
|
// corner cuts into the glow). The glow extends up to MAX_GLOW_RADIUS
|
||||||
|
// past the avatar; the outer ring extends OUTER_RING_GAP +
|
||||||
|
// OUTER_RING_MAX_WIDTH past it.
|
||||||
|
val ringPadding =
|
||||||
|
maxOf(MAX_GLOW_RADIUS.value, (OUTER_RING_GAP + OUTER_RING_MAX_WIDTH).value).dp
|
||||||
val avatarModifier =
|
val avatarModifier =
|
||||||
Modifier
|
Modifier
|
||||||
.drawBehind {
|
.border(animatedRingWidth, animatedRingColor, CircleShape)
|
||||||
if (animatedGlowAlpha > 0.001f) {
|
|
||||||
val baseRadius = size.minDimension / 2f
|
|
||||||
val extra = MAX_GLOW_RADIUS.toPx() * clampedLevel
|
|
||||||
drawCircle(
|
|
||||||
color = NEST_SPEAKING_COLOR.copy(alpha = animatedGlowAlpha),
|
|
||||||
radius = baseRadius + extra,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
if (animatedOuterRingAlpha > 0.001f && animatedOuterRingWidth > 0.dp) {
|
|
||||||
val baseRadius = size.minDimension / 2f
|
|
||||||
val strokePx = animatedOuterRingWidth.toPx()
|
|
||||||
val ringRadius = baseRadius + OUTER_RING_GAP.toPx() + strokePx / 2f
|
|
||||||
drawCircle(
|
|
||||||
color = NEST_SPEAKING_COLOR.copy(alpha = animatedOuterRingAlpha),
|
|
||||||
radius = ringRadius,
|
|
||||||
style = Stroke(width = strokePx),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}.border(animatedRingWidth, animatedRingColor, CircleShape)
|
|
||||||
.let { if (member.absent) it.alpha(0.5f) else it }
|
.let { if (member.absent) it.alpha(0.5f) else it }
|
||||||
val user =
|
val user =
|
||||||
remember(member.pubkey) {
|
remember(member.pubkey) {
|
||||||
@@ -480,58 +471,53 @@ private fun MemberCell(
|
|||||||
horizontalAlignment = Alignment.CenterHorizontally,
|
horizontalAlignment = Alignment.CenterHorizontally,
|
||||||
modifier = modifier.fillMaxWidth().padding(vertical = 4.dp),
|
modifier = modifier.fillMaxWidth().padding(vertical = 4.dp),
|
||||||
) {
|
) {
|
||||||
Box(contentAlignment = Alignment.Center) {
|
// Outer Box paints the glow halo + detached outer ring on a
|
||||||
ClickableUserPicture(
|
// canvas that's bigger than the avatar by [ringPadding]. The
|
||||||
baseUserHex = member.pubkey,
|
// inner Box keeps its tight-to-avatar bounds so badge corner
|
||||||
size = avatarSize,
|
// alignment (TopStart, TopEnd, BottomCenter, BottomEnd) still
|
||||||
accountViewModel = accountViewModel,
|
// tracks the avatar circle, not the padded outer area.
|
||||||
modifier = avatarModifier,
|
Box(
|
||||||
onClick = onClick,
|
modifier =
|
||||||
onLongClick = onLongClick,
|
Modifier.drawBehind {
|
||||||
)
|
val avatarRadiusPx = avatarSize.toPx() / 2f
|
||||||
if (isConnecting) {
|
val cx = size.width / 2f
|
||||||
CircularProgressIndicator(
|
val cy = size.height / 2f
|
||||||
modifier = Modifier.size(avatarSize - 8.dp),
|
if (animatedGlowAlpha > 0.001f) {
|
||||||
strokeWidth = 2.dp,
|
val extra = MAX_GLOW_RADIUS.toPx() * clampedLevel
|
||||||
color = MaterialTheme.colorScheme.primary,
|
drawCircle(
|
||||||
)
|
color = NEST_SPEAKING_COLOR.copy(alpha = animatedGlowAlpha),
|
||||||
}
|
radius = avatarRadiusPx + extra,
|
||||||
val role = member.role
|
center = Offset(cx, cy),
|
||||||
if (role == ROLE.HOST || role == ROLE.MODERATOR) {
|
)
|
||||||
RoleBadge(
|
}
|
||||||
role = role,
|
if (animatedOuterRingAlpha > 0.001f && animatedOuterRingWidth > 0.dp) {
|
||||||
modifier = Modifier.align(Alignment.TopStart),
|
val strokePx = animatedOuterRingWidth.toPx()
|
||||||
)
|
val ringRadius = avatarRadiusPx + OUTER_RING_GAP.toPx() + strokePx / 2f
|
||||||
}
|
drawCircle(
|
||||||
if (member.handRaised) {
|
color = NEST_SPEAKING_COLOR.copy(alpha = animatedOuterRingAlpha),
|
||||||
HandRaiseBadge(
|
radius = ringRadius,
|
||||||
modifier = Modifier.align(Alignment.TopEnd),
|
center = Offset(cx, cy),
|
||||||
)
|
style = Stroke(width = strokePx),
|
||||||
}
|
)
|
||||||
// Show the mic badge for any on-stage speaker that has
|
}
|
||||||
// an audio state to surface — currently broadcasting
|
},
|
||||||
// (`publishing=1`) OR mic-muted (`muted=1, publishing=0`).
|
contentAlignment = Alignment.Center,
|
||||||
// Gating only on `publishing` would hide the muted icon
|
) {
|
||||||
// the moment the user mutes, which is exactly when it's
|
Box(
|
||||||
// supposed to appear.
|
modifier = Modifier.padding(ringPadding),
|
||||||
if (showMicBadge && (member.publishing || member.muted == true)) {
|
contentAlignment = Alignment.Center,
|
||||||
MicStateBadge(
|
) {
|
||||||
|
AvatarAndBadges(
|
||||||
|
member = member,
|
||||||
|
avatarSize = avatarSize,
|
||||||
|
accountViewModel = accountViewModel,
|
||||||
|
avatarModifier = avatarModifier,
|
||||||
|
onClick = onClick,
|
||||||
|
onLongClick = onLongClick,
|
||||||
|
isConnecting = isConnecting,
|
||||||
|
showMicBadge = showMicBadge,
|
||||||
isSpeaking = isSpeaking,
|
isSpeaking = isSpeaking,
|
||||||
isMuted = member.muted == true,
|
|
||||||
modifier = Modifier.align(Alignment.BottomCenter),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
// Reactions float over the avatar's bottom-right corner so
|
|
||||||
// a 👏 burst no longer pushes the username down and reflows
|
|
||||||
// neighbouring cells. The mic badge sits at BottomCenter,
|
|
||||||
// so BottomEnd + a small outward offset keeps them clear.
|
|
||||||
if (reactions.isNotEmpty()) {
|
|
||||||
SpeakerReactionOverlay(
|
|
||||||
reactions = reactions,
|
reactions = reactions,
|
||||||
modifier =
|
|
||||||
Modifier
|
|
||||||
.align(Alignment.BottomEnd)
|
|
||||||
.offset(x = 6.dp, y = 6.dp),
|
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -545,6 +531,76 @@ private fun MemberCell(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Composable
|
||||||
|
private fun AvatarAndBadges(
|
||||||
|
member: RoomMember,
|
||||||
|
avatarSize: Dp,
|
||||||
|
accountViewModel: AccountViewModel,
|
||||||
|
avatarModifier: Modifier,
|
||||||
|
onClick: ((String) -> Unit)?,
|
||||||
|
onLongClick: ((String) -> Unit)?,
|
||||||
|
isConnecting: Boolean,
|
||||||
|
showMicBadge: Boolean,
|
||||||
|
isSpeaking: Boolean,
|
||||||
|
reactions: List<RoomReaction>,
|
||||||
|
) {
|
||||||
|
Box(contentAlignment = Alignment.Center) {
|
||||||
|
ClickableUserPicture(
|
||||||
|
baseUserHex = member.pubkey,
|
||||||
|
size = avatarSize,
|
||||||
|
accountViewModel = accountViewModel,
|
||||||
|
modifier = avatarModifier,
|
||||||
|
onClick = onClick,
|
||||||
|
onLongClick = onLongClick,
|
||||||
|
)
|
||||||
|
if (isConnecting) {
|
||||||
|
CircularProgressIndicator(
|
||||||
|
modifier = Modifier.size(avatarSize - 8.dp),
|
||||||
|
strokeWidth = 2.dp,
|
||||||
|
color = MaterialTheme.colorScheme.primary,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
val role = member.role
|
||||||
|
if (role == ROLE.HOST || role == ROLE.MODERATOR) {
|
||||||
|
RoleBadge(
|
||||||
|
role = role,
|
||||||
|
modifier = Modifier.align(Alignment.TopStart),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
if (member.handRaised) {
|
||||||
|
HandRaiseBadge(
|
||||||
|
modifier = Modifier.align(Alignment.TopEnd),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
// Show the mic badge for any on-stage speaker that has
|
||||||
|
// an audio state to surface — currently broadcasting
|
||||||
|
// (`publishing=1`) OR mic-muted (`muted=1, publishing=0`).
|
||||||
|
// Gating only on `publishing` would hide the muted icon
|
||||||
|
// the moment the user mutes, which is exactly when it's
|
||||||
|
// supposed to appear.
|
||||||
|
if (showMicBadge && (member.publishing || member.muted == true)) {
|
||||||
|
MicStateBadge(
|
||||||
|
isSpeaking = isSpeaking,
|
||||||
|
isMuted = member.muted == true,
|
||||||
|
modifier = Modifier.align(Alignment.BottomCenter),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
// Reactions float over the avatar's bottom-right corner so
|
||||||
|
// a 👏 burst no longer pushes the username down and reflows
|
||||||
|
// neighbouring cells. The mic badge sits at BottomCenter,
|
||||||
|
// so BottomEnd + a small outward offset keeps them clear.
|
||||||
|
if (reactions.isNotEmpty()) {
|
||||||
|
SpeakerReactionOverlay(
|
||||||
|
reactions = reactions,
|
||||||
|
modifier =
|
||||||
|
Modifier
|
||||||
|
.align(Alignment.BottomEnd)
|
||||||
|
.offset(x = 6.dp, y = 6.dp),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Hand-raise indicator overlaid on the avatar — yellow circle with
|
* Hand-raise indicator overlaid on the avatar — yellow circle with
|
||||||
* a hand glyph at the top-right, animated in a subtle vertical
|
* a hand glyph at the top-right, animated in a subtle vertical
|
||||||
|
|||||||
+199
-63
@@ -302,14 +302,26 @@ class NestViewModel(
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* [kotlin.time.TimeMark] of the most recent recycle the cliff
|
* [kotlin.time.TimeMark] of the most recent recycle the cliff
|
||||||
* detector triggered, or `null` if it has never fired. Acts as a
|
* detector triggered, or `null` if it has never fired. Combined
|
||||||
* cooldown so a single cliff event doesn't kick off multiple
|
* with [consecutiveCliffRecycles] to drive the per-attempt
|
||||||
* recycles back-to-back while the wrapper is mid-handshake on
|
* backoff schedule in [computeStalledSpeakers] — short backoff
|
||||||
* the new session — the new session has no incoming frames yet,
|
* after the first failed recycle so a single moq-rs cliff
|
||||||
* which would otherwise trip the detector again immediately.
|
* recovers within a few seconds, escalating to the
|
||||||
|
* [ROOM_AUDIO_CLIFF_BACKOFF_MAX_MS] cap if the relay stays
|
||||||
|
* stalled across multiple recycles.
|
||||||
*/
|
*/
|
||||||
private var lastCliffRecycleAt: kotlin.time.TimeMark? = null
|
private var lastCliffRecycleAt: kotlin.time.TimeMark? = null
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Number of consecutive `recycleSession()` calls the detector has
|
||||||
|
* issued without seeing a real frame in between. Reset to 0 in
|
||||||
|
* [onSpeakerActivity] when any speaker delivers a frame, so a
|
||||||
|
* recovered-then-restall pattern re-enters with attempt = 0
|
||||||
|
* (immediate recycle eligible) rather than inheriting backoff
|
||||||
|
* from a prior failed-recycle streak.
|
||||||
|
*/
|
||||||
|
private var consecutiveCliffRecycles: Int = 0
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Per-speaker catalog-fetch coroutines. Each entry is the
|
* Per-speaker catalog-fetch coroutines. Each entry is the
|
||||||
* background `subscribeCatalog` collector launched in
|
* background `subscribeCatalog` collector launched in
|
||||||
@@ -320,7 +332,19 @@ class NestViewModel(
|
|||||||
*/
|
*/
|
||||||
private val catalogJobs = mutableMapOf<String, Job>()
|
private val catalogJobs = mutableMapOf<String, Job>()
|
||||||
private var requestedSpeakers: Set<String> = emptySet()
|
private var requestedSpeakers: Set<String> = emptySet()
|
||||||
private var closed = false
|
|
||||||
|
/**
|
||||||
|
* `@Volatile` because [closed] is read by background coroutines
|
||||||
|
* (cliff-detector, observe* loops, audio-focus + network observers)
|
||||||
|
* but written from `leave()` / `onCleared()` which can fire from
|
||||||
|
* the Activity destroy callback on a different stack frame than
|
||||||
|
* the coroutine that's about to suspend in [delay] / [collect].
|
||||||
|
* Without the marker, a coroutine that resumes from cancellation
|
||||||
|
* could read a stale `closed = false` even though the user has
|
||||||
|
* already left the room — visible in the trace as
|
||||||
|
* `cliff-detector EXITED ... closed=false` after a Leave press.
|
||||||
|
*/
|
||||||
|
@Volatile private var closed = false
|
||||||
|
|
||||||
// Speaker / publisher path
|
// Speaker / publisher path
|
||||||
private var speaker: NestsSpeaker? = null
|
private var speaker: NestsSpeaker? = null
|
||||||
@@ -1389,6 +1413,7 @@ class NestViewModel(
|
|||||||
cliffDetectorJob = null
|
cliffDetectorJob = null
|
||||||
lastFrameAt.clear()
|
lastFrameAt.clear()
|
||||||
lastCliffRecycleAt = null
|
lastCliffRecycleAt = null
|
||||||
|
consecutiveCliffRecycles = 0
|
||||||
// Audio-focus observation only ends on the final teardown —
|
// Audio-focus observation only ends on the final teardown —
|
||||||
// a transient disconnect+reconnect (user retry, room swap)
|
// a transient disconnect+reconnect (user retry, room swap)
|
||||||
// keeps the bus subscription alive so a focus loss that
|
// keeps the bus subscription alive so a focus loss that
|
||||||
@@ -1469,9 +1494,17 @@ class NestViewModel(
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Mark [pubkey] as currently speaking and (re)arm a [SPEAKING_TIMEOUT_MS]
|
* Per-frame heartbeat. Called once per MoQ object received on the
|
||||||
* coroutine that clears it once they go quiet. Called once per
|
* speaker's track — i.e. once per ~20 ms regardless of whether
|
||||||
* MoQ object received on the speaker's track.
|
* that frame contains actual speech, silence, or background noise.
|
||||||
|
*
|
||||||
|
* Bumps the cliff-detector timestamp so an active stream keeps
|
||||||
|
* resetting the relay-forward-queue stall watchdog, and clears
|
||||||
|
* the per-speaker buffering overlay the first time a frame lands.
|
||||||
|
*
|
||||||
|
* NB: this does NOT mark the speaker as "speaking right now".
|
||||||
|
* That signal is energy-gated and lives in [onAudioLevel] —
|
||||||
|
* mic-on with no voice MUST NOT light up the green ring.
|
||||||
*/
|
*/
|
||||||
private fun onSpeakerActivity(pubkey: String) {
|
private fun onSpeakerActivity(pubkey: String) {
|
||||||
if (closed) return
|
if (closed) return
|
||||||
@@ -1483,12 +1516,27 @@ class NestViewModel(
|
|||||||
lastFrameAt[pubkey] =
|
lastFrameAt[pubkey] =
|
||||||
kotlin.time.TimeSource.Monotonic
|
kotlin.time.TimeSource.Monotonic
|
||||||
.markNow()
|
.markNow()
|
||||||
speakingExpiryJobs[pubkey]?.cancel()
|
// A real frame proves the most recent recycle (if any) actually
|
||||||
|
// fixed the cliff. Reset the consecutive-failed counter so a
|
||||||
|
// future re-stall starts from attempt 0 (immediate-fire) rather
|
||||||
|
// than inheriting a long backoff from the prior streak.
|
||||||
|
if (consecutiveCliffRecycles != 0) consecutiveCliffRecycles = 0
|
||||||
// First frame for this subscription — clear the buffering
|
// First frame for this subscription — clear the buffering
|
||||||
// overlay. Subsequent frames are no-ops here.
|
// overlay. Subsequent frames are no-ops here.
|
||||||
if (_uiState.value.connectingSpeakers.contains(pubkey)) {
|
if (_uiState.value.connectingSpeakers.contains(pubkey)) {
|
||||||
_uiState.update { it.copy(connectingSpeakers = (it.connectingSpeakers - pubkey).toPersistentSet()) }
|
_uiState.update { it.copy(connectingSpeakers = (it.connectingSpeakers - pubkey).toPersistentSet()) }
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Mark [pubkey] as currently speaking and (re)arm a
|
||||||
|
* [SPEAKING_TIMEOUT_MS] coroutine that clears the flag once their
|
||||||
|
* audio drops back below the threshold for that long. Called from
|
||||||
|
* [onAudioLevel] only when the decoded peak is loud enough to
|
||||||
|
* read as speech (see [SPEAKING_LEVEL_THRESHOLD]).
|
||||||
|
*/
|
||||||
|
private fun markSpeaking(pubkey: String) {
|
||||||
|
speakingExpiryJobs[pubkey]?.cancel()
|
||||||
if (!_uiState.value.speakingNow.contains(pubkey)) {
|
if (!_uiState.value.speakingNow.contains(pubkey)) {
|
||||||
_uiState.update { it.copy(speakingNow = (it.speakingNow + pubkey).toPersistentSet()) }
|
_uiState.update { it.copy(speakingNow = (it.speakingNow + pubkey).toPersistentSet()) }
|
||||||
}
|
}
|
||||||
@@ -1570,43 +1618,53 @@ class NestViewModel(
|
|||||||
announcedSpeakers = announced,
|
announcedSpeakers = announced,
|
||||||
lastFrameAt = lastFrameAt,
|
lastFrameAt = lastFrameAt,
|
||||||
lastRecycleAt = lastCliffRecycleAt,
|
lastRecycleAt = lastCliffRecycleAt,
|
||||||
|
consecutiveFailedRecycles = consecutiveCliffRecycles,
|
||||||
cliffTimeoutMs = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
|
cliffTimeoutMs = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
|
||||||
cooldownMs = ROOM_AUDIO_CLIFF_COOLDOWN_MS,
|
postRecycleGraceMs = ROOM_AUDIO_CLIFF_RECYCLE_GRACE_MS,
|
||||||
)
|
)
|
||||||
if (stalled.isEmpty()) continue
|
if (stalled.isEmpty()) continue
|
||||||
|
consecutiveCliffRecycles += 1
|
||||||
com.vitorpamplona.quartz.utils.Log.w("NestRx") {
|
com.vitorpamplona.quartz.utils.Log.w("NestRx") {
|
||||||
"cliff-detector: announced+subscribed but silent for ≥${ROOM_AUDIO_CLIFF_TIMEOUT_MS}ms — recycling session. stalled=$stalled"
|
"cliff-detector: announced+subscribed but silent for ≥${ROOM_AUDIO_CLIFF_TIMEOUT_MS}ms — recycling session " +
|
||||||
|
"(consecutive=$consecutiveCliffRecycles). stalled=$stalled"
|
||||||
}
|
}
|
||||||
com.vitorpamplona.nestsclient.trace.NestsTrace.emit("cliff_recycle") {
|
com.vitorpamplona.nestsclient.trace.NestsTrace.emit("cliff_recycle") {
|
||||||
"\"timeout_ms\":$ROOM_AUDIO_CLIFF_TIMEOUT_MS," +
|
"\"timeout_ms\":$ROOM_AUDIO_CLIFF_TIMEOUT_MS," +
|
||||||
|
"\"consecutive\":$consecutiveCliffRecycles," +
|
||||||
"\"stalled\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(stalled)}"
|
"\"stalled\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(stalled)}"
|
||||||
}
|
}
|
||||||
val recycleMark =
|
lastCliffRecycleAt =
|
||||||
kotlin.time.TimeSource.Monotonic
|
kotlin.time.TimeSource.Monotonic
|
||||||
.markNow()
|
.markNow()
|
||||||
lastCliffRecycleAt = recycleMark
|
// Note: we deliberately do NOT overwrite
|
||||||
// Reset `lastFrameAt` for every stalled pubkey so
|
// `lastFrameAt[pubkey]` with the recycle moment.
|
||||||
// the cliff timer starts counting from the recycle
|
// Keeping the real-frame timestamp lets the next
|
||||||
// moment, not from the old (pre-recycle) last
|
// tick distinguish "the recycle delivered audio,
|
||||||
// frame timestamp. Without this reset the next
|
// we then re-stalled" (lastFrameAt advances past
|
||||||
// tick after the cooldown immediately re-trips —
|
// lastCliffRecycleAt — counter resets in
|
||||||
// because `lastFrameAt[pubkey].elapsedNow()` is
|
// `onSpeakerActivity`, attempt = 0, fire-eligible
|
||||||
// still the giant pre-recycle value plus the
|
// immediately) from "the recycle did nothing,
|
||||||
// cooldown window. Production logs at commit
|
// still no frames" (lastFrameAt unchanged,
|
||||||
// ea08c43 showed exactly this: 4 recycles in
|
// counter keeps growing, backoff escalates).
|
||||||
// ~30 s eventually drove the relay into a
|
// The previous reset collapsed both cases into
|
||||||
// "subscribe stream FIN before reply" loop where
|
// a flat 30 s wait, which made a single failed
|
||||||
// it refused fresh subscribes entirely. Using
|
// recycle sound like a 30 s+ dropout to the
|
||||||
// the recycle moment as a synthetic frame event
|
// user even when retrying earlier would have
|
||||||
// gives the new session [ROOM_AUDIO_CLIFF_TIMEOUT_MS]
|
// recovered.
|
||||||
// wall-clock to deliver before the next cliff
|
try {
|
||||||
// check, matching the expectation a freshly-
|
l.recycleSession()
|
||||||
// attached subscription should clear inside
|
} catch (ce: CancellationException) {
|
||||||
// that window if the relay is healthy.
|
// User left mid-recycle — exit promptly
|
||||||
for (pubkey in stalled) {
|
// rather than letting the next delay() be
|
||||||
lastFrameAt[pubkey] = recycleMark
|
// the one to honor the cancel. `runCatching`
|
||||||
|
// would have swallowed this CE and run one
|
||||||
|
// more loop body before exiting.
|
||||||
|
throw ce
|
||||||
|
} catch (_: Throwable) {
|
||||||
|
// Any other recycle failure is best-effort;
|
||||||
|
// keep the loop running so a follow-up tick
|
||||||
|
// can re-detect the cliff and retry.
|
||||||
}
|
}
|
||||||
runCatching { l.recycleSession() }
|
|
||||||
}
|
}
|
||||||
} finally {
|
} finally {
|
||||||
com.vitorpamplona.quartz.utils.Log
|
com.vitorpamplona.quartz.utils.Log
|
||||||
@@ -1641,6 +1699,20 @@ class NestViewModel(
|
|||||||
) {
|
) {
|
||||||
if (closed) return
|
if (closed) return
|
||||||
rawAudioLevels[pubkey] = level
|
rawAudioLevels[pubkey] = level
|
||||||
|
// Energy-gated speaking detector. The MoQ track delivers a
|
||||||
|
// frame every ~20 ms while the mic is open, even when the
|
||||||
|
// speaker is silent or only picking up room noise — gating the
|
||||||
|
// green ring on "frame arrived" therefore lights it up the
|
||||||
|
// moment the mic is unmuted, not when there's actually a voice
|
||||||
|
// on it. The decoded peak amplitude (`peakAmplitude` in
|
||||||
|
// nestsClient/audio/Amplitude.kt) gives us the signal we need:
|
||||||
|
// background noise / breath stays under a few percent of full
|
||||||
|
// scale, while even a quiet voice clears [SPEAKING_LEVEL_THRESHOLD].
|
||||||
|
// The 250 ms expiry already wired up in [markSpeaking] gives
|
||||||
|
// the indicator natural hysteresis between syllables.
|
||||||
|
if (level >= SPEAKING_LEVEL_THRESHOLD) {
|
||||||
|
markSpeaking(pubkey)
|
||||||
|
}
|
||||||
startLevelEmitter()
|
startLevelEmitter()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1799,12 +1871,29 @@ sealed class BroadcastUiState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* How long a speaker stays "speaking" after their last received MoQ object.
|
* How long a speaker stays "speaking" after their last decoded frame
|
||||||
* Roughly 12 × the 20 ms Opus frame so brief packet jitter doesn't make the
|
* over the [SPEAKING_LEVEL_THRESHOLD]. Roughly 12 × the 20 ms Opus
|
||||||
* indicator flicker.
|
* frame so brief packet jitter and inter-syllable pauses don't make
|
||||||
|
* the indicator flicker between adjacent words.
|
||||||
*/
|
*/
|
||||||
const val SPEAKING_TIMEOUT_MS: Long = 250L
|
const val SPEAKING_TIMEOUT_MS: Long = 250L
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Minimum decoded peak amplitude (normalized to `[0, 1]`) that counts
|
||||||
|
* as "this person is actually speaking right now". Frames whose peak
|
||||||
|
* lands below this threshold are treated as silence / room tone /
|
||||||
|
* breath — they keep the per-speaker subscription healthy (the cliff
|
||||||
|
* heartbeat in [NestViewModel.onSpeakerActivity] still fires) but do
|
||||||
|
* NOT light up the green speaking ring.
|
||||||
|
*
|
||||||
|
* 0.06 ≈ -24 dBFS, comfortably above the typical residential-mic
|
||||||
|
* noise floor (~-40 to -30 dBFS) while still tripping on a quiet
|
||||||
|
* voice. Tuned in conjunction with [SPEAKING_TIMEOUT_MS]: a single
|
||||||
|
* loud frame is enough to arm the indicator; ≥ 250 ms below the
|
||||||
|
* threshold drops it.
|
||||||
|
*/
|
||||||
|
const val SPEAKING_LEVEL_THRESHOLD: Float = 0.06f
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* How long [NestViewModel.openSubscription] waits for the publisher's
|
* How long [NestViewModel.openSubscription] waits for the publisher's
|
||||||
* `catalog.json` to land before constructing the decoder + AudioTrack.
|
* `catalog.json` to land before constructing the decoder + AudioTrack.
|
||||||
@@ -1926,24 +2015,34 @@ const val ROOM_AUDIO_CLIFF_TIMEOUT_MS: Long = 2_500L
|
|||||||
const val ROOM_AUDIO_CLIFF_CHECK_INTERVAL_MS: Long = 1_000L
|
const val ROOM_AUDIO_CLIFF_CHECK_INTERVAL_MS: Long = 1_000L
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Cooldown after a cliff-detector-driven recycle. Suppresses
|
* Post-recycle handshake grace. NEVER recycle again within this
|
||||||
* back-to-back recycles while the wrapper is mid-handshake on the
|
* window of the previous recycle — the wrapper is still tearing
|
||||||
* fresh session AND while the relay is recovering its
|
* down + reopening the QUIC session, so a "no frames since recycle"
|
||||||
* per-subscriber forward queue.
|
* reading is just the handshake, not a relay-side cliff.
|
||||||
*
|
*
|
||||||
* Production logs at commit ea08c43 showed an 8 s cooldown was
|
* 3 s covers a typical re-handshake (drain old session UNSUB +
|
||||||
* not enough — moq-rs needed longer to recover between aggressive
|
* UNANNOUNCE + WT_CLOSE + open new WebTransport + SUBSCRIBE +
|
||||||
* recycles, and 4 recycles in ~30 s drove the relay into a
|
* SUBSCRIBE_OK). Anything shorter risks recycling inside our own
|
||||||
* "subscribe stream FIN before reply" loop that refused all
|
* handshake window; longer wastes audible silence in the recovering-
|
||||||
* subsequent subscribes. 30 s gives the relay's per-subscriber
|
* within-grace case.
|
||||||
* 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
|
const val ROOM_AUDIO_CLIFF_RECYCLE_GRACE_MS: Long = 3_000L
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Maximum backoff for the consecutive-failed-recycle schedule —
|
||||||
|
* the cap that [defaultCliffBackoffMs] saturates at.
|
||||||
|
*
|
||||||
|
* The earlier flat 30 s cooldown was motivated by production logs
|
||||||
|
* at commit ea08c43 where 4 back-to-back recycles in ~30 s drove
|
||||||
|
* moq-rs into a "subscribe stream FIN before reply" loop that
|
||||||
|
* refused all subsequent subscribes. The replacement schedule
|
||||||
|
* (5 s → 12 s → 24 s → 30 s, reset on first real frame) keeps
|
||||||
|
* the same final-state protection — by the 4th consecutive failed
|
||||||
|
* recycle we're spaced 30 s apart — while letting the *recovering*
|
||||||
|
* case (your trace: 1st recycle worked, 2nd cliff fires later)
|
||||||
|
* retry within ~5 s instead of the previous 30 s of dead air.
|
||||||
|
*/
|
||||||
|
const val ROOM_AUDIO_CLIFF_BACKOFF_MAX_MS: Long = 30_000L
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Diagnostic log frequency for the cliff detector. Emit a state-dump
|
* Diagnostic log frequency for the cliff detector. Emit a state-dump
|
||||||
@@ -1954,6 +2053,36 @@ const val ROOM_AUDIO_CLIFF_COOLDOWN_MS: Long = 30_000L
|
|||||||
*/
|
*/
|
||||||
private const val CLIFF_DIAG_LOG_EVERY: Long = 5L
|
private const val CLIFF_DIAG_LOG_EVERY: Long = 5L
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Default backoff schedule for consecutive failed recycles. Returns
|
||||||
|
* the minimum elapsed time since [NestViewModel.lastCliffRecycleAt]
|
||||||
|
* before the [attempt + 1]-th recycle is permitted.
|
||||||
|
*
|
||||||
|
* attempt 0 → 0 ms (no recycle yet — first cliff fires immediately)
|
||||||
|
* attempt 1 → 5 000 ms (one prior recycle, no audio since)
|
||||||
|
* attempt 2 → 12 000 ms
|
||||||
|
* attempt 3 → 24 000 ms
|
||||||
|
* attempt 4+ → 30 000 ms (cap, [ROOM_AUDIO_CLIFF_BACKOFF_MAX_MS])
|
||||||
|
*
|
||||||
|
* Resets to 0 on the first real frame after a recycle (see
|
||||||
|
* [NestViewModel.onSpeakerActivity]) — so a recover-then-restall
|
||||||
|
* pattern always gets the immediate-fire treatment again rather
|
||||||
|
* than inheriting backoff from the prior failed-recycle streak.
|
||||||
|
*
|
||||||
|
* Cumulative wall-clock at attempt N: 0 → 5 → 17 → 41 → 71 s.
|
||||||
|
* 4 consecutive recycles span ~41 s — meaningfully slower than the
|
||||||
|
* "4 in ~30 s" pattern that wedged moq-rs in commit ea08c43, while
|
||||||
|
* still letting the typical 1-failure case retry inside 5 s.
|
||||||
|
*/
|
||||||
|
internal fun defaultCliffBackoffMs(attempt: Int): Long =
|
||||||
|
when {
|
||||||
|
attempt <= 0 -> 0L
|
||||||
|
attempt == 1 -> 5_000L
|
||||||
|
attempt == 2 -> 12_000L
|
||||||
|
attempt == 3 -> 24_000L
|
||||||
|
else -> ROOM_AUDIO_CLIFF_BACKOFF_MAX_MS
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Pure logic of the cliff detector — extracted so headless tests can
|
* Pure logic of the cliff detector — extracted so headless tests can
|
||||||
* exercise it with a [kotlin.time.TestTimeSource] without standing up
|
* exercise it with a [kotlin.time.TestTimeSource] without standing up
|
||||||
@@ -1974,23 +2103,30 @@ private const val CLIFF_DIAG_LOG_EVERY: Long = 5L
|
|||||||
* - the elapsed time since that last object is at or past
|
* - the elapsed time since that last object is at or past
|
||||||
* [cliffTimeoutMs].
|
* [cliffTimeoutMs].
|
||||||
*
|
*
|
||||||
* Suppression:
|
* Suppression (when [lastRecycleAt] is non-null):
|
||||||
* - if [lastRecycleAt] is non-null and less than [cooldownMs] has
|
* - while elapsed since [lastRecycleAt] is below [postRecycleGraceMs],
|
||||||
* elapsed since it, returns empty (no cascading recycles while
|
* never recycle — we're still inside the previous handshake window.
|
||||||
* the wrapper is mid-handshake on a fresh session).
|
* - while elapsed is below [backoffForAttempt]([consecutiveFailedRecycles]),
|
||||||
|
* suppress as well: the per-attempt schedule throttles
|
||||||
|
* consecutive failed recycles so we don't hammer the relay.
|
||||||
|
* - the consecutive counter is reset by the caller on the first
|
||||||
|
* real frame after a recycle, so a recovered-then-restall
|
||||||
|
* case always re-enters with attempt = 0 and fires immediately.
|
||||||
*/
|
*/
|
||||||
internal fun computeStalledSpeakers(
|
internal fun computeStalledSpeakers(
|
||||||
activeSpeakers: Set<String>,
|
activeSpeakers: Set<String>,
|
||||||
announcedSpeakers: Set<String>,
|
announcedSpeakers: Set<String>,
|
||||||
lastFrameAt: Map<String, kotlin.time.TimeMark>,
|
lastFrameAt: Map<String, kotlin.time.TimeMark>,
|
||||||
lastRecycleAt: kotlin.time.TimeMark?,
|
lastRecycleAt: kotlin.time.TimeMark?,
|
||||||
|
consecutiveFailedRecycles: Int = 0,
|
||||||
cliffTimeoutMs: Long = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
|
cliffTimeoutMs: Long = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
|
||||||
cooldownMs: Long = ROOM_AUDIO_CLIFF_COOLDOWN_MS,
|
postRecycleGraceMs: Long = ROOM_AUDIO_CLIFF_RECYCLE_GRACE_MS,
|
||||||
|
backoffForAttempt: (Int) -> Long = ::defaultCliffBackoffMs,
|
||||||
): List<String> {
|
): List<String> {
|
||||||
if (lastRecycleAt != null &&
|
if (lastRecycleAt != null) {
|
||||||
lastRecycleAt.elapsedNow().inWholeMilliseconds < cooldownMs
|
val sinceRecycleMs = lastRecycleAt.elapsedNow().inWholeMilliseconds
|
||||||
) {
|
if (sinceRecycleMs < postRecycleGraceMs) return emptyList()
|
||||||
return emptyList()
|
if (sinceRecycleMs < backoffForAttempt(consecutiveFailedRecycles)) return emptyList()
|
||||||
}
|
}
|
||||||
return activeSpeakers
|
return activeSpeakers
|
||||||
.asSequence()
|
.asSequence()
|
||||||
|
|||||||
+169
-19
@@ -189,21 +189,16 @@ class CliffDetectorTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun cooldownSuppressesRecycleEvenWhenStalled() {
|
fun postRecycleGraceSuppressesEvenAtAttemptZero() {
|
||||||
// After a recycle, the wrapper opens a fresh QUIC transport.
|
// Inside the 3 s post-recycle handshake window, never recycle
|
||||||
// The new session has no `lastFrameAt` entries yet for any
|
// — the wrapper is still tearing down + reopening QUIC, so
|
||||||
// pubkey; without a cooldown we would re-trigger on the
|
// "no frames since recycle" is just the handshake, not a
|
||||||
// very next 1 s tick because the prior tick's stalled
|
// relay-side cliff.
|
||||||
// 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 ts = TestTimeSource()
|
||||||
val frameMark = ts.markNow()
|
val frameMark = ts.markNow()
|
||||||
ts += 4_500.milliseconds // past threshold
|
ts += 4_500.milliseconds
|
||||||
val recycleMark = ts.markNow() // recycle just fired
|
val recycleMark = ts.markNow()
|
||||||
ts += 5_000.milliseconds // 5 s into cooldown — well within 30 s window
|
ts += 2_000.milliseconds // inside 3 s grace
|
||||||
|
|
||||||
val result =
|
val result =
|
||||||
computeStalledSpeakers(
|
computeStalledSpeakers(
|
||||||
@@ -211,21 +206,22 @@ class CliffDetectorTest {
|
|||||||
announcedSpeakers = setOf(ALICE),
|
announcedSpeakers = setOf(ALICE),
|
||||||
lastFrameAt = mapOf(ALICE to frameMark),
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
lastRecycleAt = recycleMark,
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 0,
|
||||||
)
|
)
|
||||||
assertTrue(result.isEmpty())
|
assertTrue(result.isEmpty())
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun cooldownReleasesAfterTimeoutPasses() {
|
fun attemptZeroFiresImmediatelyOnceGracePasses() {
|
||||||
// Once cooldown elapses, a still-stalled subscription
|
// After a recovered-then-restall pattern: counter has been
|
||||||
// becomes eligible to recycle again. Important for the
|
// reset by `onSpeakerActivity`, so even though there's a
|
||||||
// case where the recycle didn't actually fix the cliff —
|
// prior recycleMark, the next cliff fires as soon as the
|
||||||
// we want a second attempt rather than getting wedged.
|
// 3 s grace passes — no extended backoff.
|
||||||
val ts = TestTimeSource()
|
val ts = TestTimeSource()
|
||||||
val frameMark = ts.markNow()
|
val frameMark = ts.markNow()
|
||||||
ts += 4_500.milliseconds
|
ts += 4_500.milliseconds
|
||||||
val recycleMark = ts.markNow()
|
val recycleMark = ts.markNow()
|
||||||
ts += 30_001.milliseconds // 1 ms past 30 s cooldown
|
ts += 3_500.milliseconds // past 3 s grace, attempt=0 schedule = 0 ms
|
||||||
|
|
||||||
val result =
|
val result =
|
||||||
computeStalledSpeakers(
|
computeStalledSpeakers(
|
||||||
@@ -233,10 +229,164 @@ class CliffDetectorTest {
|
|||||||
announcedSpeakers = setOf(ALICE),
|
announcedSpeakers = setOf(ALICE),
|
||||||
lastFrameAt = mapOf(ALICE to frameMark),
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
lastRecycleAt = recycleMark,
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 0,
|
||||||
)
|
)
|
||||||
assertEquals(listOf(ALICE), result)
|
assertEquals(listOf(ALICE), result)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun attemptOneBackoffSuppressesUntilFiveSeconds() {
|
||||||
|
// First failed recycle: schedule says wait 5 s before next.
|
||||||
|
// 4 s in: still suppressed.
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 4_500.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 4_000.milliseconds
|
||||||
|
|
||||||
|
val result =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 1,
|
||||||
|
)
|
||||||
|
assertTrue(result.isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun attemptOneBackoffReleasesAtFiveSeconds() {
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 4_500.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 5_000.milliseconds // exactly at attempt-1 boundary
|
||||||
|
|
||||||
|
val result =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 1,
|
||||||
|
)
|
||||||
|
assertEquals(listOf(ALICE), result)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun attemptTwoBackoffSuppressesUntilTwelveSeconds() {
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 4_500.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 11_000.milliseconds
|
||||||
|
|
||||||
|
val result =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 2,
|
||||||
|
)
|
||||||
|
assertTrue(result.isEmpty())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun attemptTwoBackoffReleasesPastTwelveSeconds() {
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 4_500.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 12_500.milliseconds
|
||||||
|
|
||||||
|
val result =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 2,
|
||||||
|
)
|
||||||
|
assertEquals(listOf(ALICE), result)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun attemptFourCapsAtThirtySecondMax() {
|
||||||
|
// Fourth and beyond consecutive failed recycle: backoff
|
||||||
|
// saturates at the 30 s cap (matching the original flat
|
||||||
|
// cooldown — by this point we ARE the moq-rs-protection
|
||||||
|
// case the old constant existed for).
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 4_500.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 25_000.milliseconds // past attempt-3 (24 s), inside attempt-4 cap (30 s)
|
||||||
|
|
||||||
|
val resultStillSuppressed =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 4,
|
||||||
|
)
|
||||||
|
assertTrue(resultStillSuppressed.isEmpty())
|
||||||
|
|
||||||
|
ts += 6_000.milliseconds // now past 30 s cap
|
||||||
|
val resultReleased =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 4,
|
||||||
|
)
|
||||||
|
assertEquals(listOf(ALICE), resultReleased)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun customBackoffFunctionIsHonored() {
|
||||||
|
// A test can override the schedule (e.g. shorter intervals
|
||||||
|
// for unit tests that don't want to march through 30 s of
|
||||||
|
// virtual time) without mutating the production constants.
|
||||||
|
val ts = TestTimeSource()
|
||||||
|
val frameMark = ts.markNow()
|
||||||
|
ts += 3_000.milliseconds
|
||||||
|
val recycleMark = ts.markNow()
|
||||||
|
ts += 1_000.milliseconds // past grace=500, past tightBackoff(1)=750
|
||||||
|
|
||||||
|
val tightBackoff = { attempt: Int -> if (attempt <= 0) 0L else 750L }
|
||||||
|
val result =
|
||||||
|
computeStalledSpeakers(
|
||||||
|
activeSpeakers = setOf(ALICE),
|
||||||
|
announcedSpeakers = setOf(ALICE),
|
||||||
|
lastFrameAt = mapOf(ALICE to frameMark),
|
||||||
|
lastRecycleAt = recycleMark,
|
||||||
|
consecutiveFailedRecycles = 1,
|
||||||
|
postRecycleGraceMs = 500L,
|
||||||
|
backoffForAttempt = tightBackoff,
|
||||||
|
)
|
||||||
|
assertEquals(listOf(ALICE), result)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun defaultBackoffSchedulePinsValues() {
|
||||||
|
// Pin the production schedule so a future tweak is visible
|
||||||
|
// in code review. attempt 0 → immediate, 1 → 5 s, 2 → 12 s,
|
||||||
|
// 3 → 24 s, 4+ → 30 s cap. Cumulative wall-clock to the
|
||||||
|
// Nth recycle: 0, 5, 17, 41, 71 s — slower than the
|
||||||
|
// 4-recycles-in-30 s pattern that wedged moq-rs in
|
||||||
|
// commit ea08c43.
|
||||||
|
assertEquals(0L, defaultCliffBackoffMs(0))
|
||||||
|
assertEquals(5_000L, defaultCliffBackoffMs(1))
|
||||||
|
assertEquals(12_000L, defaultCliffBackoffMs(2))
|
||||||
|
assertEquals(24_000L, defaultCliffBackoffMs(3))
|
||||||
|
assertEquals(30_000L, defaultCliffBackoffMs(4))
|
||||||
|
assertEquals(30_000L, defaultCliffBackoffMs(10))
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun customTimeoutsAreHonored() {
|
fun customTimeoutsAreHonored() {
|
||||||
// Defaults are wired into the production VM, but the function
|
// Defaults are wired into the production VM, but the function
|
||||||
|
|||||||
+50
-17
@@ -55,6 +55,7 @@ import kotlin.test.assertEquals
|
|||||||
import kotlin.test.assertFalse
|
import kotlin.test.assertFalse
|
||||||
import kotlin.test.assertIs
|
import kotlin.test.assertIs
|
||||||
import kotlin.test.assertTrue
|
import kotlin.test.assertTrue
|
||||||
|
import kotlin.time.Duration.Companion.seconds
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Drives [NestViewModel] with a fake [NestsListenerConnector] and
|
* Drives [NestViewModel] with a fake [NestsListenerConnector] and
|
||||||
@@ -72,6 +73,16 @@ import kotlin.test.assertTrue
|
|||||||
*/
|
*/
|
||||||
@OptIn(ExperimentalCoroutinesApi::class)
|
@OptIn(ExperimentalCoroutinesApi::class)
|
||||||
class NestViewModelTest {
|
class NestViewModelTest {
|
||||||
|
// Tracks every VM created by `newViewModel` so each test can dispose them
|
||||||
|
// before runTest exits. `connect()` starts an infinite cliff-detector
|
||||||
|
// loop in viewModelScope; without an explicit teardown that loop spins
|
||||||
|
// forever under runTest's virtual scheduler and the test wedges until
|
||||||
|
// the real-time runTest deadline (60 s × N tests => :commons:jvmTest
|
||||||
|
// hangs). [runVmTest] wraps each test body with a finally that calls
|
||||||
|
// `disconnect()` (which cancels cliffDetectorJob via teardown()) so the
|
||||||
|
// scheduler can become idle and runTest can return.
|
||||||
|
private val createdVms = mutableListOf<NestViewModel>()
|
||||||
|
|
||||||
@BeforeTest
|
@BeforeTest
|
||||||
fun setupMainDispatcher() {
|
fun setupMainDispatcher() {
|
||||||
Dispatchers.setMain(UnconfinedTestDispatcher())
|
Dispatchers.setMain(UnconfinedTestDispatcher())
|
||||||
@@ -79,12 +90,34 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@AfterTest
|
@AfterTest
|
||||||
fun resetMainDispatcher() {
|
fun resetMainDispatcher() {
|
||||||
|
// Belt-and-braces: if a test crashes before its finally runs, this
|
||||||
|
// still tears down so the next test starts clean.
|
||||||
|
createdVms.forEach { runCatching { it.disconnect() } }
|
||||||
|
createdVms.clear()
|
||||||
Dispatchers.resetMain()
|
Dispatchers.resetMain()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Test wrapper that guarantees every VM created via [newViewModel] is
|
||||||
|
* disconnected before [runTest] tries to drain the test scheduler. The
|
||||||
|
* 10-second [runTest] timeout is a safety net — a healthy test in this
|
||||||
|
* class finishes in milliseconds; if we trip the timeout it almost
|
||||||
|
* certainly means a new code path is leaving viewModelScope coroutines
|
||||||
|
* running and needs its own teardown call.
|
||||||
|
*/
|
||||||
|
private fun runVmTest(body: suspend TestScope.() -> Unit) =
|
||||||
|
runTest(timeout = 10.seconds) {
|
||||||
|
try {
|
||||||
|
body()
|
||||||
|
} finally {
|
||||||
|
createdVms.forEach { runCatching { it.disconnect() } }
|
||||||
|
createdVms.clear()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun connectShowsConnectingThenConnected() =
|
fun connectShowsConnectingThenConnected() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
|
|
||||||
@@ -103,7 +136,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun listenerFailedSurfacesAsUiFailed() =
|
fun listenerFailedSurfacesAsUiFailed() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
|
|
||||||
@@ -117,7 +150,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun connectorThrowsBecomesUiFailed() =
|
fun connectorThrowsBecomesUiFailed() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { throw NestsException("dns blew up") }
|
val vm = newViewModel { throw NestsException("dns blew up") }
|
||||||
|
|
||||||
vm.connect()
|
vm.connect()
|
||||||
@@ -129,7 +162,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun setMutedFlipsUiStateAndIsRetained() =
|
fun setMutedFlipsUiStateAndIsRetained() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
|
|
||||||
assertFalse(vm.uiState.value.isMuted)
|
assertFalse(vm.uiState.value.isMuted)
|
||||||
@@ -141,7 +174,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onStageNowDefaultsTrueAndSetOnStageFlipsIt() =
|
fun onStageNowDefaultsTrueAndSetOnStageFlipsIt() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
|
|
||||||
// Defaults to true so a freshly-joined speaker advertises
|
// Defaults to true so a freshly-joined speaker advertises
|
||||||
@@ -156,7 +189,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onPresenceEventPopulatesPresencesMapAndDedupesByPubkey() =
|
fun onPresenceEventPopulatesPresencesMapAndDedupesByPubkey() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
|
|
||||||
val alice = "a".repeat(64)
|
val alice = "a".repeat(64)
|
||||||
@@ -187,7 +220,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun evictStalePresencesDropsOldPeers() =
|
fun evictStalePresencesDropsOldPeers() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
val alice = "a".repeat(64)
|
val alice = "a".repeat(64)
|
||||||
val bob = "b".repeat(64)
|
val bob = "b".repeat(64)
|
||||||
@@ -219,7 +252,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onChatEventAccumulatesMessagesSortedByCreatedAt() =
|
fun onChatEventAccumulatesMessagesSortedByCreatedAt() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
val alice = "a".repeat(64)
|
val alice = "a".repeat(64)
|
||||||
|
|
||||||
@@ -255,7 +288,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onChatEventDedupesByEventId() =
|
fun onChatEventDedupesByEventId() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
val alice = "a".repeat(64)
|
val alice = "a".repeat(64)
|
||||||
val msg =
|
val msg =
|
||||||
@@ -281,7 +314,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onReactionEventGroupsByTargetAndEvictsOnTick() =
|
fun onReactionEventGroupsByTargetAndEvictsOnTick() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val vm = newViewModel { FakeNestsListener() }
|
val vm = newViewModel { FakeNestsListener() }
|
||||||
val alice = "a".repeat(64)
|
val alice = "a".repeat(64)
|
||||||
val bob = "b".repeat(64)
|
val bob = "b".repeat(64)
|
||||||
@@ -326,7 +359,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onKickFlipsWasKickedAndDisconnects() =
|
fun onKickFlipsWasKickedAndDisconnects() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
vm.connect()
|
vm.connect()
|
||||||
@@ -343,7 +376,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun onKickIsIdempotent() =
|
fun onKickIsIdempotent() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
vm.connect()
|
vm.connect()
|
||||||
@@ -377,7 +410,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun connectIsIdempotentWhileConnecting() =
|
fun connectIsIdempotentWhileConnecting() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
var connectCalls = 0
|
var connectCalls = 0
|
||||||
val vm =
|
val vm =
|
||||||
@@ -395,7 +428,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun disconnectReturnsToIdleAndClosesListener() =
|
fun disconnectReturnsToIdleAndClosesListener() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
|
|
||||||
@@ -411,7 +444,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun connectingStepMapsThroughToUiStep() =
|
fun connectingStepMapsThroughToUiStep() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
|
|
||||||
@@ -429,7 +462,7 @@ class NestViewModelTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun speakingNowClearsOnTeardown() =
|
fun speakingNowClearsOnTeardown() =
|
||||||
runTest {
|
runVmTest {
|
||||||
val fakeListener = FakeNestsListener()
|
val fakeListener = FakeNestsListener()
|
||||||
val vm = newViewModel { fakeListener }
|
val vm = newViewModel { fakeListener }
|
||||||
|
|
||||||
@@ -466,7 +499,7 @@ class NestViewModelTest {
|
|||||||
// Wire to the test's backgroundScope so close calls run during
|
// Wire to the test's backgroundScope so close calls run during
|
||||||
// the test rather than escaping to the real GlobalScope.
|
// the test rather than escaping to the real GlobalScope.
|
||||||
cleanupScope = backgroundScope,
|
cleanupScope = backgroundScope,
|
||||||
)
|
).also { createdVms.add(it) }
|
||||||
|
|
||||||
private class FakeNestsListener : NestsListener {
|
private class FakeNestsListener : NestsListener {
|
||||||
private val mutable = MutableStateFlow<NestsListenerState>(NestsListenerState.Idle)
|
private val mutable = MutableStateFlow<NestsListenerState>(NestsListenerState.Idle)
|
||||||
|
|||||||
@@ -1,151 +0,0 @@
|
|||||||
# Stream priority for moq-lite group uni streams (T11.3 follow-up)
|
|
||||||
|
|
||||||
**Status**: deferred — flagged in T11 (drop `bestEffort=true`), not landed.
|
|
||||||
|
|
||||||
## Why
|
|
||||||
|
|
||||||
`T11` removed `bestEffort=true` from `MoqLiteSession.openGroupStream`
|
|
||||||
because the QUIC contract it relied on (drop lost ranges silently
|
|
||||||
without `RESET_STREAM`) created undeliverable streams that wedge
|
|
||||||
peer reassembly buffers — the actual user-visible bug was a 30 s
|
|
||||||
silent dropout per lost packet on lossy networks.
|
|
||||||
|
|
||||||
`bestEffort=true` was incidentally providing a "newer-groups-skip-
|
|
||||||
queued-retransmits" effect: the writer never queued retransmits in
|
|
||||||
the first place, so under congestion the loss budget naturally
|
|
||||||
biased toward dropping older lost ranges. With `bestEffort=true`
|
|
||||||
gone, all retransmits are now queued, and under sustained
|
|
||||||
congestion the writer drains streams in `streamRoundRobinStart`
|
|
||||||
order — which can mean the listener catches up on a stale group
|
|
||||||
when a fresh one would be more useful.
|
|
||||||
|
|
||||||
The kixelated reference (`moq-rs`'s `Publisher::serve_group` in
|
|
||||||
`rs/moq-lite/src/lite/publisher.rs:347-406`) addresses this by
|
|
||||||
calling `stream.set_priority(priority.current())` on every group
|
|
||||||
stream and biasing the writer toward higher-priority (newer)
|
|
||||||
streams. We should do the same.
|
|
||||||
|
|
||||||
## Target shape
|
|
||||||
|
|
||||||
### `:quic` API
|
|
||||||
|
|
||||||
Add a stable priority hook to `QuicStream`:
|
|
||||||
|
|
||||||
```kotlin
|
|
||||||
class QuicStream(...) {
|
|
||||||
@Volatile
|
|
||||||
var priority: Int = 0
|
|
||||||
// Higher drains first under contention. Default 0 = unchanged
|
|
||||||
// round-robin behaviour.
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
Expose at the WebTransport layer:
|
|
||||||
|
|
||||||
```kotlin
|
|
||||||
interface WebTransportWriteStream {
|
|
||||||
fun setPriority(priority: Int)
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
### `QuicConnectionWriter` — the load-bearing change
|
|
||||||
|
|
||||||
Today's send-frame loop (`QuicConnectionWriter.kt:411-416`):
|
|
||||||
|
|
||||||
```kotlin
|
|
||||||
val streamsView = conn.streamsListLocked()
|
|
||||||
val start = conn.streamRoundRobinStart % streamsView.size
|
|
||||||
for (i in streamsView.indices) {
|
|
||||||
val stream = streamsView[(start + i) % streamsView.size]
|
|
||||||
// ...
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
Replace with priority-then-round-robin:
|
|
||||||
|
|
||||||
```kotlin
|
|
||||||
val streamsView = conn.streamsListLocked()
|
|
||||||
val sorted = streamsView.sortedByDescending { it.priority }
|
|
||||||
val start = conn.streamRoundRobinStart % sorted.size
|
|
||||||
for (i in sorted.indices) {
|
|
||||||
val stream = sorted[(start + i) % sorted.size]
|
|
||||||
// ...
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
Sort is stable; same-priority streams retain insertion order, so the
|
|
||||||
existing round-robin behaviour holds within a priority tier. Higher-
|
|
||||||
priority streams always come first in the iteration.
|
|
||||||
|
|
||||||
Cost: O(N log N) per drain pass, where N is active local-initiated
|
|
||||||
streams. N is small (1–10 in the moq-lite audio path); the
|
|
||||||
allocation of `sorted` per pass is the only real cost. If that
|
|
||||||
shows up in profiling, switch to `kotlin.collections.IntArray`-
|
|
||||||
backed indirect sort or maintain a priority-sorted list incrementally
|
|
||||||
on `setPriority` calls.
|
|
||||||
|
|
||||||
### moq-lite wiring
|
|
||||||
|
|
||||||
`MoqLiteSession.openGroupStream:1022-1037`:
|
|
||||||
|
|
||||||
```kotlin
|
|
||||||
internal suspend fun openGroupStream(
|
|
||||||
subscribeId: Long,
|
|
||||||
sequence: Long,
|
|
||||||
): WebTransportWriteStream {
|
|
||||||
val uni = transport.openUniStream()
|
|
||||||
uni.setPriority(sequence.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
|
||||||
uni.write(Varint.encode(MoqLiteDataType.Group.code))
|
|
||||||
uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
|
|
||||||
return uni
|
|
||||||
}
|
|
||||||
```
|
|
||||||
|
|
||||||
Newer groups have higher sequence → higher priority → drain first
|
|
||||||
under congestion. The saturation conversion handles broadcasts that
|
|
||||||
run long enough for `sequence` to exceed `Int.MAX_VALUE` (≈ 71
|
|
||||||
years at our 1 group/sec production cadence; defensive only).
|
|
||||||
|
|
||||||
## Test
|
|
||||||
|
|
||||||
Add a unit test in `:quic` that builds a `QuicConnection`, opens
|
|
||||||
two streams, sets `.priority = 0` on the first and `.priority = 1`
|
|
||||||
on the second, queues bytes on both, and verifies the higher-
|
|
||||||
priority stream's bytes hit the wire first. Doesn't need to fake
|
|
||||||
real flow-control backpressure — pinning the iteration order via
|
|
||||||
the writer's emitted-frames tape is sufficient to catch the
|
|
||||||
regression case ("a later refactor accidentally re-introduces
|
|
||||||
round-robin order").
|
|
||||||
|
|
||||||
A second test in `:nestsClient` that opens 3 group streams and
|
|
||||||
asserts later-sequence streams drain before earlier ones is
|
|
||||||
nice-to-have but secondary; the `:quic`-level test pins the load-
|
|
||||||
bearing invariant.
|
|
||||||
|
|
||||||
## Why deferred
|
|
||||||
|
|
||||||
The change is small in lines but touches the QUIC writer's hot
|
|
||||||
path. Rebasing a future `:quic` retransmit / pacing change onto a
|
|
||||||
priority-sorted iteration order is doable but the diff conflicts
|
|
||||||
get noisy. Better landed as a focused PR after the current
|
|
||||||
interop-work series stabilises.
|
|
||||||
|
|
||||||
Risk profile:
|
|
||||||
- **Bugs**: subtle starvation risk if a high-priority stream
|
|
||||||
always has `streamRemaining > 0` — round-robin tiebreaker within
|
|
||||||
a priority tier mitigates this for same-priority case, but the
|
|
||||||
cross-tier case needs deliberate thought (do we want strict
|
|
||||||
priority or weighted?).
|
|
||||||
- **Performance**: per-pass sort allocation. Negligible for N≤10
|
|
||||||
but worth measuring if N grows.
|
|
||||||
- **Compat**: streams without an explicit priority default to 0,
|
|
||||||
matching today's behaviour. Existing tests should pass unchanged.
|
|
||||||
|
|
||||||
## When to land
|
|
||||||
|
|
||||||
After the catalog interop series is fully verified (production
|
|
||||||
audio with no degradation under realistic conditions), and ideally
|
|
||||||
alongside a `:quic` perf review pass that touches the same code
|
|
||||||
path. Not blocking on any audio bug today — `T11`'s reliable-
|
|
||||||
delivery fix is the actual correctness change; this is the
|
|
||||||
spec-aligned hardening.
|
|
||||||
+7
@@ -1037,6 +1037,13 @@ class MoqLiteSession internal constructor(
|
|||||||
// 150 ms late still falls inside hang's default ~200 ms
|
// 150 ms late still falls inside hang's default ~200 ms
|
||||||
// jitter buffer) and avoids the dropout entirely.
|
// jitter buffer) and avoids the dropout entirely.
|
||||||
val uni = transport.openUniStream()
|
val uni = transport.openUniStream()
|
||||||
|
// Mirror moq-rs `Publisher::serve_group`
|
||||||
|
// (`rs/moq-lite/src/lite/publisher.rs`): newer groups get higher
|
||||||
|
// priority so the QUIC writer drains fresh audio first when
|
||||||
|
// retransmits queue up under loss. Saturating cast guards a
|
||||||
|
// theoretical broadcast that runs long enough for `sequence` to
|
||||||
|
// exceed Int.MAX_VALUE (≈ 71 years at 1 group/sec); defensive only.
|
||||||
|
uni.setPriority(sequence.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||||
uni.write(Varint.encode(MoqLiteDataType.Group.code))
|
uni.write(Varint.encode(MoqLiteDataType.Group.code))
|
||||||
uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
|
uni.write(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
|
||||||
return uni
|
return uni
|
||||||
|
|||||||
+4
@@ -162,6 +162,8 @@ class FakeBidiStream internal constructor(
|
|||||||
override suspend fun finish() {
|
override suspend fun finish() {
|
||||||
write.close()
|
write.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun setPriority(priority: Int) = Unit
|
||||||
}
|
}
|
||||||
|
|
||||||
class FakeReadStream internal constructor(
|
class FakeReadStream internal constructor(
|
||||||
@@ -186,4 +188,6 @@ private class ChannelWriteStream(
|
|||||||
override suspend fun finish() {
|
override suspend fun finish() {
|
||||||
channel.close()
|
channel.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun setPriority(priority: Int) = Unit
|
||||||
}
|
}
|
||||||
|
|||||||
+16
@@ -112,6 +112,22 @@ interface WebTransportWriteStream {
|
|||||||
|
|
||||||
/** Half-close the write side (FIN). No further writes after this call. */
|
/** Half-close the write side (FIN). No further writes after this call. */
|
||||||
suspend fun finish()
|
suspend fun finish()
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Hint to the transport about this stream's drain priority relative
|
||||||
|
* to other streams on the same session. Higher value drains first
|
||||||
|
* under congestion; same-priority streams keep round-robin order.
|
||||||
|
* Default 0 = unchanged round-robin behaviour.
|
||||||
|
*
|
||||||
|
* moq-lite uses this to bias the writer toward newer group streams
|
||||||
|
* (sequence-numbered, fresher audio) so a backlog of retransmits on
|
||||||
|
* an older group doesn't starve the listener of fresh frames. See
|
||||||
|
* `Publisher::serve_group` in `rs/moq-lite/src/lite/publisher.rs`.
|
||||||
|
*
|
||||||
|
* Implementations that don't model priority (e.g. the in-memory
|
||||||
|
* fake) MAY treat this as a no-op.
|
||||||
|
*/
|
||||||
|
fun setPriority(priority: Int)
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
+18
@@ -365,6 +365,10 @@ private class QuicBidiStreamAdapter(
|
|||||||
stream.send.finish()
|
stream.send.finish()
|
||||||
driver.wakeup()
|
driver.wakeup()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun setPriority(priority: Int) {
|
||||||
|
stream.priority = priority
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private class QuicReadStreamAdapter(
|
private class QuicReadStreamAdapter(
|
||||||
@@ -392,6 +396,10 @@ private class QuicUniWriteStreamAdapter(
|
|||||||
stream.send.finish()
|
stream.send.finish()
|
||||||
driver.wakeup()
|
driver.wakeup()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override fun setPriority(priority: Int) {
|
||||||
|
stream.priority = priority
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Adapter for a WT peer-initiated uni stream whose prefix has been stripped. */
|
/** Adapter for a WT peer-initiated uni stream whose prefix has been stripped. */
|
||||||
@@ -431,4 +439,14 @@ private class StrippedWtBidiStreamAdapter(
|
|||||||
?: error("peer-initiated bidi stream has no finish — demux didn't wire one")
|
?: error("peer-initiated bidi stream has no finish — demux didn't wire one")
|
||||||
finish()
|
finish()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* No-op: peer-initiated bidi streams arrive through the demux as a
|
||||||
|
* [com.vitorpamplona.quic.webtransport.StrippedWtStream] which exposes
|
||||||
|
* only `send`/`finish` closures, not the underlying [QuicStream]. The
|
||||||
|
* moq-lite priority use case targets locally-opened uni group streams
|
||||||
|
* only, so this path doesn't need to model priority — see the
|
||||||
|
* [WebTransportWriteStream.setPriority] contract.
|
||||||
|
*/
|
||||||
|
override fun setPriority(priority: Int) = Unit
|
||||||
}
|
}
|
||||||
|
|||||||
+66
-38
@@ -493,46 +493,74 @@ private fun buildApplicationPacket(
|
|||||||
// insertion-ordered and stays in sync with the streams map.
|
// insertion-ordered and stays in sync with the streams map.
|
||||||
val streamsView = conn.streamsListLocked()
|
val streamsView = conn.streamsListLocked()
|
||||||
if (streamsView.isNotEmpty()) {
|
if (streamsView.isNotEmpty()) {
|
||||||
val start = conn.streamRoundRobinStart % streamsView.size
|
// Strict priority across tiers, round-robin within each tier.
|
||||||
for (i in streamsView.indices) {
|
// Higher-priority streams (e.g. moq-lite newer-sequence group
|
||||||
if (packetBudget <= 64) break
|
// streams) ALWAYS drain ahead of lower-priority ones; the
|
||||||
val stream = streamsView[(start + i) % streamsView.size]
|
// rotating start-index only rotates among same-priority peers.
|
||||||
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
|
// This is the spec-aligned shape — applying the rotation
|
||||||
// Skip if both stream and connection have no credit; FIN-only
|
// globally over the sorted list would flip cross-tier order on
|
||||||
// (zero-byte) chunks may still go through because they don't
|
// alternating drains, defeating the priority hint entirely.
|
||||||
// consume credit.
|
//
|
||||||
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
|
// Default priority is 0; if every stream is at the default, all
|
||||||
val effectiveCap = minOf(streamRemaining, connBudget)
|
// streams form a single tier and iteration order matches the
|
||||||
val maxBytes =
|
// pre-priority round-robin behaviour exactly.
|
||||||
minOf(packetBudget - 32, effectiveCap.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
//
|
||||||
val chunk = stream.send.takeChunk(maxBytes = maxBytes) ?: continue
|
// Cost: O(N log N) per drain pass plus one transient sorted
|
||||||
if (chunk.data.isNotEmpty() || chunk.fin) {
|
// list. N is small (1–10 in the moq-lite audio path); if it
|
||||||
frames +=
|
// ever grows enough to matter, switch to an indirect index sort
|
||||||
StreamFrame(
|
// or maintain an incrementally-sorted view on setPriority.
|
||||||
streamId = stream.streamId,
|
val sorted =
|
||||||
offset = chunk.offset,
|
if (streamsView.size > 1) streamsView.sortedByDescending { it.priority } else streamsView
|
||||||
data = chunk.data,
|
val rotation = conn.streamRoundRobinStart
|
||||||
fin = chunk.fin,
|
var tierStart = 0
|
||||||
explicitLength = true,
|
outer@ while (tierStart < sorted.size) {
|
||||||
)
|
// Walk the contiguous run of same-priority streams.
|
||||||
// Step C of the deferred-follow-ups pass: track this
|
val tierPriority = sorted[tierStart].priority
|
||||||
// STREAM emission so RFC 9002 retransmit can re-queue
|
var tierEnd = tierStart + 1
|
||||||
// the byte range on loss. SendBuffer.markLost (commit B)
|
while (tierEnd < sorted.size && sorted[tierEnd].priority == tierPriority) tierEnd++
|
||||||
// moves the range from in-flight back to the retransmit
|
val tierSize = tierEnd - tierStart
|
||||||
// queue, and the next takeChunk replays it.
|
val tierRotation = if (tierSize > 1) rotation % tierSize else 0
|
||||||
tokens +=
|
for (k in 0 until tierSize) {
|
||||||
RecoveryToken.Stream(
|
if (packetBudget <= 64) break@outer
|
||||||
streamId = stream.streamId,
|
val stream = sorted[tierStart + ((tierRotation + k) % tierSize)]
|
||||||
offset = chunk.offset,
|
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
|
||||||
length = chunk.data.size.toLong(),
|
// Skip if both stream and connection have no credit; FIN-only
|
||||||
fin = chunk.fin,
|
// (zero-byte) chunks may still go through because they don't
|
||||||
)
|
// consume credit.
|
||||||
packetBudget -= chunk.data.size + 32
|
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
|
||||||
connBudget -= chunk.data.size
|
val effectiveCap = minOf(streamRemaining, connBudget)
|
||||||
conn.sendConnectionFlowConsumed += chunk.data.size
|
val maxBytes =
|
||||||
|
minOf(packetBudget - 32, effectiveCap.coerceAtMost(Int.MAX_VALUE.toLong()).toInt())
|
||||||
|
val chunk = stream.send.takeChunk(maxBytes = maxBytes) ?: continue
|
||||||
|
if (chunk.data.isNotEmpty() || chunk.fin) {
|
||||||
|
frames +=
|
||||||
|
StreamFrame(
|
||||||
|
streamId = stream.streamId,
|
||||||
|
offset = chunk.offset,
|
||||||
|
data = chunk.data,
|
||||||
|
fin = chunk.fin,
|
||||||
|
explicitLength = true,
|
||||||
|
)
|
||||||
|
// Step C of the deferred-follow-ups pass: track this
|
||||||
|
// STREAM emission so RFC 9002 retransmit can re-queue
|
||||||
|
// the byte range on loss. SendBuffer.markLost (commit B)
|
||||||
|
// moves the range from in-flight back to the retransmit
|
||||||
|
// queue, and the next takeChunk replays it.
|
||||||
|
tokens +=
|
||||||
|
RecoveryToken.Stream(
|
||||||
|
streamId = stream.streamId,
|
||||||
|
offset = chunk.offset,
|
||||||
|
length = chunk.data.size.toLong(),
|
||||||
|
fin = chunk.fin,
|
||||||
|
)
|
||||||
|
packetBudget -= chunk.data.size + 32
|
||||||
|
connBudget -= chunk.data.size
|
||||||
|
conn.sendConnectionFlowConsumed += chunk.data.size
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
tierStart = tierEnd
|
||||||
}
|
}
|
||||||
conn.streamRoundRobinStart = (start + 1) % streamsView.size
|
conn.streamRoundRobinStart = (rotation + 1) % streamsView.size
|
||||||
}
|
}
|
||||||
|
|
||||||
if (frames.isEmpty()) return null
|
if (frames.isEmpty()) return null
|
||||||
|
|||||||
@@ -47,6 +47,27 @@ class QuicStream(
|
|||||||
val send = SendBuffer(bestEffort = bestEffort)
|
val send = SendBuffer(bestEffort = bestEffort)
|
||||||
val receive = ReceiveBuffer()
|
val receive = ReceiveBuffer()
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Send-side scheduling priority. The connection writer's drain loop
|
||||||
|
* iterates streams by descending priority; same-priority streams keep
|
||||||
|
* their existing round-robin order. Higher value = drains first under
|
||||||
|
* congestion. Default 0 matches pre-priority round-robin behaviour for
|
||||||
|
* every existing call site.
|
||||||
|
*
|
||||||
|
* Used by moq-lite group streams: the publisher assigns each new group
|
||||||
|
* a priority equal to its sequence number so that newer groups
|
||||||
|
* (fresher audio) drain ahead of older ones when retransmits queue up
|
||||||
|
* on a lossy link. Mirrors `Publisher::serve_group` in
|
||||||
|
* `rs/moq-lite/src/lite/publisher.rs` (`stream.set_priority`).
|
||||||
|
*
|
||||||
|
* `@Volatile` because callers (e.g. moq-lite's openGroupStream)
|
||||||
|
* assign from arbitrary coroutines while the writer reads it under
|
||||||
|
* [com.vitorpamplona.quic.connection.QuicConnection.lock] during a
|
||||||
|
* drain pass.
|
||||||
|
*/
|
||||||
|
@Volatile
|
||||||
|
var priority: Int = 0
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Bytes received and confirmed contiguous, exposed as a flow to the consumer.
|
* Bytes received and confirmed contiguous, exposed as a flow to the consumer.
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -327,6 +327,46 @@ class InMemoryQuicPipe(
|
|||||||
*/
|
*/
|
||||||
fun buildServerApplicationPacket(frames: List<com.vitorpamplona.quic.frame.Frame>): ByteArray? = buildServerApplicationDatagram(frames)
|
fun buildServerApplicationPacket(frames: List<com.vitorpamplona.quic.frame.Frame>): ByteArray? = buildServerApplicationDatagram(frames)
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Decrypt the application-level (1-RTT) packet inside a client-emitted
|
||||||
|
* datagram and return its frames in wire order. Test-only helper for
|
||||||
|
* assertions that depend on per-frame ordering inside a packet (e.g.
|
||||||
|
* stream priority scheduling). Walks past any coalesced long-header
|
||||||
|
* packets (Initial / Handshake) at the front of the datagram, since the
|
||||||
|
* client may still flush ACKs at those levels post-handshake. Returns
|
||||||
|
* null if no short-header packet is present or decryption fails.
|
||||||
|
*/
|
||||||
|
fun decryptClientApplicationFrames(datagram: ByteArray): List<com.vitorpamplona.quic.frame.Frame>? {
|
||||||
|
if (datagram.isEmpty()) return null
|
||||||
|
var offset = 0
|
||||||
|
while (offset < datagram.size) {
|
||||||
|
val first = datagram[offset].toInt() and 0xFF
|
||||||
|
if ((first and 0x80) == 0) {
|
||||||
|
val proto = serverApplicationRx ?: return null
|
||||||
|
val parsed =
|
||||||
|
ShortHeaderPacket.parseAndDecrypt(
|
||||||
|
bytes = datagram,
|
||||||
|
offset = offset,
|
||||||
|
dcidLen = serverScid.length,
|
||||||
|
aead = proto.aead,
|
||||||
|
key = proto.key,
|
||||||
|
iv = proto.iv,
|
||||||
|
hp = proto.hp,
|
||||||
|
hpKey = proto.hpKey,
|
||||||
|
largestReceivedInSpace = applicationPnSpace.largestReceived,
|
||||||
|
) ?: return null
|
||||||
|
applicationPnSpace.observeInbound(parsed.packet.packetNumber, 0L)
|
||||||
|
return decodeFrames(parsed.packet.payload)
|
||||||
|
}
|
||||||
|
// Long header — skip past it using the encoded length field so
|
||||||
|
// we can inspect any short-header packet that was coalesced
|
||||||
|
// after it.
|
||||||
|
val peeked = LongHeaderPacket.peekHeader(datagram, offset) ?: return null
|
||||||
|
offset += peeked.totalLength
|
||||||
|
}
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
private fun buildServerInitialPacket(crypto: ByteArray): ByteArray {
|
private fun buildServerInitialPacket(crypto: ByteArray): ByteArray {
|
||||||
val proto =
|
val proto =
|
||||||
PacketProtection(
|
PacketProtection(
|
||||||
|
|||||||
+80
@@ -110,6 +110,86 @@ class QuicConnectionWriterTest {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun writer_drains_higher_priority_streams_before_lower_priority() {
|
||||||
|
// T11.3 follow-up: priority must dominate iteration order on
|
||||||
|
// EVERY drain, not just the first. The naive "sort by priority
|
||||||
|
// then apply the existing rotating start globally" shape looks
|
||||||
|
// right at a glance but the rotating start advances on every
|
||||||
|
// drain, so cross-tier ordering flips on alternate drains and
|
||||||
|
// the priority hint is silently defeated. The correct shape is
|
||||||
|
// strict priority across tiers + round-robin only within a
|
||||||
|
// tier, which we pin here by draining TWICE and asserting the
|
||||||
|
// higher-priority stream's StreamFrame lands first in BOTH
|
||||||
|
// packets. Low-priority stream is opened FIRST so insertion-
|
||||||
|
// order can't accidentally pass for priority ordering.
|
||||||
|
runBlocking {
|
||||||
|
val (client, pipe) = connectedClient()
|
||||||
|
val low = client.openBidiStream()
|
||||||
|
val high = client.openBidiStream()
|
||||||
|
low.priority = 0
|
||||||
|
high.priority = 10
|
||||||
|
|
||||||
|
repeat(2) { round ->
|
||||||
|
low.send.enqueue(ByteArray(200) { 0xAA.toByte() })
|
||||||
|
high.send.enqueue(ByteArray(200) { 0xBB.toByte() })
|
||||||
|
|
||||||
|
val datagram = drainOutbound(client, nowMillis = 0L)
|
||||||
|
assertNotNull(datagram, "drain on round $round must emit a packet")
|
||||||
|
val frames = pipe.decryptClientApplicationFrames(datagram)
|
||||||
|
assertNotNull(frames, "decrypt must succeed on round $round")
|
||||||
|
val streamFrames = frames.filterIsInstance<StreamFrame>()
|
||||||
|
assertEquals(
|
||||||
|
2,
|
||||||
|
streamFrames.size,
|
||||||
|
"expected one StreamFrame per stream on round $round, got $streamFrames",
|
||||||
|
)
|
||||||
|
assertEquals(
|
||||||
|
high.streamId,
|
||||||
|
streamFrames[0].streamId,
|
||||||
|
"higher-priority stream must drain first on round $round; " +
|
||||||
|
"saw ${streamFrames.map { it.streamId }}",
|
||||||
|
)
|
||||||
|
assertEquals(low.streamId, streamFrames[1].streamId, "low second on round $round")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun writer_round_robins_within_a_priority_tier() {
|
||||||
|
// Regression guard for the tier-local round-robin: same-priority
|
||||||
|
// streams must still rotate so an early-opened stream doesn't
|
||||||
|
// monopolise a packet's stream-frame slot indefinitely. We open
|
||||||
|
// three streams at the default (0) priority and verify the
|
||||||
|
// rotating start advances by one per drain, matching the
|
||||||
|
// pre-priority behaviour.
|
||||||
|
runBlocking {
|
||||||
|
val (client, pipe) = connectedClient()
|
||||||
|
val a = client.openBidiStream()
|
||||||
|
val b = client.openBidiStream()
|
||||||
|
val c = client.openBidiStream()
|
||||||
|
// All default priority — single tier, three streams.
|
||||||
|
val expectedRotation =
|
||||||
|
listOf(
|
||||||
|
listOf(a.streamId, b.streamId, c.streamId),
|
||||||
|
listOf(b.streamId, c.streamId, a.streamId),
|
||||||
|
listOf(c.streamId, a.streamId, b.streamId),
|
||||||
|
)
|
||||||
|
for ((round, expected) in expectedRotation.withIndex()) {
|
||||||
|
a.send.enqueue(ByteArray(64) { 0xA1.toByte() })
|
||||||
|
b.send.enqueue(ByteArray(64) { 0xB2.toByte() })
|
||||||
|
c.send.enqueue(ByteArray(64) { 0xC3.toByte() })
|
||||||
|
|
||||||
|
val datagram = drainOutbound(client, nowMillis = 0L)
|
||||||
|
assertNotNull(datagram, "drain $round must emit a packet")
|
||||||
|
val frames = pipe.decryptClientApplicationFrames(datagram)
|
||||||
|
assertNotNull(frames, "decrypt must succeed on round $round")
|
||||||
|
val ids = frames.filterIsInstance<StreamFrame>().map { it.streamId }
|
||||||
|
assertEquals(expected, ids, "round $round round-robin order")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun writer_respects_connection_level_send_credit_cap() {
|
fun writer_respects_connection_level_send_credit_cap() {
|
||||||
// Audit-4 #9: pre-fix the writer ignored sendConnectionFlowCredit
|
// Audit-4 #9: pre-fix the writer ignored sendConnectionFlowCredit
|
||||||
|
|||||||
Reference in New Issue
Block a user