diff --git a/android/README.md b/android/README.md index 549aedc..d654bec 100644 --- a/android/README.md +++ b/android/README.md @@ -61,8 +61,15 @@ the one core for the process through `CoreHost`: own headers) plus the gapless follow-up as a second playlist item; `SetNext` replaces everything after the current item; the transition is detected from `onMediaItemTransition(AUTO)`, reported as `Ended` (played item) then `TransitionedToNext`, and the played item removed. The player holds `C.WAKE_MODE_NETWORK` (wake + - Wi-Fi lock) so streams keep going with the screen off. Network errors (connection failed/timeout) - are retried with `prepare()` on a ~1 min backoff and reported non-fatal; only then fatal. `PreBuffer` prepares a + Wi-Fi lock) so streams keep going with the screen off. Errors a later `prepare()` can fix (the + connection, a timeout, a 5xx from the server, a stuck player, a core stream handle that went away + under the player) are retried with `prepare()` on a ~1 min backoff and reported non-fatal; while + the device is offline the retry waits instead of spending attempts, and the network coming back + (`NetworkMonitor` -> `ExoBackend.onConnectivityChanged`) retries at once; only when the attempts run + out is the error fatal and the core's own retry/skip takes over. `Paused` is reported only when the + player no longer means to play (`playWhenReady` false); a stall with it still set is `Buffering`, + and the service never stops itself under a player that is buffering or recovering + (`ExoBackend.isBusy`). `PreBuffer` prepares a second silent ExoPlayer at the requested position (`PreBufferReady` when READY); `DiscardPreBuffer` releases it. `gain_db` is applied as `10^(gain/20) * masterVolume` clamped to 1.0 — Media3 has no gain stage, so positive gain is an approximation (documented in `ExoBackend`). @@ -74,7 +81,9 @@ the one core for the process through `CoreHost`: loader threads): a seek is a new open at the position; the core fetches from the server, caches a whole read and serves later plays and seeks from disk. An unknown/expired token is `ERROR_CODE_IO_FILE_NOT_FOUND`, never retried (`CoreStreamLoadErrorPolicy`), so the backend reports - a fatal error and the core resolves a fresh token. The factory routes the `hocket-stream` scheme to + a fatal error and the core resolves a fresh token; a handle the core closed under the player (idle + while the buffer was full) is retried with a fresh open at the same position, which the cache + serves. The factory routes the `hocket-stream` scheme to it and everything else (`file:`, a direct server URL for the fake core) to `DefaultDataSource`; there is no `CacheDataSource` (the core caches). `NativeCore` implements the `CoreStreams` seam. - Reports back: `Ready`, `Playing`, `Paused`, `Buffering`, `Position` every 750 ms while playing and diff --git a/android/playback/src/main/java/app/hocket/playback/ExoBackend.kt b/android/playback/src/main/java/app/hocket/playback/ExoBackend.kt index 51136e1..416c6bc 100644 --- a/android/playback/src/main/java/app/hocket/playback/ExoBackend.kt +++ b/android/playback/src/main/java/app/hocket/playback/ExoBackend.kt @@ -8,11 +8,13 @@ import androidx.media3.common.C import androidx.media3.common.MediaItem import androidx.media3.common.PlaybackException import androidx.media3.common.Player +import androidx.media3.datasource.HttpDataSource import androidx.media3.exoplayer.DefaultLoadControl import androidx.media3.exoplayer.ExoPlayer import androidx.media3.exoplayer.source.DefaultMediaSourceFactory import androidx.media3.exoplayer.source.MediaSource import app.hocket.core.Commands +import app.hocket.core.CoreStreamException import app.hocket.core.CoreStreams import app.hocket.core.api.BackendCommand import app.hocket.core.api.BackendReport @@ -62,9 +64,17 @@ import kotlin.math.pow * then resumes by itself when focus returns and reports `Playing`). * - The player holds a wake lock and a Wi-Fi lock while playing ([C.WAKE_MODE_NETWORK]): without * them a stream stalls once the screen is off and the CPU or Wi-Fi radio sleeps. - * - Network failures leave ExoPlayer idle with an error, so they are retried here with backoff - * ([recoveryDelayMs]; `prepare()` resumes at the same position) and reported as non-fatal; only when - * the retries run out is the error fatal and the core's own retry/skip takes over. + * - `Paused` is reported only when the player no longer means to play (`playWhenReady` false: a + * pause, becoming noisy, a lasting focus loss). A stall with `playWhenReady` still set (an empty + * buffer, a load being retried, an error being recovered from) is `Buffering`, never `Paused`: the + * core would otherwise show the item paused, and the service would count the player as idle and + * could stop itself under a playback that was only rebuffering. + * - Errors a later `prepare()` can fix ([isRecoverable]: the connection, a timeout, a 5xx from the + * server, a stream handle the core closed, a stuck player) leave ExoPlayer idle with an error, so + * they are retried here with backoff ([recoveryDelayMs]; `prepare()` resumes at the same position) + * and reported as non-fatal. While the device is offline ([onConnectivityChanged]) no attempt is + * spent: the retry waits, and the network coming back retries at once with a fresh schedule. Only + * when the attempts run out is the error fatal and the core's own retry/skip takes over. */ class ExoBackend( private val context: Context, @@ -73,6 +83,9 @@ class ExoBackend( /** The running core's stream reader, read on each open (a restarted service has a new core). */ private val streams: () -> CoreStreams? = { CoreHost.current as? CoreStreams }, ) { + /** What a player that stopped playing while it still has an item is reporting. */ + internal enum class Stall { Paused, Buffering, Nothing } + internal companion object { private const val TAG = "ExoBackend" private const val POSITION_INTERVAL_MS = 750L @@ -81,10 +94,49 @@ class ExoBackend( /** Backoff before recovery attempt [attempt] (0-based) of a network error, or null when out of attempts. */ internal fun recoveryDelayMs(attempt: Int): Long? = RECOVERY_DELAYS_MS.getOrNull(attempt) - /** Errors a later `prepare()` can fix: the connection, not the media. */ + /** Errors a later `prepare()` can fix by their code alone: the connection, not the media. */ internal fun isRecoverable(errorCode: Int): Boolean = errorCode == PlaybackException.ERROR_CODE_IO_NETWORK_CONNECTION_FAILED || - errorCode == PlaybackException.ERROR_CODE_IO_NETWORK_CONNECTION_TIMEOUT + errorCode == PlaybackException.ERROR_CODE_IO_NETWORK_CONNECTION_TIMEOUT || + errorCode == PlaybackException.ERROR_CODE_TIMEOUT + + /** + * Whether a later `prepare()` can fix [error]: a network failure or timeout, a stuck player + * (`ERROR_CODE_TIMEOUT`), a 5xx from the server (the core's stream status or a direct HTTP + * source), or a core stream handle that went away under the player (closed as idle, too many + * open, unknown): a re-open serves it again. A 4xx, an unknown token (the core re-resolves), a + * core that shut down, and anything about the media itself are not. + */ + internal fun isRecoverable(error: PlaybackException): Boolean { + if (isRecoverable(error.errorCode)) return true + var t: Throwable? = error + while (t != null) { + when (t) { + is CoreStreamException -> return when (t.kind) { + CoreStreamException.Kind.Network, CoreStreamException.Kind.Closed, CoreStreamException.Kind.UnknownHandle, + CoreStreamException.Kind.TooManyHandles, CoreStreamException.Kind.NoServer -> true + CoreStreamException.Kind.Status -> (t.httpStatus ?: 0) >= 500 + else -> false + } + is HttpDataSource.InvalidResponseCodeException -> return t.responseCode >= 500 + is HttpDataSource.CleartextNotPermittedException -> return false + is HttpDataSource.HttpDataSourceException -> return true + } + t = t.cause + } + return false + } + + /** + * [Player.Listener.onIsPlayingChanged] `false`: [Stall.Paused] when the player no longer + * means to play, [Stall.Buffering] while it does but has nothing to play yet; idle (an error, + * a stop) and ended are reported by their own callbacks. + */ + internal fun stall(playWhenReady: Boolean, playbackState: Int): Stall = when { + !playWhenReady -> if (playbackState == Player.STATE_IDLE || playbackState == Player.STATE_ENDED) Stall.Nothing else Stall.Paused + playbackState == Player.STATE_BUFFERING -> Stall.Buffering + else -> Stall.Nothing + } } val player: ExoPlayer = ExoPlayer.Builder(context) @@ -103,6 +155,11 @@ class ExoBackend( private var recoveryJob: Job? = null /** Recovery attempts for the current error streak; reset once the player is ready again. */ private var recoveryAttempts = 0 + /** The key a pending recovery is for. */ + private var recoveryKey: String? = null + /** Connectivity as the service last reported it; a retry waits while this is false. */ + @Volatile + private var online = true /** The key of the last `Load`, so an error raised after the playlist emptied is still attributable. */ private var lastLoadedKey: String? = null /** Keys by media id, so reports name the queue key the core gave us. */ @@ -117,7 +174,7 @@ class ExoBackend( val key = currentKey() ?: return when (playbackState) { Player.STATE_READY -> { - recoveryAttempts = 0 + cancelRecovery() val duration = player.duration.takeIf { it != C.TIME_UNSET }?.toUInt() report(BackendReport.Ready(BackendReportReadyInner(key, duration))) report(BackendReport.Buffering(BackendReportBufferingInner(key, false))) @@ -140,17 +197,22 @@ class ExoBackend( startPositionLoop() } else { stopPositionLoop() - if (player.playbackState != Player.STATE_ENDED && !player.isLoading) { - report(BackendReport.Paused(BackendReportPausedInner(key, position()))) + when (stall(player.playWhenReady, player.playbackState)) { + Stall.Paused -> report(BackendReport.Paused(BackendReportPausedInner(key, position()))) + Stall.Buffering -> report(BackendReport.Buffering(BackendReportBufferingInner(key, true))) + Stall.Nothing -> Unit } } } override fun onPlayWhenReadyChanged(playWhenReady: Boolean, reason: Int) { - if (!playWhenReady && reason == Player.PLAY_WHEN_READY_CHANGE_REASON_AUDIO_FOCUS_LOSS) { + if (playWhenReady) return + if (reason == Player.PLAY_WHEN_READY_CHANGE_REASON_AUDIO_FOCUS_LOSS) { report(BackendReport.AudioFocusLost(BackendReportAudioFocusLostInner(transient = false))) } - if (!playWhenReady && reason == Player.PLAY_WHEN_READY_CHANGE_REASON_AUDIO_BECOMING_NOISY) { + // The player stopped meaning to play, whoever asked (a pause while buffering never + // reaches onIsPlayingChanged, so this is the only place that sees it). + if (player.playbackState == Player.STATE_BUFFERING || player.playbackState == Player.STATE_READY) { currentKey()?.let { report(BackendReport.Paused(BackendReportPausedInner(it, position()))) } } } @@ -176,27 +238,72 @@ class ExoBackend( val key = errorKey() ?: return val message = error.errorCodeName + ": " + (error.message ?: "") // Any player error leaves ExoPlayer idle: "non-fatal" only holds if we bring it back. - val delayMs = if (isRecoverable(error.errorCode)) recoveryDelayMs(recoveryAttempts) else null + val delayMs = if (isRecoverable(error)) recoveryDelayMs(recoveryAttempts) else null if (delayMs == null) { - recoveryAttempts = 0 + cancelRecovery() report(BackendReport.Error(BackendReportErrorInner(key, message, true))) return } recoveryAttempts++ + Log.w(TAG, "playback error on $key: ${error.errorCodeName}; retrying in ${delayMs}ms (attempt $recoveryAttempts)") report(BackendReport.Error(BackendReportErrorInner(key, message, false))) report(BackendReport.Buffering(BackendReportBufferingInner(key, true))) - recoveryJob?.cancel() - recoveryJob = scope.launch { - delay(delayMs) - if (player.playerError != null && currentKey() == key) { - Log.i(TAG, "retrying $key after ${error.errorCodeName} (attempt $recoveryAttempts)") - player.prepare() - } - } + scheduleRecovery(key, delayMs) } }) } + private fun scheduleRecovery(key: String, delayMs: Long) { + recoveryJob?.cancel() + recoveryKey = key + recoveryJob = scope.launch { + delay(delayMs) + attemptRecovery(key) + } + } + + /** + * Re-prepare after a recoverable error, if the error still stands and the item is still the one + * that failed. Offline, the attempt is not spent: it waits for [onConnectivityChanged], with the + * longest delay as a safety net should that never come. + */ + private fun attemptRecovery(key: String) { + if (player.playerError == null || currentKey() != key) { + recoveryKey = null + return + } + if (!online) { + Log.i(TAG, "offline; waiting for connectivity before retrying $key") + scheduleRecovery(key, RECOVERY_DELAYS_MS.last()) + return + } + Log.i(TAG, "retrying $key (attempt $recoveryAttempts)") + player.prepare() + } + + /** + * Connectivity as the service sees it. Coming back online while a recovery is pending retries at + * once, with the attempts reset (the failures so far were the old network's). + */ + fun onConnectivityChanged(online: Boolean) { + this.online = online + if (!online) return + val key = recoveryKey ?: return + if (recoveryJob?.isActive != true) return + recoveryJob?.cancel() + recoveryAttempts = 0 + attemptRecovery(key) + } + + /** + * True while the player still means to play something: playing, buffering, or recovering from + * an error. The service must not stop itself (releasing the player) on such a player, whatever + * the media session says about it. + */ + fun isBusy(): Boolean = + recoveryJob?.isActive == true || + (player.playWhenReady && (player.playbackState == Player.STATE_BUFFERING || player.playbackState == Player.STATE_READY)) + fun handle(command: BackendCommand) { when (command) { is BackendCommand.Load -> { @@ -281,6 +388,7 @@ class ExoBackend( private fun cancelRecovery() { recoveryJob?.cancel() recoveryJob = null + recoveryKey = null recoveryAttempts = 0 } diff --git a/android/playback/src/main/java/app/hocket/playback/HocketStreamDataSource.kt b/android/playback/src/main/java/app/hocket/playback/HocketStreamDataSource.kt index c031e76..0c0e6b2 100644 --- a/android/playback/src/main/java/app/hocket/playback/HocketStreamDataSource.kt +++ b/android/playback/src/main/java/app/hocket/playback/HocketStreamDataSource.kt @@ -200,8 +200,10 @@ class CoreStreamDataSourceFactory( /** * `DefaultLoadErrorHandlingPolicy`, except that core stream failures a retry cannot fix fail at once: * an unknown/expired token (the core must resolve a fresh one), an offset past the end, and a core - * that has shut down or closed the handle. Media3's default retries everything but - * `FileNotFoundException`-typed errors, which would re-open an expired token three times. + * that has shut down. Media3's default retries everything but `FileNotFoundException`-typed errors, + * which would re-open an expired token three times. A handle the core closed under the player + * (idle for ten minutes while the player's buffer was full, or reaped) is retried: the retry opens a + * fresh handle at the same position, which the cache serves. */ class CoreStreamLoadErrorPolicy : DefaultLoadErrorHandlingPolicy() { override fun getRetryDelayMsFor(loadErrorInfo: LoadErrorHandlingPolicy.LoadErrorInfo): Long { @@ -214,7 +216,6 @@ class CoreStreamLoadErrorPolicy : DefaultLoadErrorHandlingPolicy() { CoreStreamException.Kind.UnknownToken, CoreStreamException.Kind.RangeNotSatisfiable, CoreStreamException.Kind.ShutDown, - CoreStreamException.Kind.Closed, ) fun coreStreamFailure(e: Throwable?): CoreStreamException? { diff --git a/android/playback/src/main/java/app/hocket/playback/Monitors.kt b/android/playback/src/main/java/app/hocket/playback/Monitors.kt index 69a3732..d885800 100644 --- a/android/playback/src/main/java/app/hocket/playback/Monitors.kt +++ b/android/playback/src/main/java/app/hocket/playback/Monitors.kt @@ -23,7 +23,12 @@ import java.security.MessageDigest * network id so transcoding profiles can vary per network. The id is a hash of the SSID when it is * readable (needs location permission on 8.1+), otherwise the transport type. */ -class NetworkMonitor(private val context: Context, private val dispatch: (Command) -> Unit) { +class NetworkMonitor( + private val context: Context, + private val dispatch: (Command) -> Unit, + /** Called with whether there is any connectivity, on every change of that (and once at start). */ + private val onConnectivity: (online: Boolean) -> Unit = {}, +) { private val cm = context.getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager private var last: NetworkState? = null @@ -78,9 +83,12 @@ class NetworkMonitor(private val context: Context, private val dispatch: (Comman private fun publish() { val state = current() - if (state != last) { + val was = last + if (state != was) { last = state dispatch(Commands.setNetworkState(state)) + val online = state.kind != NetworkKind.Offline + if (was == null || (was.kind != NetworkKind.Offline) != online) onConnectivity(online) } } } diff --git a/android/playback/src/main/java/app/hocket/playback/PlaybackService.kt b/android/playback/src/main/java/app/hocket/playback/PlaybackService.kt index b751745..f98d294 100644 --- a/android/playback/src/main/java/app/hocket/playback/PlaybackService.kt +++ b/android/playback/src/main/java/app/hocket/playback/PlaybackService.kt @@ -43,7 +43,8 @@ import kotlinx.coroutines.launch * explicitly in [onCreate]. Without that there was no notification, no lock-screen or quick-settings * controls, and no foreground promotion, so a backgrounded process could be killed mid-song. * When nothing has played for - * [IDLE_TIMEOUT_MS] and no UI client is bound, it stops itself. + * [IDLE_TIMEOUT_MS] and no UI client is bound, it stops itself; a player that is still buffering + * or recovering from a network error ([ExoBackend.isBusy]) keeps it, whatever the session says. * - Exposes a [LocalBinder] so the app process can obtain the [CoreHandle] by binding. Only the * app's own bind counts as a UI client: the service is exported for Media3, so the bind intent * carries [CoreHost.bindToken], which no other process can know. @@ -97,7 +98,7 @@ class PlaybackService : MediaLibraryService() { val player = CoreSessionPlayer(Looper.getMainLooper(), ::dispatch, media = browser, artwork = { t -> t.coverArt?.let { ArtworkProvider.uri(this, it) } }) bridge = MediaSessionBridge(this, player, ::dispatch, launch, browser, scope, ExternalControl(this)) addSession(bridge.session) - network = NetworkMonitor(this, ::dispatch) + network = NetworkMonitor(this, ::dispatch, onConnectivity = { online -> main.post { backend.onConnectivityChanged(online) } }) battery = BatterySaverMonitor(this, ::dispatch) // Subscribed before the snapshot is requested, so its `Started` cannot be missed. scope.launch(start = CoroutineStart.UNDISPATCHED) { core.events.collect { event -> main.post { onEvent(event) } } } @@ -173,9 +174,12 @@ class PlaybackService : MediaLibraryService() { idleJob?.cancel() idleJob = scope.launch { delay(IDLE_TIMEOUT_MS) - if (!bridge.player.state.isPlaying && boundClients == 0) { + if (!bridge.player.state.isPlaying && boundClients == 0 && !backend.isBusy()) { Log.i(TAG, "Idle for ${IDLE_TIMEOUT_MS / 1000}s with no clients; stopping") pauseAllPlayersAndStopSelf() + } else if (boundClients == 0) { + // Still playing, buffering or recovering: look again later. + scheduleIdleStop() } } } @@ -207,7 +211,8 @@ class PlaybackService : MediaLibraryService() { override fun onTaskRemoved(rootIntent: Intent?) { // The UI unbinds shortly after (ForegroundBinder); with nothing playing the service then goes. - if (!bridge.player.state.isPlaying) pauseAllPlayersAndStopSelf() + // A player that is still buffering or recovering from a network error is not "nothing". + if (!bridge.player.state.isPlaying && !backend.isBusy()) pauseAllPlayersAndStopSelf() } override fun onDestroy() { diff --git a/android/playback/src/test/java/app/hocket/playback/ExoBackendTest.kt b/android/playback/src/test/java/app/hocket/playback/ExoBackendTest.kt index e9fb2a5..3de910e 100644 --- a/android/playback/src/test/java/app/hocket/playback/ExoBackendTest.kt +++ b/android/playback/src/test/java/app/hocket/playback/ExoBackendTest.kt @@ -17,6 +17,11 @@ import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.cancel import org.junit.After import androidx.media3.common.PlaybackException +import androidx.media3.datasource.DataSourceException +import androidx.media3.datasource.DataSpec +import androidx.media3.datasource.HttpDataSource +import androidx.media3.exoplayer.ExoPlaybackException +import app.hocket.core.CoreStreamException import org.junit.Assert.assertEquals import org.junit.Assert.assertFalse import org.junit.Assert.assertNull @@ -83,6 +88,61 @@ class ExoBackendTest { assertEquals("b", backend.player.currentMediaItem?.mediaId) } + private fun sourceError(cause: Throwable, code: Int): PlaybackException = + ExoPlaybackException.createForSource(java.io.IOException(cause), code) + + @Test + fun errorsALaterPrepareCanFixAreRecoverableTheRestAreNot() { + fun core(kind: CoreStreamException.Kind, status: Int? = null, code: Int = PlaybackException.ERROR_CODE_IO_UNSPECIFIED) = + sourceError(DataSourceException(CoreStreamException(kind, status), code), code) + // The connection, a stalled server, a handle the core closed: retried. + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.Network, code = PlaybackException.ERROR_CODE_IO_NETWORK_CONNECTION_FAILED))) + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.Status, 503, PlaybackException.ERROR_CODE_IO_BAD_HTTP_STATUS))) + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.Closed))) + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.UnknownHandle))) + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.TooManyHandles))) + assertTrue(ExoBackend.isRecoverable(core(CoreStreamException.Kind.NoServer))) + assertTrue("a stuck player is re-prepared", ExoBackend.isRecoverable(PlaybackException.ERROR_CODE_TIMEOUT)) + val http5xx = HttpDataSource.InvalidResponseCodeException(502, "Bad Gateway", null, emptyMap(), DataSpec(android.net.Uri.parse("https://music.example/s")), ByteArray(0)) + assertTrue(ExoBackend.isRecoverable(sourceError(http5xx, PlaybackException.ERROR_CODE_IO_BAD_HTTP_STATUS))) + // The media, the token, a 4xx, a core that is gone: fatal, the core decides. + assertFalse(ExoBackend.isRecoverable(core(CoreStreamException.Kind.Status, 404, PlaybackException.ERROR_CODE_IO_BAD_HTTP_STATUS))) + assertFalse(ExoBackend.isRecoverable(core(CoreStreamException.Kind.UnknownToken, code = PlaybackException.ERROR_CODE_IO_FILE_NOT_FOUND))) + assertFalse(ExoBackend.isRecoverable(core(CoreStreamException.Kind.ShutDown))) + assertFalse(ExoBackend.isRecoverable(core(CoreStreamException.Kind.ErrorEnvelope, code = PlaybackException.ERROR_CODE_IO_BAD_HTTP_STATUS))) + assertFalse(ExoBackend.isRecoverable(ExoPlaybackException.createForUnexpected(RuntimeException("decoder"), PlaybackException.ERROR_CODE_DECODING_FAILED))) + } + + @Test + fun aStallIsPausedOnlyWhenThePlayerNoLongerMeansToPlay() { + // Rebuffering, a load being retried, an error being recovered from: still buffering. + assertEquals(ExoBackend.Stall.Buffering, ExoBackend.stall(true, androidx.media3.common.Player.STATE_BUFFERING)) + assertEquals(ExoBackend.Stall.Nothing, ExoBackend.stall(true, androidx.media3.common.Player.STATE_IDLE)) + assertEquals(ExoBackend.Stall.Nothing, ExoBackend.stall(true, androidx.media3.common.Player.STATE_ENDED)) + assertEquals(ExoBackend.Stall.Nothing, ExoBackend.stall(true, androidx.media3.common.Player.STATE_READY)) + // A pause, becoming noisy, a lasting focus loss. + assertEquals(ExoBackend.Stall.Paused, ExoBackend.stall(false, androidx.media3.common.Player.STATE_READY)) + assertEquals(ExoBackend.Stall.Paused, ExoBackend.stall(false, androidx.media3.common.Player.STATE_BUFFERING)) + assertEquals(ExoBackend.Stall.Nothing, ExoBackend.stall(false, androidx.media3.common.Player.STATE_ENDED)) + assertEquals(ExoBackend.Stall.Nothing, ExoBackend.stall(false, androidx.media3.common.Player.STATE_IDLE)) + } + + @Test + fun thePlayerIsBusyWhileItStillMeansToPlayAndIdleAfterStop() { + assertFalse(backend.isBusy()) + backend.handle(BackendCommand.Load(BackendCommandLoadInner(source("a"), null, 0u, true))) + assertTrue("loading with play: the service must not stop it", backend.isBusy()) + backend.handle(BackendCommand.Pause) + assertFalse("paused: idle as far as the service is concerned", backend.isBusy()) + backend.handle(BackendCommand.Play) + assertTrue(backend.isBusy()) + backend.handle(BackendCommand.Stop) + assertFalse(backend.isBusy()) + backend.onConnectivityChanged(false) + backend.onConnectivityChanged(true) + assertFalse("connectivity with no pending recovery changes nothing", backend.isBusy()) + } + @Test fun networkErrorsAreRetriedWithBackoffThenGiveUp() { assertTrue(ExoBackend.isRecoverable(PlaybackException.ERROR_CODE_IO_NETWORK_CONNECTION_FAILED)) diff --git a/android/playback/src/test/java/app/hocket/playback/HocketStreamDataSourceTest.kt b/android/playback/src/test/java/app/hocket/playback/HocketStreamDataSourceTest.kt index 60d517c..67d8919 100644 --- a/android/playback/src/test/java/app/hocket/playback/HocketStreamDataSourceTest.kt +++ b/android/playback/src/test/java/app/hocket/playback/HocketStreamDataSourceTest.kt @@ -158,6 +158,16 @@ class HocketStreamDataSourceTest { HocketStreamDataSource.toDataSourceException(CoreStreamException(CoreStreamException.Kind.Network, detail = "reset")), 1, ) assertTrue("a network hiccup is retried", policy.getRetryDelayMsFor(network) != C.TIME_UNSET) + for (kind in listOf(CoreStreamException.Kind.Closed, CoreStreamException.Kind.UnknownHandle, CoreStreamException.Kind.TooManyHandles)) { + val gone = androidx.media3.exoplayer.upstream.LoadErrorHandlingPolicy.LoadErrorInfo( + info.loadEventInfo, info.mediaLoadData, HocketStreamDataSource.toDataSourceException(CoreStreamException(kind)), 1, + ) + assertTrue("a handle the core closed is re-opened by a retry ($kind)", policy.getRetryDelayMsFor(gone) != C.TIME_UNSET) + } + val dead = androidx.media3.exoplayer.upstream.LoadErrorHandlingPolicy.LoadErrorInfo( + info.loadEventInfo, info.mediaLoadData, HocketStreamDataSource.toDataSourceException(CoreStreamException(CoreStreamException.Kind.ShutDown)), 1, + ) + assertEquals("a core that shut down is not retried", C.TIME_UNSET, policy.getRetryDelayMsFor(dead)) source.close() // nothing opened: no close call, no transfer end assertFalse(streams.calls.any { it.startsWith("close") }) } diff --git a/crates/hocket-core/src/audio/scripted.rs b/crates/hocket-core/src/audio/scripted.rs index b2bdabe..7d42824 100644 --- a/crates/hocket-core/src/audio/scripted.rs +++ b/crates/hocket-core/src/audio/scripted.rs @@ -108,6 +108,9 @@ struct State { /// Report a gapless boundary the way Media3 does: `TransitionedToNext` /// alone, with no `Ended` for the item that finished. transition_only: bool, + /// Lose the preloaded follow-up at the boundary: `Ended` and nothing + /// after it (a `SetNext` the player had not applied yet). + drop_next: bool, } /// See the module docs. @@ -163,6 +166,12 @@ impl ScriptedBackend { self.state.lock().transition_only = on; } + /// At the next boundary, end without moving on to the follow-up that + /// was set (as a player does when it had already run out of items). + pub fn set_drop_next(&self, on: bool) { + self.state.lock().drop_next = on; + } + /// The device list [`PlaybackBackend::output_devices`] returns. pub fn set_devices(&self, devices: Vec) { self.state.lock().devices = devices.clone(); @@ -270,6 +279,9 @@ impl ScriptedBackend { if !(st.transition_only && has_next) { out.push(BackendReport::Ended { key: ended_key }); } + if st.drop_next { + st.next = None; + } match st.next.take() { Some(next) if !st.failing.contains(&next.key) diff --git a/crates/hocket-core/src/core/actor.rs b/crates/hocket-core/src/core/actor.rs index 1b53950..d8b514a 100644 --- a/crates/hocket-core/src/core/actor.rs +++ b/crates/hocket-core/src/core/actor.rs @@ -46,6 +46,9 @@ pub(crate) const ENGINE_TICK_MS: f64 = 1_000.0; pub(crate) const AUTOPLAY_BATCH: u32 = 10; /// Fatal load failures of one item before it is marked unavailable. pub(crate) const LOAD_RETRIES: u32 = 1; +/// How long an item adopted at a gapless boundary waits for the backend's +/// `TransitionedToNext` before it is loaded explicitly. +pub(crate) const ADOPTION_TIMEOUT_MS: f64 = 2_000.0; /// Consecutive unavailable items before playback stops trying. pub(crate) const MAX_CONSECUTIVE_SKIPS: u32 = 3; /// Artwork size for the media session (small in battery saver). @@ -856,6 +859,7 @@ impl Actor { self.engine_input(Input::Tick); } self.prefetch_tick(now); + self.playback_tick(now); self.prime_next(); self.cache_tick(now); // Sleep timer. diff --git a/crates/hocket-core/src/core/handlers/playback.rs b/crates/hocket-core/src/core/handlers/playback.rs index 95df5d3..7d8f79c 100644 --- a/crates/hocket-core/src/core/handlers/playback.rs +++ b/crates/hocket-core/src/core/handlers/playback.rs @@ -7,7 +7,8 @@ use crate::audio::sleep::SleepAction; use crate::connect::engine::Input; use crate::connect::wire::{SessionOp, TransportCommand}; use crate::core::actor::{ - Actor, LOAD_RETRIES, MAX_CONSECUTIVE_SKIPS, MEDIA_SESSION_ART, MEDIA_SESSION_ART_SMALL, + Actor, ADOPTION_TIMEOUT_MS, LOAD_RETRIES, MAX_CONSECUTIVE_SKIPS, MEDIA_SESSION_ART, + MEDIA_SESSION_ART_SMALL, }; use crate::core::state::{pending_scrobbles_key, PendingScrobble}; use crate::core::Internal; @@ -219,6 +220,7 @@ impl Actor { self.playback.next = None; self.playback.next_doc_key = None; self.playback.awaiting_transition = false; + self.playback.adopted_at = Some(now); self.playback.loaded = true; } else { let Some(source) = self.media_source_for(&item.key, &track) else { @@ -238,6 +240,7 @@ impl Actor { self.playback.position_ms = position_ms; self.playback.position_at = now; self.playback.awaiting_transition = false; + self.playback.adopted_at = None; self.playback.loaded = true; self.playback.next = None; self.playback.next_doc_key = None; @@ -542,6 +545,36 @@ impl Actor { self.playback.next = None; self.playback.next_doc_key = None; self.playback.awaiting_transition = false; + self.playback.adopted_at = None; + } + + /// Every tick: an item adopted at a gapless boundary that the backend + /// never confirmed is loaded explicitly (see `Playback::adopted_at`). + pub(crate) fn playback_tick(&mut self, now: f64) { + let Some(at) = self.playback.adopted_at else { + return; + }; + if now - at < ADOPTION_TIMEOUT_MS { + return; + } + self.playback.adopted_at = None; + if !self.playback.loaded || !self.owns_transport() { + return; + } + let Some(item) = self.doc().and_then(|d| d.current.clone()) else { + return; + }; + if self.playback.doc_key.as_ref() != Some(&item.key) { + return; + } + self.log( + "info", + "the backend did not move on to the preloaded item; loading it", + ); + let play = self.playback.want_playing || self.playback.playing; + self.playback.next = None; + self.playback.next_doc_key = None; + self.load_item(&item, 0, play, None); } pub(crate) fn set_playing(&mut self, playing: bool) { @@ -768,6 +801,7 @@ impl Actor { // in `load_item`. Late arrival means the doc did not follow // (repeat one, queue edited): reload what the document says. if self.playback.backend_key.as_ref() == Some(&key) { + self.playback.adopted_at = None; self.playback.playing = true; self.playback.position_ms = 0; self.playback.position_at = self.now(); @@ -804,9 +838,17 @@ impl Actor { return; } if is_next && !current(&key, self) { - // The preloaded item failed to start after the current one ended. - self.playback.awaiting_transition = false; + // The preloaded follow-up failed. Forget it: it is resolved + // afresh (with the usual retry and skip) when it comes up. self.playback.next = None; + self.playback.next_doc_key = None; + if !self.playback.awaiting_transition { + // The current item is still playing: leave it alone. + return; + } + // The current one ended and the follow-up never started: + // move on explicitly. + self.playback.awaiting_transition = false; if let Some(item) = self.doc().and_then(|d| d.current.clone()) { if self.playback.doc_key.as_ref() == Some(&item.key) { self.playback.load_failures += 1; diff --git a/crates/hocket-core/src/core/handlers/prefetch.rs b/crates/hocket-core/src/core/handlers/prefetch.rs index 3e33c94..8355a50 100644 --- a/crates/hocket-core/src/core/handlers/prefetch.rs +++ b/crates/hocket-core/src/core/handlers/prefetch.rs @@ -22,10 +22,14 @@ //! a playback stream is pulling from the server). A queue change is //! debounced ([`PREFETCH_DEBOUNCE_MS`]); a fetch no longer among the next //! two is cancelled (its temp file removed) and the new ones start. +//! - A fetch that fails (or completes without leaving a cache entry: no +//! room, a write error) is tried again after a growing backoff +//! ([`PREFETCH_RETRY_MS`]) while the track stays wanted, and at once when +//! the network comes back; a queue change clears the bookkeeping. //! //! [`Actor::start_prefetch`] is the reusable entry point for one track. -use std::collections::HashSet; +use std::collections::HashMap; use tokio_util::sync::CancellationToken; @@ -43,6 +47,23 @@ pub const PREFETCH_AHEAD: usize = 2; pub const PREFETCH_DEBOUNCE_MS: f64 = 2_000.0; /// Share of the cache budget prefetch may fill. pub const PREFETCH_BUDGET_SHARE: f64 = 0.25; +/// Backoff before a failed prefetch of a track is tried again, per attempt +/// (the last one repeats). +pub const PREFETCH_RETRY_MS: [f64; 5] = [15_000.0, 30_000.0, 60_000.0, 120_000.0, 300_000.0]; + +/// A track whose fetch failed while it stayed wanted. +#[derive(Debug, Clone, PartialEq)] +pub(crate) struct PrefetchFailure { + pub attempts: u32, + /// When it may be tried again (infinite while a retry runs). + pub retry_at: f64, +} + +impl PrefetchFailure { + fn backoff_ms(attempts: u32) -> f64 { + PREFETCH_RETRY_MS[(attempts as usize).min(PREFETCH_RETRY_MS.len() - 1)] + } +} /// How one prefetch ended. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -71,9 +92,9 @@ pub(crate) struct PrefetchState { pub targets: Vec, pub running: Option, pub generation: u64, - /// Tracks whose fetch failed while they stayed wanted (not retried - /// until the wanted set changes). - pub failed: HashSet, + /// Tracks whose fetch failed while they stayed wanted, with their + /// retry schedule (cleared when the wanted set changes). + pub failed: HashMap, } impl Actor { @@ -105,6 +126,24 @@ impl Actor { self.apply_prefetch(); } } + // A failed fetch whose backoff has elapsed. + if self.prefetch.running.is_none() + && self.prefetch.changed_at.is_none() + && self.prefetch.failed.values().any(|f| f.retry_at <= now) + { + self.apply_prefetch(); + } + } + + /// The network is back (or changed): failed fetches are due again now + /// rather than after their backoff. + pub(crate) fn prefetch_retry_now(&mut self) { + let now = self.now(); + for f in self.prefetch.failed.values_mut() { + if f.retry_at.is_finite() { + f.retry_at = now; + } + } } /// Whether this device prefetches audio right now. @@ -232,9 +271,13 @@ impl Actor { /// Act on the wanted set: protect it, cancel a fetch that left it, /// start the next one needed. pub(crate) fn apply_prefetch(&mut self) { + let now = self.now(); let targets = self.prefetch.wanted.clone(); self.prefetch.targets = targets.clone(); self.protect_loaded_tracks(); + self.prefetch + .failed + .retain(|id, _| targets.iter().any(|(_, t)| t == id)); if let Some(r) = &self.prefetch.running { if !targets.contains(&r.key) { r.cancel.cancel(); @@ -245,18 +288,28 @@ impl Actor { return; } for (_, id) in targets { - if self.prefetch.failed.contains(&id) { + if self + .prefetch + .failed + .get(&id) + .is_some_and(|f| f.retry_at > now) + { continue; } let Some(track) = self.track_or_bare(&id) else { continue; }; if !self.prefetch_needed(&track) { + self.prefetch.failed.remove(&id); continue; } if self.start_prefetch(&track, None) { + if let Some(f) = self.prefetch.failed.get_mut(&id) { + f.retry_at = f64::INFINITY; + } return; } + self.prefetch.failed.remove(&id); } } @@ -328,14 +381,49 @@ impl Actor { return; } self.prefetch.running = None; - if outcome == PrefetchOutcome::Failed { - self.log("debug", format!("prefetch of {track_id} failed")); - self.prefetch.failed.insert(track_id); + // A read that reached the end without completing the entry (no + // room, a write error, a stream the cache refused) would otherwise + // start over at once, for as long as the track stays wanted. + let uncached = outcome == PrefetchOutcome::Done + && self + .track_or_bare(&track_id) + .is_some_and(|t| self.prefetch_needed(&t)); + match outcome { + PrefetchOutcome::Failed => self.note_prefetch_failure(&track_id, "failed"), + PrefetchOutcome::Done if uncached => { + self.note_prefetch_failure(&track_id, "left no cache entry") + } + PrefetchOutcome::Done => { + self.prefetch.failed.remove(&track_id); + } + PrefetchOutcome::Cancelled => {} } self.apply_prefetch(); self.prime_next(); } + fn note_prefetch_failure(&mut self, track_id: &str, what: &str) { + let now = self.now(); + let f = self + .prefetch + .failed + .entry(track_id.to_string()) + .or_insert(PrefetchFailure { + attempts: 0, + retry_at: now, + }); + let delay = PrefetchFailure::backoff_ms(f.attempts); + f.attempts += 1; + f.retry_at = now + delay; + self.log( + "debug", + format!( + "prefetch of {track_id} {what}; trying again in {}s", + delay / 1000.0 + ), + ); + } + /// Stop any prefetch (shutdown). pub(crate) fn cancel_prefetch(&mut self) { if let Some(r) = self.prefetch.running.take() { diff --git a/crates/hocket-core/src/core/handlers/servers.rs b/crates/hocket-core/src/core/handlers/servers.rs index d916121..fc356b6 100644 --- a/crates/hocket-core/src/core/handlers/servers.rs +++ b/crates/hocket-core/src/core/handlers/servers.rs @@ -704,6 +704,8 @@ impl Actor { } self.mark_prefetch_check(); if state.kind != NetworkKind::Offline && (was_offline || self.network.is_some()) { + // Prefetches that failed on the old network are due now. + self.prefetch_retry_now(); self.last_outbox_retry = self.now(); self.schedule_flush(); self.maybe_start_sync(false, false); diff --git a/crates/hocket-core/src/core/state.rs b/crates/hocket-core/src/core/state.rs index 91f050e..e774027 100644 --- a/crates/hocket-core/src/core/state.rs +++ b/crates/hocket-core/src/core/state.rs @@ -80,6 +80,14 @@ pub(crate) struct Playback { pub fade_gain: Option, /// Waiting for a gapless `TransitionedToNext` for this doc key. pub awaiting_transition: bool, + /// When the document's current item was adopted as the backend's + /// preloaded follow-up without a `Load` (a gapless boundary), until + /// the backend confirms with `TransitionedToNext`. Past + /// [`crate::core::actor::ADOPTION_TIMEOUT_MS`] the item is loaded + /// explicitly: a backend that ended without a follow-up in place (a + /// `SetNext` still in flight, a preload it dropped) would otherwise + /// leave the session on an item nothing plays. + pub adopted_at: Option, /// Silenced by a transient audio-focus loss: the backend still means to /// play and resumes on its own (reporting `Playing`) when focus returns. pub focus_suspended: bool, diff --git a/crates/hocket-core/src/core/stream_reader.rs b/crates/hocket-core/src/core/stream_reader.rs index 7492b61..019fe12 100644 --- a/crates/hocket-core/src/core/stream_reader.rs +++ b/crates/hocket-core/src/core/stream_reader.rs @@ -73,6 +73,12 @@ pub const MAX_HANDLES: usize = 32; pub const HANDLE_IDLE_MS: f64 = 10.0 * 60.0 * 1000.0; /// Largest chunk one read returns. pub const MAX_READ: usize = 256 * 1024; +/// How long a read waits for the server's next bytes (or a fetch for its +/// response) before the handle fails with a network error. Counted only +/// while a read is waiting: a player that has filled its buffer and stops +/// reading for minutes (Media3 loads in bursts) must not find the stream +/// timed out behind its back, which the HTTP client's own read timeout did. +pub const UPSTREAM_WAIT: std::time::Duration = std::time::Duration::from_secs(60); #[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] pub enum StreamError { @@ -144,7 +150,9 @@ pub trait StreamUpstream: Send + Sync + 'static { /// Production upstream: reqwest + rustls, no redirects (a redirect would /// carry the auth query elsewhere), no overall timeout (a track streams for -/// as long as it plays) but connect and read timeouts. +/// as long as it plays) and no read timeout either: the reader applies +/// [`UPSTREAM_WAIT`] only while it is actually waiting for bytes, whereas +/// the client's read timeout also ran while nobody read the body. pub struct ReqwestUpstream { client: reqwest::Client, } @@ -154,7 +162,6 @@ impl ReqwestUpstream { let client = reqwest::Client::builder() .user_agent(concat!("hocket/", env!("CARGO_PKG_VERSION"))) .connect_timeout(std::time::Duration::from_secs(15)) - .read_timeout(std::time::Duration::from_secs(60)) .redirect(reqwest::redirect::Policy::none()) .use_rustls_tls() .build() @@ -924,7 +931,10 @@ impl CacheSource { self.seg = Some(seg); return Err(StreamError::Closed); } - c = seg.body.next() => c, + c = tokio::time::timeout(UPSTREAM_WAIT, seg.body.next()) => match c { + Ok(c) => c, + Err(_) => Some(Err("no data from the server for 60s".into())), + }, }; match next { Some(Ok(chunk)) => { @@ -1064,10 +1074,12 @@ impl CacheSource { let fetch = self.sh.upstream.fetch(UpstreamRequest { url, range }); let up = tokio::select! { _ = cancel.cancelled() => return Err(StreamError::Closed), - r = fetch => r.map_err(|e| { - tracing::warn!(target: "hocket_core", error = %e, track = %self.token.track_id, "stream upstream"); - StreamError::Network(e) - })?, + r = tokio::time::timeout(UPSTREAM_WAIT, fetch) => r + .unwrap_or_else(|_| Err("no response from the server for 60s".into())) + .map_err(|e| { + tracing::warn!(target: "hocket_core", error = %e, track = %self.token.track_id, "stream upstream"); + StreamError::Network(e) + })?, }; if up.status == 416 { return Err(StreamError::RangeNotSatisfiable); diff --git a/crates/hocket-core/tests/actor_playback.rs b/crates/hocket-core/tests/actor_playback.rs index f80e7ef..d59a58e 100644 --- a/crates/hocket-core/tests/actor_playback.rs +++ b/crates/hocket-core/tests/actor_playback.rs @@ -506,3 +506,68 @@ async fn transient_focus_loss_resumes_when_focus_returns() { ); assert!(!t.backend.is_playing()); } + +/// A backend that ends an item without moving on to the follow-up it was +/// given (a `SetNext` it had not applied yet, a preload it lost): the +/// session still adopts the follow-up at the boundary, and when no +/// `TransitionedToNext` confirms it, the core loads it explicitly instead of +/// sitting on an item nothing plays. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn an_adopted_item_the_backend_never_starts_is_loaded_explicitly() { + let t = synced_core("i").await; + t.backend.set_drop_next(true); + let sid = t.server_id.clone(); + t.run(Command::PlayTracks { + server_id: sid, + track_ids: vec!["t0".into(), "t1".into(), "t2".into()], + start_index: 0, + label: "Sel".into(), + shuffle: false, + }) + .await; + t.run_for(500.0).await; + assert_eq!(t.current_track_id().await.as_deref(), Some("t0")); + t.backend.clear_log(); + + // t0 (200 s) ends; the backend reports `Ended` and nothing after it. + t.run_for(200_500.0).await; + assert_eq!(t.current_track_id().await.as_deref(), Some("t1")); + assert!( + !t.backend + .log() + .iter() + .any(|c| matches!(c, ScriptedCall::Load { .. })), + "adopted first, not reloaded at once" + ); + assert!(!t.backend.is_playing(), "the backend has nothing loaded"); + + // After the adoption timeout the core loads what the document says. + t.run_for(3_000.0).await; + // The explicit load carries the document's key for t1 (the preload was + // keyed before the item was materialised). + let t1_keys: Vec = t + .sources + .lock() + .iter() + .filter(|s| s.track.id == "t1") + .map(|s| s.key.clone()) + .collect(); + assert!(!t1_keys.is_empty(), "t1 was resolved"); + assert!( + t.backend.log().iter().any( + |c| matches!(c, ScriptedCall::Load { key, play: true, .. } if t1_keys.contains(key)) + ), + "{:?}", + t.backend.log() + ); + assert!(t.backend.is_playing()); + assert_eq!(t.current_track_id().await.as_deref(), Some("t1")); + let snap = t.snapshot().await; + assert!(snap.transport.position.is_playing); + assert!(snap.transport.position.position_ms < 5_000); + + // And playback carries on the same way through the next boundary. + t.run_for(205_000.0).await; + assert_eq!(t.current_track_id().await.as_deref(), Some("t2")); + assert!(t.backend.is_playing()); +} diff --git a/crates/hocket-core/tests/actor_prefetch.rs b/crates/hocket-core/tests/actor_prefetch.rs index 59409fc..cbefb0b 100644 --- a/crates/hocket-core/tests/actor_prefetch.rs +++ b/crates/hocket-core/tests/actor_prefetch.rs @@ -300,3 +300,52 @@ async fn only_the_transport_owner_prefetches_and_a_handoff_moves_it() { a.core.shutdown().await; b.core.shutdown().await; } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_failed_prefetch_is_tried_again_after_a_backoff() { + let t = started("retry", seeded_server(4, 100.0)).await; + t.upstream + .set_behaviour("t1", UpstreamBehaviour::ErrorEnvelope); + play_album(&t, 4, 0).await; + // t1 fails, t2 is fetched anyway. + wait_cached(&t, "t2").await; + run_real(&t, 3000.0).await; + assert_eq!(offline(&t, "t1").await, OfflineState::None); + assert_eq!(t.upstream.stream_requests("t1"), 1, "one failed attempt"); + // The server recovers: nothing before the first backoff has elapsed... + t.upstream.set_behaviour("t1", UpstreamBehaviour::Serve); + run_real(&t, 5_000.0).await; + assert_eq!( + t.upstream.stream_requests("t1"), + 1, + "not before the backoff" + ); + // ...then the retry lands and the track is cached. + wait_cached(&t, "t1").await; + assert_eq!(t.upstream.stream_requests("t1"), 2); + t.core.shutdown().await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn the_network_coming_back_retries_a_failed_prefetch_at_once() { + let t = started("netretry", seeded_server(4, 100.0)).await; + t.upstream + .set_behaviour("t1", UpstreamBehaviour::ErrorEnvelope); + play_album(&t, 4, 0).await; + wait_cached(&t, "t2").await; + run_real(&t, 3000.0).await; + assert_eq!(t.upstream.stream_requests("t1"), 1); + t.upstream.set_behaviour("t1", UpstreamBehaviour::Serve); + // A network change (Wi-Fi again, another id) makes the failed fetch due now. + t.run(Command::SetNetworkState { + state: NetworkState { + kind: NetworkKind::Wifi, + metered: false, + network_id: Some("home".into()), + }, + }) + .await; + wait_cached(&t, "t1").await; + assert_eq!(t.upstream.stream_requests("t1"), 2); + t.core.shutdown().await; +}