diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/ExoPlayerPool.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/ExoPlayerPool.kt index 5100afe6d4..6e9c18c89c 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/ExoPlayerPool.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/ExoPlayerPool.kt @@ -146,8 +146,7 @@ class ExoPlayerPool( // Count only once the player exists — incrementing first leaves a phantom decoder on the // books if builder.build() throws. val player = coldPool.poll() ?: builder.build(context) - liveDecoders.incrementAndGet() - Log.d(PLAYBACK_DIAG_TAG) { "DECODERS acquire -> ${liveDecoders.get()} / budget $poolSize" } + acquireDecoder() return player } @@ -171,13 +170,20 @@ class ExoPlayerPool( } } + // acquireDecoder/releaseDecoder are the only two places the shared counter moves, so the + // DECODERS trace always logs the value the mutation itself produced (a re-read could show a + // sibling pool's concurrent move instead). + private fun acquireDecoder() { + val now = liveDecoders.incrementAndGet() + Log.d(PLAYBACK_DIAG_TAG) { "DECODERS acquire -> $now / budget $poolSize" } + } + // Floored at zero. An under-count is the harmful direction: it lets ensureDecoderHeadroom // hand out more concurrent decoders than the device grants, which surfaces as MediaCodec // NO_MEMORY and "can't play this video". An over-count only costs the warm cache. - private fun releaseDecoder(): Int { + private fun releaseDecoder() { val now = liveDecoders.updateAndGet { (it - 1).coerceAtLeast(0) } Log.d(PLAYBACK_DIAG_TAG) { "DECODERS release -> $now / budget $poolSize" } - return now } private fun evictOldestWarm(): Boolean { @@ -208,6 +214,12 @@ class ExoPlayerPool( null } + /** + * Queues the return on this pool's single main-thread scope. Contract: returns queued before + * [destroy] is called run before destroy's teardown — both serialize on the same scope and + * [mutex] — so a caller may tear its sessions down synchronously and then destroy the pool + * (see MediaSessionPool.destroy). + */ fun releasePlayerAsync(player: ExoPlayer) { scope.launch { releasePlayer(player) @@ -342,22 +354,7 @@ class ExoPlayerPool( } coldPool.clear() - // Remove from the shared registry only after this pool has decremented its own - // players, so livePools.isEmpty() becomes true only once EVERY pool has finished - // tearing down. Removing synchronously at the top of destroy() (the previous - // approach) let the first pool's body see an empty list while a sibling's - // decoders were still counted, firing a false drift warning on every ordinary - // two-pool shutdown. - livePools.remove(this@ExoPlayerPool) - - if (livePools.isEmpty()) { - // Last pool out. A non-zero value here is now a genuine accounting leak — - // a player acquired and never returned — worth seeing. - val stranded = liveDecoders.getAndSet(0) - if (stranded != 0) { - Log.w(PLAYBACK_DIAG_TAG) { "decoder accounting drift at teardown: $stranded" } - } - } + poolFinishedTeardown(this@ExoPlayerPool) } }.invokeOnCompletion { scope.cancel() @@ -378,5 +375,20 @@ class ExoPlayerPool( // Every pool that hasn't been destroy()'d, so a pool starved of headroom can reclaim a // warm player from a sibling instead of overshooting the shared ceiling. private val livePools = ConcurrentLinkedQueue() + + // A pool calls this only after decrementing its own players, so livePools.isEmpty() + // means EVERY pool has finished tearing down — deregistering earlier would let the first + // pool out see a sibling's still-counted decoders as drift. The last pool out resets the + // shared counter; a non-zero value at that point is a genuine accounting leak (a player + // acquired and never returned), worth seeing. + private fun poolFinishedTeardown(pool: ExoPlayerPool) { + livePools.remove(pool) + if (livePools.isEmpty()) { + val stranded = liveDecoders.getAndSet(0) + if (stranded != 0) { + Log.w(PLAYBACK_DIAG_TAG) { "decoder accounting drift at teardown: $stranded" } + } + } + } } } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/MediaSessionPool.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/MediaSessionPool.kt index 9f8cec496e..f6a81cc141 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/MediaSessionPool.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/MediaSessionPool.kt @@ -52,8 +52,8 @@ class SessionListener( val session: MediaSession, val playerListener: Player.Listener, ) { - // Set once, by the retire funnel. Guards the re-entrant path session.release() -> - // onDisconnected -> releaseSession -> registry.drop() landing on this same entry. + // Set once, by the retire funnel, making retireSession idempotent no matter how many + // paths reach the same entry. val retired = AtomicBoolean(false) fun removeListeners() { @@ -133,18 +133,21 @@ class MediaSessionPool( } /** - * The one place a session's player goes back to the pool. + * The one place a session's player goes back to the pool — the registry's onDropped and + * newSession's failure unwind both land here. * - * The CAS is not redundant even though [SessionRegistry] signals each drop exactly once: - * session.release() below can reach onDisconnected -> releaseSession -> registry.drop() on - * the same entry, and that re-entrant path is what the guard stops. Do not remove it. + * [SessionRegistry] removes an entry from both of its maps before signaling, so its paths + * already deliver each entry at most once; the [SessionListener.retired] CAS keeps the funnel + * idempotent regardless, because a double retire would double-return the player — an + * under-count, the harmful direction (toward MediaCodec NO_MEMORY). * * Both orderings below are load-bearing, not tidiness: * - removeListeners() before release(), because releasing a session can fire * onIsPlayingChanged, which would re-enter setPlaying while we are still inside * the registry's entryRemoved callback. * - releasePlayerAsync() before release(), so a throw from session teardown cannot - * strand the player. + * strand the player. releasePlayerAsync only queues the return, so the listener and + * session are detached before the queued release ever runs. */ private fun retireSession(entry: SessionListener) { if (!entry.retired.compareAndSet(false, true)) return @@ -172,15 +175,13 @@ class MediaSessionPool( ): MediaSession { val player = exoPlayerPool.acquirePlayer(context, preferredMediaId) - // newSession owns the player until registry.register() hands ownership over. Anything that - // throws in between — MediaSession.Builder.build(), reset(), or the PendingIntent in - // bindSessionActivity() — would otherwise strand a counted player with no owner. - // - // The throw sites named above are all *after* the session exists, so returning the player - // alone is not enough: a live MediaSession bound to it, and a listener attached to it, - // must come off first, or the pool re-issues a player that something else still holds. - var session: MediaSession? = null - var connector: MediaSessionExoPlayerConnector? = null + // newSession owns the player until registry.register() hands ownership over. Anything + // that throws in between — reset(), or the PendingIntent in bindSessionActivity() — + // would otherwise strand a counted player with no owner. The unwind routes through + // retireSession so there is exactly one teardown sequence; removing a never-added + // listener there is a no-op. + var entry: SessionListener? = null + var handedOff = false try { val mediaSession = @@ -191,10 +192,9 @@ class MediaSessionPool( setId(id) setCallback(globalCallback) }.build() - session = mediaSession val listener = MediaSessionExoPlayerConnector(mediaSession, this) - connector = listener + entry = SessionListener(mediaSession, listener) mediaSession.player.addListener(listener) @@ -207,16 +207,19 @@ class MediaSessionPool( // notification opens the originating nostr URI. bindSessionActivity(mediaSession, mediaSession.player.currentMediaItem) - registry.register(mediaSession.id, SessionListener(mediaSession, listener)) + // Past this point the registry owns the session. register() inserts the entry before + // it can fire an eviction's onDropped, so even if register() itself throws (a + // displaced entry's retire failing), the new entry is already registered and a later + // sweep retires it — the catch must not also return this player, which would + // double-return it. + handedOff = true + registry.register(mediaSession.id, entry) return mediaSession } catch (e: Throwable) { - // Unwind in reverse. releasePlayerAsync only *queues* the return, so the session and - // listener are detached synchronously before the queued release ever runs — same - // rationale as the ordering inside retireSession. - exoPlayerPool.releasePlayerAsync(player) - connector?.let { player.removeListener(it) } - session?.release() + if (!handedOff) { + entry?.let { retireSession(it) } ?: exoPlayerPool.releasePlayerAsync(player) + } throw e } } @@ -237,14 +240,13 @@ class MediaSessionPool( } fun releaseSession(session: MediaSession) { - // The registry removes the entry and then signals; retireSession does the teardown. - // Nothing here depends on a cache removal firing a callback — that dependency is what - // stranded players whose entry had already been evicted while playing. + // The registry removes the entry and then signals; retireSession does the teardown — + // explicitly, whether the session is idle or playing, never by relying on a cache + // removal firing a callback (a no-op once the entry was evicted while playing). // - // Unlike the old code this does NOT release the session when the id is unknown. Not in the - // registry means already retired (retireSession released it) or never owned, so releasing - // again would be exactly the double-release being removed here. newSession's failure path - // produces such a session. + // An unknown id drops nothing: not in the registry means already retired (retireSession + // released it) or never owned (newSession's failure path), and releasing such a session + // again would be a double-release. registry.drop(session.id) cleanupUnused() } @@ -265,14 +267,17 @@ class MediaSessionPool( } fun destroy() { - // Synchronous: retireSession uses the non-suspending releasePlayerAsync, so no coroutine - // is needed here. Ordering still holds — releasePlayerAsync and ExoPlayerPool.destroy() - // both launch on ExoPlayerPool's single main-thread scope and serialize on its one Mutex, - // so queued returns run first. The old version launched its teardown and then cancelled - // the scope on the same frame, so the teardown never ran at all. - registry.dropAll() - exoPlayerPool.destroy() - scope.cancel() + // Synchronous on purpose: retireSession uses the non-suspending releasePlayerAsync, and + // releasePlayerAsync's contract guarantees returns queued here run before the pool's own + // teardown — so every player is back before exoPlayerPool.destroy() sweeps. + try { + registry.dropAll() + } finally { + // The player pool must tear down even if a retire throws mid-sweep; skipping it + // would leak every pooled player, not just the session that failed. + exoPlayerPool.destroy() + scope.cancel() + } } fun getSession( @@ -288,9 +293,7 @@ class MediaSessionPool( fun playingContent() = registry.playingEntries() - // Widened from cache-only to cache-or-playing. Inert at its only caller - // (PlaybackService.onUpdateNotification), which is inside a `playing.isEmpty()` branch - // covering both pools, so no session is playing when it runs. + // Resolves the session whether it is idle or playing. fun getSession(id: String) = registry.get(id)?.session class MediaSessionCallback( @@ -322,10 +325,9 @@ class MediaSessionPool( val pool: MediaSessionPool, ) : Player.Listener { override fun onIsPlayingChanged(isPlaying: Boolean) { - // Moves the existing entry. Allocating a fresh SessionListener here (the old - // behaviour) meant one session could have two or three live wrappers at once, so a - // drop of one could hand its player back while another still pointed at that live - // session. + // Moves the one existing entry between tiers. A session must never have more than + // one live wrapper: dropping one wrapper would hand the player back while another + // still pointed at the live session. pool.setPlaying(mediaSession.id, isPlaying) } } diff --git a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistry.kt b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistry.kt index 768a420458..9b47453081 100644 --- a/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistry.kt +++ b/amethyst/src/main/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistry.kt @@ -97,9 +97,10 @@ internal class SessionRegistry( fun get(id: String): T? = playing[id] ?: idle.get(id) - fun idleSnapshot(): List = idle.snapshot().values.toList() + // snapshot() already materializes a detached map, so its values need no second copy. + fun idleSnapshot(): Collection = idle.snapshot().values - // Copies, like idleSnapshot(). Handing out playing.values would be a live view, and a caller + // Detached, like idleSnapshot(). Handing out playing.values would be a live view, and a caller // that iterated it while a drop fired would get a ConcurrentModificationException. A type whose // job is owning reachability should not ship that footgun. fun playingEntries(): List = playing.values.toList() @@ -110,13 +111,15 @@ internal class SessionRegistry( * nothing when the entry was already evicted while playing. */ fun drop(id: String): T? { - val entry = playing.remove(id) ?: idle.get(id) ?: return null + val playingEntry = playing.remove(id) suppressDropSignal = true - try { - idle.remove(id) - } finally { - suppressDropSignal = false - } + val idleEntry = + try { + idle.remove(id) + } finally { + suppressDropSignal = false + } + val entry = playingEntry ?: idleEntry ?: return null onDropped(entry) return entry } diff --git a/amethyst/src/test/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistryTest.kt b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistryTest.kt index 3fea991a07..35202d9a7f 100644 --- a/amethyst/src/test/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistryTest.kt +++ b/amethyst/src/test/java/com/vitorpamplona/amethyst/service/playback/playerPool/SessionRegistryTest.kt @@ -140,7 +140,7 @@ class SessionRegistryTest { registry.register("b", "B") registry.setPlaying("a", true) assertEquals(setOf("A", "B"), registry.idleSnapshot().toSet()) - assertEquals(listOf("A"), registry.playingEntries().toList()) + assertEquals(listOf("A"), registry.playingEntries()) } // The replacement guard uses !== (identity, not equality). This test pins that distinction: