diff --git a/android/app/src/main/java/com/fabledsword/minstrel/api/endpoints/CastApi.kt b/android/app/src/main/java/com/fabledsword/minstrel/api/endpoints/CastApi.kt index e291f4a5..0e2fb3d6 100644 --- a/android/app/src/main/java/com/fabledsword/minstrel/api/endpoints/CastApi.kt +++ b/android/app/src/main/java/com/fabledsword/minstrel/api/endpoints/CastApi.kt @@ -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, ) diff --git a/android/app/src/main/java/com/fabledsword/minstrel/player/MinstrelForwardingPlayer.kt b/android/app/src/main/java/com/fabledsword/minstrel/player/MinstrelForwardingPlayer.kt index 1556bb42..af40b975 100644 --- a/android/app/src/main/java/com/fabledsword/minstrel/player/MinstrelForwardingPlayer.kt +++ b/android/app/src/main/java/com/fabledsword/minstrel/player/MinstrelForwardingPlayer.kt @@ -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 -> { diff --git a/android/app/src/main/java/com/fabledsword/minstrel/player/PlayerFactory.kt b/android/app/src/main/java/com/fabledsword/minstrel/player/PlayerFactory.kt index 0b655d86..2cc0cc0e 100644 --- a/android/app/src/main/java/com/fabledsword/minstrel/player/PlayerFactory.kt +++ b/android/app/src/main/java/com/fabledsword/minstrel/player/PlayerFactory.kt @@ -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) }, ), ) } diff --git a/android/app/src/main/java/com/fabledsword/minstrel/player/StreamTokenProvider.kt b/android/app/src/main/java/com/fabledsword/minstrel/player/StreamTokenProvider.kt index ffff44bf..45e8b281 100644 --- a/android/app/src/main/java/com/fabledsword/minstrel/player/StreamTokenProvider.kt +++ b/android/app/src/main/java/com/fabledsword/minstrel/player/StreamTokenProvider.kt @@ -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), + ) } diff --git a/android/app/src/main/java/com/fabledsword/minstrel/player/output/SonosQueueLoader.kt b/android/app/src/main/java/com/fabledsword/minstrel/player/output/SonosQueueLoader.kt index 950b5024..355abbfd 100644 --- a/android/app/src/main/java/com/fabledsword/minstrel/player/output/SonosQueueLoader.kt +++ b/android/app/src/main/java/com/fabledsword/minstrel/player/output/SonosQueueLoader.kt @@ -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 = 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, + queue: List, 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, + queue: List, 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, ): 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, 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, index: Int, prerender: Boolean = false) = + streamTokens.mint(queue[index].id, asAlbum = playingAsAlbum(queue, index), prerender = prerender) + private fun commonPrefixLength(a: List, b: List): 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, 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 } } diff --git a/android/app/src/test/java/com/fabledsword/minstrel/player/output/SonosAlbumPlayTest.kt b/android/app/src/test/java/com/fabledsword/minstrel/player/output/SonosAlbumPlayTest.kt new file mode 100644 index 00000000..d5e28f69 --- /dev/null +++ b/android/app/src/test/java/com/fabledsword/minstrel/player/output/SonosAlbumPlayTest.kt @@ -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) = 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)))) + } +} diff --git a/internal/api/cast_token.go b/internal/api/cast_token.go index a1594874..8930f2c5 100644 --- a/internal/api/cast_token.go +++ b/internal/api/cast_token.go @@ -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,11 +143,11 @@ 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. - h.leveled.Prerender(library.LeveledSource{ - TrackID: req.TrackID, Path: track.FilePath, DurationMs: track.DurationMs, - }, g) + if req.Prerender { + h.leveled.Prerender(library.LeveledSource{ + TrackID: req.TrackID, Path: track.FilePath, DurationMs: track.DurationMs, + }, g) + } } } diff --git a/internal/library/leveled.go b/internal/library/leveled.go index 4fb982b7..b7abdd11 100644 --- a/internal/library/leveled.go +++ b/internal/library/leveled.go @@ -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) @@ -172,11 +177,12 @@ func NewLeveledRenderer(dir string, settings *LoudnessSettingsService, logger *s return nil, fmt.Errorf("leveled cache: %w", err) } return &LeveledRenderer{ - dir: dir, - settings: settings, - logger: logger, - render: ffmpegRenderLeveled, - probe: ffprobeLeveledSource, + 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 { diff --git a/internal/library/leveled_test.go b/internal/library/leveled_test.go index ff18a2e1..7b8c49c2 100644 --- a/internal/library/leveled_test.go +++ b/internal/library/leveled_test.go @@ -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})