feat(android): Sonos/UPnP queue plays leveled URLs, rendered a track ahead (M464 #5002)
release / go (push) Successful in 2m49s
release / web (push) Successful in 2m21s
release / govulncheck (push) Successful in 25s
release / integration (push) Successful in 5m54s
release / android (push) Successful in 8m10s
release / Build signed APK (releases and dev) (push) Successful in 8m35s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 1m42s
release / Verify release artifacts (tag releases only) (push) Skipped

Every URL the Sonos queue loader sends is minted with level=true and
the track's album-play verdict from its neighbours in the queue; the
server returns the plain stream when leveling is off or changes
nothing. The playing track and the one after it are rendered ahead,
and each time the renderer moves on, the next is.

Server: a mint no longer prerenders on its own. A queue load mints
every track, which would have started an ffmpeg render per track at
once. The request now carries prerender, and at most two prerenders
run at a time; past that they are dropped, since a fetch renders on
demand anyway.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-06 19:38:21 -04:00
co-authored by Claude Opus 5.5
parent f34423a0e0
commit 2e36e70268
9 changed files with 216 additions and 38 deletions
@@ -29,11 +29,20 @@ interface CastApi {
* Request body. [expSeconds] is clamped server-side to [60, 86400];
* the 21_600 default (6h) is long enough to play through any typical
* track without re-minting mid-playback.
*
* [level] asks for the leveled stream (M464 #5001): the track rendered at
* the user's loudness gain, which the server works out from their setting.
* [asAlbum] says the track plays among its album in order, which picks
* album gain in auto mode. [prerender] says the speaker will fetch it
* soon, so the server renders it ahead.
*/
@Serializable
data class StreamTokenRequest(
val trackId: String,
val expSeconds: Int = 21_600,
val level: Boolean = false,
val asAlbum: Boolean = false,
val prerender: Boolean = false,
)
/**
@@ -53,4 +62,6 @@ data class StreamTokenResponse(
val url: String,
val mime: String = "audio/mpeg",
val title: String = "",
/** [url] is the leveled stream; false when leveling is off or changes nothing. */
val leveled: Boolean = false,
)
@@ -93,6 +93,8 @@ class MinstrelForwardingPlayer(
* one. Diagnostics-only; see [TransportObservation].
*/
val onTransport: (TransportObservation) -> Unit = {},
/** The renderer's 1-based queue position, every poll. */
val onRendererTrack: (trackNumber: Int) -> Unit = {},
)
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
@@ -622,6 +624,7 @@ class MinstrelForwardingPlayer(
trackNumber = info.track,
)
syncLocalCursorToRemote(sonosTrack = info.track, trackUri = info.trackUri)
events.onRendererTrack(info.track)
val transport = active.avTransport.getTransportInfo()
when (transport.state) {
TransportState.PLAYING -> {
@@ -24,6 +24,7 @@ import com.fabledsword.minstrel.cache.audiocache.CacheConfig
import com.fabledsword.minstrel.player.gain.GainAudioProcessor
import com.fabledsword.minstrel.player.gain.ReplayGainStore
import com.fabledsword.minstrel.player.output.ActiveUpnpHolder
import com.fabledsword.minstrel.player.output.SonosQueueLoader
import dagger.hilt.android.qualifiers.ApplicationContext
import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.flow.MutableSharedFlow
@@ -65,6 +66,7 @@ class PlayerFactory @Inject constructor(
private val serverHealth: com.fabledsword.minstrel.connectivity.NetworkStatusController,
private val authStore: AuthStore,
private val replayGains: ReplayGainStore,
private val sonosQueue: SonosQueueLoader,
) {
private val cacheDir: File = File(context.cacheDir, "audio_cache").apply { mkdirs() }
@@ -129,6 +131,7 @@ class PlayerFactory @Inject constructor(
onStalled = { trackId -> stallEventsInternal.tryEmit(trackId) },
onQueueTruncated = { queueRepairInternal.tryEmit(Unit) },
onTransport = { transportInternal.tryEmit(it) },
onRendererTrack = { sonosQueue.onRendererTrack(it) },
),
)
}
@@ -19,6 +19,13 @@ import javax.inject.Singleton
class StreamTokenProvider @Inject constructor(retrofit: Retrofit) {
private val api: CastApi = retrofit.create()
suspend fun mint(trackId: String): StreamTokenResponse =
api.streamToken(StreamTokenRequest(trackId = trackId))
/**
* Mints a URL for a speaker, leveled when the user's setting calls for
* it (M464 #5002). The server answers with the plain stream when it does
* not, so every speaker URL asks.
*/
suspend fun mint(trackId: String, asAlbum: Boolean = false, prerender: Boolean = false): StreamTokenResponse =
api.streamToken(
StreamTokenRequest(trackId = trackId, level = true, asAlbum = asAlbum, prerender = prerender),
)
}
@@ -4,6 +4,8 @@ import com.fabledsword.minstrel.di.ApplicationScope
import com.fabledsword.minstrel.models.TrackRef
import com.fabledsword.minstrel.player.RemotePlayerState
import com.fabledsword.minstrel.player.StreamTokenProvider
import com.fabledsword.minstrel.player.gain.AlbumPosition
import com.fabledsword.minstrel.player.gain.GainMath
import com.fabledsword.minstrel.player.output.upnp.AVTransportClient
import com.fabledsword.minstrel.player.output.upnp.SoapFaultException
import com.fabledsword.minstrel.player.output.upnp.bareUdn
@@ -24,6 +26,11 @@ import javax.inject.Singleton
* with its own failure modes, and it had grown large enough to hide one:
* every write here is a SOAP call that can fail individually, and until
* [verifyQueueLength] nothing ever read the result back.
*
* Every URL sent is a leveled one (M464 #5002): the server renders the track
* at the user's loudness gain, or hands back the plain stream when leveling
* is off. Renders are slow enough to matter, so the track playing and the
* one after it are rendered ahead; the rest render when the speaker asks.
*/
@Singleton
class SonosQueueLoader @Inject constructor(
@@ -32,6 +39,13 @@ class SonosQueueLoader @Inject constructor(
private val activeUpnpHolder: ActiveUpnpHolder,
private val remoteState: RemotePlayerState,
) {
// The queue as last sent to the renderer, for looking up the track after
// the one it is playing.
@Volatile private var sent: List<TrackRef> = emptyList()
// The renderer track the next one was last prerendered for.
@Volatile private var prerenderedAfter = 0
suspend fun load(
transport: AVTransportClient,
route: OutputRoute,
@@ -46,8 +60,10 @@ class SonosQueueLoader @Inject constructor(
"UPnP select: add %d initial tracks (currentIndex=%d, totalQueue=%d)",
initialBatch.size, currentIndex, queue.size,
)
initialBatch.forEachIndexed { idx, ref ->
val token = streamTokens.mint(ref.id)
sent = queue
prerenderedAfter = currentIndex + 1
initialBatch.indices.forEach { idx ->
val token = mint(queue, idx, prerender = idx == currentIndex)
transport.addURIToQueue(
uri = token.url,
mime = token.mime,
@@ -64,12 +80,12 @@ class SonosQueueLoader @Inject constructor(
Timber.w("UPnP select: Play")
transport.play()
Timber.w("UPnP select: initial done; backgrounding remainder")
val remaining = queue.drop(initialEnd)
// Verify even when there is no tail to append: the initial batch is
// sent the same way and can be dropped the same way.
scope.launch {
if (remaining.isNotEmpty()) {
extendQueueOnSonos(transport, route, remaining, initialEnd)
prerender(queue, currentIndex + 1)
if (initialEnd < queue.size) {
extendQueueOnSonos(transport, route, queue, initialEnd)
}
verifyQueueLength(transport, route, queue)
}
@@ -84,15 +100,16 @@ class SonosQueueLoader @Inject constructor(
private suspend fun extendQueueOnSonos(
transport: AVTransportClient,
route: OutputRoute,
tracks: List<TrackRef>,
queue: List<TrackRef>,
startPosition: Int,
) {
val count = queue.size - startPosition
Timber.w(
"UPnP extend: appending %d tracks starting at position %d",
tracks.size, startPosition + 1,
count, startPosition + 1,
)
val succeeded = appendTracksToQueue(transport, route, tracks, startPosition)
Timber.w("UPnP extend: done (%d / %d appended)", succeeded, tracks.size)
val succeeded = appendTracksToQueue(transport, route, queue, startPosition)
Timber.w("UPnP extend: done (%d / %d appended)", succeeded, count)
}
/**
@@ -136,18 +153,17 @@ class SonosQueueLoader @Inject constructor(
Timber.w("UPnP verify: renderer holds %d tracks, queue intact", nrTracks)
return
}
val missing = fullQueue.drop(nrTracks)
Timber.w(
"UPnP verify: %s holds %d of %d tracks; appending %d missing (round %d)",
route.name, nrTracks, fullQueue.size, missing.size, round + 1,
route.name, nrTracks, fullQueue.size, fullQueue.size - nrTracks, round + 1,
)
appendTracksToQueue(transport, route, missing, nrTracks)
appendTracksToQueue(transport, route, fullQueue, nrTracks)
}
Timber.w("UPnP verify: gave up repairing queue length on %s", route.name)
}
/**
* Append [tracks] at [startPosition] (0-based), returning how many landed.
* Append [queue] from [startPosition] (0-based) on, returning how many landed.
* Tolerates individual AddURIToQueue failures — log and continue so some
* tracks loaded is better than zero tracks loaded — and stops early after
* [EXTEND_ABORT_AFTER_FAILURES] consecutive ones.
@@ -155,20 +171,20 @@ class SonosQueueLoader @Inject constructor(
private suspend fun appendTracksToQueue(
transport: AVTransportClient,
route: OutputRoute,
tracks: List<TrackRef>,
queue: List<TrackRef>,
startPosition: Int,
): Int {
var consecutiveFailures = 0
var succeeded = 0
var aborted = false
for ((i, ref) in tracks.withIndex()) {
for (i in 0 until queue.size - startPosition) {
if (aborted) break
if (activeUpnpHolder.active.value?.routeId != route.id) {
Timber.w("UPnP extend: cancelled at offset %d (route changed)", i)
aborted = true
} else {
val outcome = runCatching {
val token = streamTokens.mint(ref.id)
val token = mint(queue, startPosition + i)
transport.addURIToQueue(
uri = token.url,
mime = token.mime,
@@ -220,6 +236,7 @@ class SonosQueueLoader @Inject constructor(
newQueue: List<TrackRef>,
): Boolean {
val newIds = newQueue.map { it.id }
sent = newQueue
if (oldIds == newIds) return true
val prefixLen = commonPrefixLength(oldIds, newIds)
val suffixLen = commonSuffixLength(
@@ -271,8 +288,7 @@ class SonosQueueLoader @Inject constructor(
prefixLen + 1,
)
for (i in 0 until addedCount) {
val ref = newQueue[prefixLen + i]
val token = streamTokens.mint(ref.id)
val token = mint(newQueue, prefixLen + i)
transport.addURIToQueue(
uri = token.url,
mime = token.mime,
@@ -283,6 +299,27 @@ class SonosQueueLoader @Inject constructor(
}
}
/**
* Called on every poll with the renderer's 1-based track number. When it
* moves, the track after it is rendered ahead, so the speaker's fetch of
* it finds the render ready.
*/
fun onRendererTrack(trackNumber: Int) {
if (trackNumber <= 0 || trackNumber == prerenderedAfter) return
prerenderedAfter = trackNumber
val queue = sent
scope.launch { prerender(queue, trackNumber) }
}
private suspend fun prerender(queue: List<TrackRef>, index: Int) {
if (index !in queue.indices) return
runCatching { mint(queue, index, prerender = true) }
.onFailure { Timber.w(it, "UPnP prerender failed for %s", queue[index].id) }
}
private suspend fun mint(queue: List<TrackRef>, index: Int, prerender: Boolean = false) =
streamTokens.mint(queue[index].id, asAlbum = playingAsAlbum(queue, index), prerender = prerender)
private fun commonPrefixLength(a: List<String>, b: List<String>): Int {
val limit = minOf(a.size, b.size)
for (i in 0 until limit) {
@@ -299,18 +336,33 @@ class SonosQueueLoader @Inject constructor(
return limit
}
private companion object {
companion object {
/**
* Whether [queue]'s track at [index] plays among its album in order,
* judged by its neighbours in the renderer's queue, which plays
* straight through.
*/
internal fun playingAsAlbum(queue: List<TrackRef>, index: Int): Boolean =
GainMath.playingAsAlbum(
prev = queue.getOrNull(index - 1)?.let(::position),
cur = position(queue[index]),
next = queue.getOrNull(index + 1)?.let(::position),
)
private fun position(t: TrackRef) =
AlbumPosition(t.albumId.ifEmpty { null }, t.discNumber, t.trackNumber)
// Abort the append loop after this many consecutive AddURIToQueue
// failures; Sonos rate-limits burst adds and a wall of failures means
// it has stopped accepting, not that the next one might land.
const val EXTEND_ABORT_AFTER_FAILURES = 3
const val EXTEND_THROTTLE_MS = 50L
private const val EXTEND_ABORT_AFTER_FAILURES = 3
private const val EXTEND_THROTTLE_MS = 50L
// Verify/repair passes after a queue load. Two: one to catch the
// common case (a rate-limit burst dropped a chunk), one to catch a
// repair that itself got rate-limited. Beyond that the renderer is
// refusing for a reason retrying won't fix, and the stall watchdog
// becomes the backstop.
const val VERIFY_ROUNDS = 2
private const val VERIFY_ROUNDS = 2
}
}
@@ -0,0 +1,43 @@
package com.fabledsword.minstrel.player.output
import com.fabledsword.minstrel.models.TrackRef
import org.junit.jupiter.api.Test
import kotlin.test.assertEquals
/**
* The Sonos queue plays straight through, so a track's album play is judged
* from its neighbours in the list sent (M464 #5002).
*/
class SonosAlbumPlayTest {
private fun t(album: String, n: Int?, disc: Int? = null) =
TrackRef(id = "$album-$disc-$n", title = "", albumId = album, artistId = "", trackNumber = n, discNumber = disc)
private fun asAlbum(queue: List<TrackRef>) = queue.indices.map { SonosQueueLoader.playingAsAlbum(queue, it) }
@Test
fun `an album in order plays as an album throughout`() {
assertEquals(listOf(true, true, true), asAlbum(listOf(t("a", 1), t("a", 2), t("a", 3))))
}
@Test
fun `a mix of albums plays as singles`() {
assertEquals(listOf(false, false, false), asAlbum(listOf(t("a", 1), t("b", 1), t("c", 4))))
}
@Test
fun `album order runs across discs`() {
assertEquals(listOf(true, true), asAlbum(listOf(t("a", 12, disc = 1), t("a", 1, disc = 2))))
}
@Test
fun `shuffled tracks of one album are not album play`() {
assertEquals(listOf(false, false), asAlbum(listOf(t("a", 3), t("a", 1))))
}
@Test
fun `unnumbered tracks and missing albums are never album play`() {
assertEquals(listOf(false, false), asAlbum(listOf(t("a", null), t("a", null))))
assertEquals(listOf(false, false), asAlbum(listOf(t("", 1), t("", 2))))
}
}
+7 -2
View File
@@ -26,6 +26,11 @@ type castTokenRequest struct {
// client holding the queue knows it.
Level bool `json:"level,omitempty"`
AsAlbum bool `json:"asAlbum,omitempty"`
// Prerender says the speaker will fetch this track soon: the current
// track at a queue load, or the next one as it starts. Only those are
// rendered ahead; a queue load mints every track and renders none of
// the rest until the speaker asks.
Prerender bool `json:"prerender,omitempty"`
}
type castTokenResponse struct {
@@ -138,13 +143,13 @@ func (h *handlers) handleCastStreamToken(w http.ResponseWriter, r *http.Request)
path = leveledStreamPath(trackUUID) + leveledQuery(g, token, exp)
mime = "audio/flac"
leveled = true
// The speaker fetches the URL shortly; start rendering now so
// the fetch finds the file ready.
if req.Prerender {
h.leveled.Prerender(library.LeveledSource{
TrackID: req.TrackID, Path: track.FilePath, DurationMs: track.DurationMs,
}, g)
}
}
}
// Behind a TLS-terminating reverse proxy, r.TLS is nil even though
// the public-facing URL is https://. UPnP devices (Sonos especially)
+23 -4
View File
@@ -161,6 +161,11 @@ type LeveledRenderer struct {
group singleflight.Group
evictMu sync.Mutex
// prerenders bounds the background renders running at once. A fetch
// renders on demand regardless, so a prerender that finds no slot is
// dropped, not queued.
prerenders chan struct{}
// render and probe are ffmpeg and ffprobe; tests replace them.
render func(ctx context.Context, src, dst string, g LeveledGain, s leveledSource) error
probe func(ctx context.Context, src string) (leveledSource, error)
@@ -175,6 +180,7 @@ func NewLeveledRenderer(dir string, settings *LoudnessSettingsService, logger *s
dir: dir,
settings: settings,
logger: logger,
prerenders: make(chan struct{}, maxLeveledPrerenders),
render: ffmpegRenderLeveled,
probe: ffprobeLeveledSource,
}, nil
@@ -217,17 +223,30 @@ func (r *LeveledRenderer) Path(ctx context.Context, src LeveledSource, g Leveled
}
}
// Prerender starts rendering src at g in the background, so the renderer's
// fetch finds it ready. Used when a leveled URL is handed out: the speaker
// asks for the next track shortly before it plays.
func (r *LeveledRenderer) Prerender(src LeveledSource, g LeveledGain) {
// maxLeveledPrerenders is how many prerenders may run at once. A speaker
// needs the track it is about to play and the one after; more than that is a
// client asking for too much, and the excess is dropped.
const maxLeveledPrerenders = 2
// Prerender starts rendering src at g in the background, so the speaker's
// fetch finds it ready. It reports false, and does nothing, when the
// prerender slots are full.
func (r *LeveledRenderer) Prerender(src LeveledSource, g LeveledGain) bool {
select {
case r.prerenders <- struct{}{}:
default:
r.logger.Debug("leveled: prerender dropped, slots full", "track", src.TrackID)
return false
}
go func() {
defer func() { <-r.prerenders }()
ctx, cancel := context.WithTimeout(context.Background(), loudnessTimeout(src.DurationMs))
defer cancel()
if _, err := r.Path(ctx, src, g); err != nil {
r.logger.Warn("leveled: prerender failed", "track", src.TrackID, "err", err)
}
}()
return true
}
func (r *LeveledRenderer) renderTo(ctx context.Context, src, dst string, g LeveledGain) error {
+35
View File
@@ -195,6 +195,41 @@ func TestLeveledRenderer_CoalescesAndCaches(t *testing.T) {
}
}
// Prerenders past the slot limit are dropped rather than piling up ffmpeg
// processes; a slot frees when its render ends.
func TestLeveledRenderer_PrerenderSlots(t *testing.T) {
r, renders, _ := testRenderer(t, 2048, 10)
gate := make(chan struct{})
r.render = func(_ context.Context, _, dst string, _ LeveledGain, _ leveledSource) error {
<-gate
renders.Add(1)
return os.WriteFile(dst, []byte("flac"), 0o600)
}
g := LeveledGain{CentiDB: -300}
for i := range maxLeveledPrerenders {
if !r.Prerender(writeSource(t, string(rune('a'+i))), g) {
t.Fatalf("prerender %d refused with slots free", i)
}
}
if r.Prerender(writeSource(t, "z"), g) {
t.Fatal("a prerender past the limit was accepted")
}
close(gate)
deadline := time.Now().Add(5 * time.Second)
for !r.Prerender(writeSource(t, "y"), g) {
if time.Now().After(deadline) {
t.Fatal("no slot freed after the renders finished")
}
time.Sleep(10 * time.Millisecond)
}
for renders.Load() != maxLeveledPrerenders+1 {
if time.Now().After(deadline) {
t.Fatalf("renders = %d, want %d", renders.Load(), maxLeveledPrerenders+1)
}
time.Sleep(10 * time.Millisecond)
}
}
func TestLeveledRenderer_MissingSource(t *testing.T) {
r, _, _ := testRenderer(t, 2048, 10)
_, err := r.Path(context.Background(), LeveledSource{TrackID: "x", Path: "/nonexistent/x.flac"}, LeveledGain{CentiDB: 100})