Merge remote-tracking branch 'origin/main' into claude/fix-green-circle-ui-H7zsK

# Conflicts:
#	commons/src/commonMain/kotlin/com/vitorpamplona/amethyst/commons/viewmodels/NestViewModel.kt
This commit is contained in:
Claude
2026-05-06 22:23:33 +00:00
11 changed files with 564 additions and 264 deletions
@@ -302,14 +302,26 @@ class NestViewModel(
/**
* [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.
* detector triggered, or `null` if it has never fired. Combined
* with [consecutiveCliffRecycles] to drive the per-attempt
* backoff schedule in [computeStalledSpeakers] — short backoff
* after the first failed recycle so a single moq-rs cliff
* 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
/**
* 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
* background `subscribeCatalog` collector launched in
@@ -320,7 +332,19 @@ class NestViewModel(
*/
private val catalogJobs = mutableMapOf<String, Job>()
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
private var speaker: NestsSpeaker? = null
@@ -1389,6 +1413,7 @@ class NestViewModel(
cliffDetectorJob = null
lastFrameAt.clear()
lastCliffRecycleAt = null
consecutiveCliffRecycles = 0
// 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
@@ -1491,6 +1516,11 @@ class NestViewModel(
lastFrameAt[pubkey] =
kotlin.time.TimeSource.Monotonic
.markNow()
// 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
// overlay. Subsequent frames are no-ops here.
if (_uiState.value.connectingSpeakers.contains(pubkey)) {
@@ -1588,43 +1618,53 @@ class NestViewModel(
announcedSpeakers = announced,
lastFrameAt = lastFrameAt,
lastRecycleAt = lastCliffRecycleAt,
consecutiveFailedRecycles = consecutiveCliffRecycles,
cliffTimeoutMs = ROOM_AUDIO_CLIFF_TIMEOUT_MS,
cooldownMs = ROOM_AUDIO_CLIFF_COOLDOWN_MS,
postRecycleGraceMs = ROOM_AUDIO_CLIFF_RECYCLE_GRACE_MS,
)
if (stalled.isEmpty()) continue
consecutiveCliffRecycles += 1
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") {
"\"timeout_ms\":$ROOM_AUDIO_CLIFF_TIMEOUT_MS," +
"\"consecutive\":$consecutiveCliffRecycles," +
"\"stalled\":${com.vitorpamplona.nestsclient.trace.jsonArrStr(stalled)}"
}
val recycleMark =
lastCliffRecycleAt =
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
// Note: we deliberately do NOT overwrite
// `lastFrameAt[pubkey]` with the recycle moment.
// Keeping the real-frame timestamp lets the next
// tick distinguish "the recycle delivered audio,
// we then re-stalled" (lastFrameAt advances past
// lastCliffRecycleAt — counter resets in
// `onSpeakerActivity`, attempt = 0, fire-eligible
// immediately) from "the recycle did nothing,
// still no frames" (lastFrameAt unchanged,
// counter keeps growing, backoff escalates).
// The previous reset collapsed both cases into
// a flat 30 s wait, which made a single failed
// recycle sound like a 30 s+ dropout to the
// user even when retrying earlier would have
// recovered.
try {
l.recycleSession()
} catch (ce: CancellationException) {
// User left mid-recycle — exit promptly
// rather than letting the next delay() be
// 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 {
com.vitorpamplona.quartz.utils.Log
@@ -1975,24 +2015,34 @@ const val ROOM_AUDIO_CLIFF_TIMEOUT_MS: Long = 2_500L
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.
* Post-recycle handshake grace. NEVER recycle again within this
* window of the previous recycle the wrapper is still tearing
* down + reopening the QUIC session, so a "no frames since recycle"
* reading is just the handshake, not a relay-side cliff.
*
* 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.
* 3 s covers a typical re-handshake (drain old session UNSUB +
* UNANNOUNCE + WT_CLOSE + open new WebTransport + SUBSCRIBE +
* SUBSCRIBE_OK). Anything shorter risks recycling inside our own
* handshake window; longer wastes audible silence in the recovering-
* within-grace case.
*/
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
@@ -2003,6 +2053,36 @@ const val ROOM_AUDIO_CLIFF_COOLDOWN_MS: Long = 30_000L
*/
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
* exercise it with a [kotlin.time.TestTimeSource] without standing up
@@ -2023,23 +2103,30 @@ private const val CLIFF_DIAG_LOG_EVERY: Long = 5L
* - 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).
* Suppression (when [lastRecycleAt] is non-null):
* - while elapsed since [lastRecycleAt] is below [postRecycleGraceMs],
* never recycle — we're still inside the previous handshake window.
* - 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(
activeSpeakers: Set<String>,
announcedSpeakers: Set<String>,
lastFrameAt: Map<String, kotlin.time.TimeMark>,
lastRecycleAt: kotlin.time.TimeMark?,
consecutiveFailedRecycles: Int = 0,
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> {
if (lastRecycleAt != null &&
lastRecycleAt.elapsedNow().inWholeMilliseconds < cooldownMs
) {
return emptyList()
if (lastRecycleAt != null) {
val sinceRecycleMs = lastRecycleAt.elapsedNow().inWholeMilliseconds
if (sinceRecycleMs < postRecycleGraceMs) return emptyList()
if (sinceRecycleMs < backoffForAttempt(consecutiveFailedRecycles)) return emptyList()
}
return activeSpeakers
.asSequence()
@@ -189,21 +189,16 @@ class CliffDetectorTest {
}
@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.
fun postRecycleGraceSuppressesEvenAtAttemptZero() {
// Inside the 3 s post-recycle handshake window, never recycle
// — the wrapper is still tearing down + reopening QUIC, so
// "no frames since recycle" is just the handshake, not a
// relay-side cliff.
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
ts += 4_500.milliseconds
val recycleMark = ts.markNow()
ts += 2_000.milliseconds // inside 3 s grace
val result =
computeStalledSpeakers(
@@ -211,21 +206,22 @@ class CliffDetectorTest {
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = recycleMark,
consecutiveFailedRecycles = 0,
)
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.
fun attemptZeroFiresImmediatelyOnceGracePasses() {
// After a recovered-then-restall pattern: counter has been
// reset by `onSpeakerActivity`, so even though there's a
// prior recycleMark, the next cliff fires as soon as the
// 3 s grace passes — no extended backoff.
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
ts += 3_500.milliseconds // past 3 s grace, attempt=0 schedule = 0 ms
val result =
computeStalledSpeakers(
@@ -233,10 +229,164 @@ class CliffDetectorTest {
announcedSpeakers = setOf(ALICE),
lastFrameAt = mapOf(ALICE to frameMark),
lastRecycleAt = recycleMark,
consecutiveFailedRecycles = 0,
)
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
fun customTimeoutsAreHonored() {
// Defaults are wired into the production VM, but the function
@@ -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 (110 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.
@@ -1037,6 +1037,13 @@ class MoqLiteSession internal constructor(
// 150 ms late still falls inside hang's default ~200 ms
// jitter buffer) and avoids the dropout entirely.
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(MoqLiteCodec.encodeGroupHeader(MoqLiteGroupHeader(subscribeId, sequence)))
return uni
@@ -162,6 +162,8 @@ class FakeBidiStream internal constructor(
override suspend fun finish() {
write.close()
}
override fun setPriority(priority: Int) = Unit
}
class FakeReadStream internal constructor(
@@ -186,4 +188,6 @@ private class ChannelWriteStream(
override suspend fun finish() {
channel.close()
}
override fun setPriority(priority: Int) = Unit
}
@@ -112,6 +112,22 @@ interface WebTransportWriteStream {
/** Half-close the write side (FIN). No further writes after this call. */
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)
}
/**
@@ -365,6 +365,10 @@ private class QuicBidiStreamAdapter(
stream.send.finish()
driver.wakeup()
}
override fun setPriority(priority: Int) {
stream.priority = priority
}
}
private class QuicReadStreamAdapter(
@@ -392,6 +396,10 @@ private class QuicUniWriteStreamAdapter(
stream.send.finish()
driver.wakeup()
}
override fun setPriority(priority: Int) {
stream.priority = priority
}
}
/** 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")
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
}
@@ -410,46 +410,74 @@ private fun buildApplicationPacket(
// insertion-ordered and stays in sync with the streams map.
val streamsView = conn.streamsListLocked()
if (streamsView.isNotEmpty()) {
val start = conn.streamRoundRobinStart % streamsView.size
for (i in streamsView.indices) {
if (packetBudget <= 64) break
val stream = streamsView[(start + i) % streamsView.size]
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
// Skip if both stream and connection have no credit; FIN-only
// (zero-byte) chunks may still go through because they don't
// consume credit.
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
val effectiveCap = minOf(streamRemaining, connBudget)
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
// Strict priority across tiers, round-robin within each tier.
// Higher-priority streams (e.g. moq-lite newer-sequence group
// streams) ALWAYS drain ahead of lower-priority ones; the
// rotating start-index only rotates among same-priority peers.
// This is the spec-aligned shape — applying the rotation
// globally over the sorted list would flip cross-tier order on
// alternating drains, defeating the priority hint entirely.
//
// Default priority is 0; if every stream is at the default, all
// streams form a single tier and iteration order matches the
// pre-priority round-robin behaviour exactly.
//
// Cost: O(N log N) per drain pass plus one transient sorted
// list. N is small (110 in the moq-lite audio path); if it
// ever grows enough to matter, switch to an indirect index sort
// or maintain an incrementally-sorted view on setPriority.
val sorted =
if (streamsView.size > 1) streamsView.sortedByDescending { it.priority } else streamsView
val rotation = conn.streamRoundRobinStart
var tierStart = 0
outer@ while (tierStart < sorted.size) {
// Walk the contiguous run of same-priority streams.
val tierPriority = sorted[tierStart].priority
var tierEnd = tierStart + 1
while (tierEnd < sorted.size && sorted[tierEnd].priority == tierPriority) tierEnd++
val tierSize = tierEnd - tierStart
val tierRotation = if (tierSize > 1) rotation % tierSize else 0
for (k in 0 until tierSize) {
if (packetBudget <= 64) break@outer
val stream = sorted[tierStart + ((tierRotation + k) % tierSize)]
val streamRemaining = (stream.sendCredit - stream.send.sentOffset).coerceAtLeast(0L)
// Skip if both stream and connection have no credit; FIN-only
// (zero-byte) chunks may still go through because they don't
// consume credit.
if (streamRemaining <= 0L && connBudget <= 0L && !stream.send.finPending) continue
val effectiveCap = minOf(streamRemaining, connBudget)
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
@@ -47,6 +47,27 @@ class QuicStream(
val send = SendBuffer(bestEffort = bestEffort)
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.
*
@@ -327,6 +327,46 @@ class InMemoryQuicPipe(
*/
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 {
val proto =
PacketProtection(
@@ -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
fun writer_respects_connection_level_send_credit_cap() {
// Audit-4 #9: pre-fix the writer ignored sendConnectionFlowCredit