diff --git a/app/src/main/java/org/libremediaconverter/work/ConversionWorker.kt b/app/src/main/java/org/libremediaconverter/work/ConversionWorker.kt index 5d6994d..a31249f 100644 --- a/app/src/main/java/org/libremediaconverter/work/ConversionWorker.kt +++ b/app/src/main/java/org/libremediaconverter/work/ConversionWorker.kt @@ -12,6 +12,7 @@ import androidx.work.WorkerParameters import androidx.work.hasKeyWithValueOfType import androidx.work.workDataOf import com.arthenica.ffmpegkit.FFmpegKitConfig +import com.google.common.util.concurrent.ListenableFuture import kotlinx.coroutines.CancellationException import org.libremediaconverter.convert.ConversionDependencies import org.libremediaconverter.convert.InputQuery @@ -231,16 +232,58 @@ class ConversionWorker(context: Context, params: WorkerParameters) : CoroutineWo private var lastNotified = 0L + /** + * Reports how far along the conversion is, through WorkManager rather than around it. + * + * This used to call `NotificationManager.notify(NOTIFICATION_ID, …)` directly — on the very id + * WorkManager owns through [setForeground], with a notification built `setOngoing(true)`. Two + * owners of one id is a race, and on a Pixel 10 Pro XL it was lost on attempt 3 of 12 while + * cancelling a `BEST`-tier job: WorkManager tore the notification down at +300 ms and a tick + * still in flight put it back at +700 ms, where it stayed for ten minutes with no app process + * left to cancel it. The orphan's record carried `flags=ONGOING_EVENT|ONLY_ALERT_ONCE` and no + * `FOREGROUND_SERVICE` — which is what identifies the poster, since WorkManager's own goes out + * with that flag. + * + * Two things stop that, and they are independent on purpose: + * + * - **Nothing is published once the worker is stopped.** A tick arriving after the stop has + * nobody left to report to, and posting one is exactly the resurrection above. + * - **What is published goes through [setForegroundAsync].** The foreground notification is + * WorkManager's to post and to withdraw; sharing the id with it was the defect, not merely + * the mechanism of it. `setForegroundAsync` also refuses on its own for work that has + * already finished, which closes the window `isStopped` can only narrow. + * + * The `~1/sec` throttle is unchanged in effect and unchanged in reason: FFmpeg's statistics + * callback and Media3's progress polling both fire several times a second, and pushing every + * one of them janks the system UI. Routing them through WorkManager does not make them cheap. + * + * The future is deliberately not awaited — this is called from an engine callback, which is + * not a coroutine — and both of its failure modes are benign, so a refusal is logged rather + * than propagated. The initial [setForeground] in [doWork] is a different matter entirely and + * stays a suspending call inside the `try`: a denial *there* is what [FailureOutcome] turns + * into a retry rather than a terminal failure. + */ private fun publishProgress(displayName: String, percent: Int) { + if (isStopped) return setProgressAsync(workDataOf(KEY_PROGRESS to percent)) - // Throttle to ~1/sec: progress arrives several times a second and pushing every - // update janks the system UI. val now = System.currentTimeMillis() - if (now - lastNotified >= NOTIFICATION_INTERVAL_MS) { - lastNotified = now - applicationContext.getSystemService(android.app.NotificationManager::class.java) - .notify(NOTIFICATION_ID, notifications.build(id, displayName, percent)) - } + if (now - lastNotified < NOTIFICATION_INTERVAL_MS) return + lastNotified = now + val posted = setForegroundAsync(foregroundInfo(displayName, percent, indeterminate = false)) + posted.addListener({ logIfRefused(posted) }, Runnable::run) + } + + /** + * Notes a progress update WorkManager would not take, without making it the job's problem. + * + * The one that actually happens is a job finishing between [publishProgress]'s `isStopped` + * check and the update reaching the task thread: `WorkForegroundUpdater` refuses to post for + * work whose state is already terminal, which is the behaviour being relied on rather than + * worked around. Logged so it is greppable instead of vanishing into an unobserved future. + */ + private fun logIfRefused(posted: ListenableFuture) { + runCatching { posted.get() } + .onFailure { Log.d(TAG, "Progress update refused; the job is already finishing.", it) } } private fun isCancellation(e: Throwable): Boolean = e is CancellationException || isStopped diff --git a/app/src/test/java/org/libremediaconverter/work/ProgressNotificationTest.kt b/app/src/test/java/org/libremediaconverter/work/ProgressNotificationTest.kt new file mode 100644 index 0000000..8dc51d9 --- /dev/null +++ b/app/src/test/java/org/libremediaconverter/work/ProgressNotificationTest.kt @@ -0,0 +1,213 @@ +package org.libremediaconverter.work + +import android.app.Application +import android.app.Notification +import android.app.NotificationManager +import android.content.Context +import android.net.Uri +import androidx.media3.common.util.UnstableApi +import androidx.work.Data +import androidx.work.ForegroundInfo +import androidx.work.WorkInfo +import androidx.work.testing.TestForegroundUpdater +import androidx.work.testing.TestListenableWorkerBuilder +import androidx.work.workDataOf +import com.google.common.util.concurrent.ListenableFuture +import kotlinx.coroutines.runBlocking +import org.junit.After +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.libremediaconverter.convert.ConversionDependencies +import org.libremediaconverter.convert.SoftwareTranscoder +import org.libremediaconverter.convert.installTestWorkManager +import org.libremediaconverter.model.ConversionRequest +import org.libremediaconverter.model.DeviceCodecs +import org.libremediaconverter.model.EnginePreference +import org.libremediaconverter.model.InputProbe +import org.libremediaconverter.model.OutputFormat +import org.robolectric.RobolectricTestRunner +import org.robolectric.RuntimeEnvironment +import org.robolectric.Shadows.shadowOf +import java.io.File +import java.util.UUID + +/** + * That progress goes through WorkManager rather than around it. + * + * `publishProgress` called `NotificationManager.notify(1001, …)` directly, on the very id + * WorkManager owns through `setForeground`, with a notification built `setOngoing(true)`. Two + * owners of one id is a race, and on a Pixel 10 Pro XL it was lost on attempt 3 of 12 while + * cancelling a `BEST`-tier job: + * + * ``` + * attempt 3: terminal state = CANCELLED + * attempt 3: +300ms active=0 id1001=false ongoing=null <- WorkManager tore it down + * attempt 3: +700ms active=1 id1001=true ongoing=true <- a progress tick put it back + * attempt 3: +5000ms active=1 id1001=true ongoing=true + * ``` + * + * Still there ten minutes later with no app process at all. The record carried + * `flags=ONGOING_EVENT|ONLY_ALERT_ONCE` and **no `FOREGROUND_SERVICE`**, which is what proves the + * direct `notify` posted it rather than `setForeground` — WorkManager's own post carries that flag. + * Whether the orphan could be swiped away was never established and is not what these tests are + * about; the resurrection is, and it is reproduced. + */ +@UnstableApi +@RunWith(RobolectricTestRunner::class) +class ProgressNotificationTest { + + private lateinit var app: Application + private lateinit var updater: RecordingForegroundUpdater + + @Before + fun setUp() { + app = RuntimeEnvironment.getApplication() + updater = RecordingForegroundUpdater() + ConversionDependencies.publisher = { AlwaysRoomPublisher(app) } + ConversionDependencies.probe = { _, _ -> InputProbe() } + ConversionDependencies.deviceCodecs = { DeviceCodecs.PERMISSIVE } + // The notification's cancel action is a WorkManager PendingIntent, so the worker would + // fail for that reason rather than the one under test without this. + installTestWorkManager(app, Data.EMPTY) + } + + @After + fun tearDown() { + ConversionDependencies.reset() + } + + @Test + fun `progress on a running worker updates WorkManager's own notification, not one of ours`() { + runBlocking { workerReporting { onProgress -> onProgress(PERCENT) }.doWork() } + + // The positive half: the update really happened, on the id WorkManager is holding, and it + // carries the percentage. Asserting only that nothing was posted directly would pass just + // as well against a `publishProgress` that had been deleted. + val progressUpdates = updater.infos.drop(1) + assertEquals("one throttled progress update expected", 1, progressUpdates.size) + assertEquals(updater.infos.first().notificationId, progressUpdates.single().notificationId) + assertEquals(PERCENT, progressUpdates.single().notification.extras.getInt(Notification.EXTRA_PROGRESS)) + + // And the negative half: nothing reached the notification manager under its own steam. + // Asserted over the whole manager rather than one id, so a renamed constant cannot make + // this pass by asking about a notification nobody posts. + assertEquals("the worker must post no notification of its own", 0, postedNotifications()) + } + + @Test + fun `a progress update that lands after the worker is stopped puts nothing back`() { + // The device sequence, in one worker: WorkManager has torn the notification down, and a + // tick that was already in flight arrives afterwards. `lastNotified` is still 0 here, so + // this tick is one the throttle would have let through -- which is what makes the test + // about `isStopped` rather than about timing. + val worker = workerReporting { onProgress -> + stopped(WorkInfo.STOP_REASON_CANCELLED_BY_APP) + onProgress(PERCENT) + } + + runBlocking { worker.doWork() } + + assertEquals("a stopped worker must publish nothing", 1, updater.infos.size) + assertEquals("and must resurrect nothing", 0, postedNotifications()) + } + + @Test + fun `progress arriving several times a second is still throttled to one update`() { + // The throttle is not decoration: FFmpeg's statistics callback and Media3's progress + // polling both fire several times a second, and pushing every one of them janks the + // system UI. Routing progress through `setForeground` does not make that cheaper. + runBlocking { + workerReporting { onProgress -> repeat(TICKS) { onProgress(it) } }.doWork() + } + + assertTrue( + "$TICKS ticks inside one throttle window must not be ${updater.infos.size - 1} updates", + updater.infos.size - 1 == 1, + ) + } + + /** + * A worker routed to the software engine, whose engine is [report] and a written output. + * + * `FORCE_SOFTWARE` because it is the one preference that decides without consulting the input, + * and a `file://` URI because a `content://` one would send the worker through FFmpegKit's SAF + * bridge, which is native. [report] is handed the worker's own progress callback, and runs with + * the worker as its receiver so a test can stop it mid-transcode. + */ + private fun workerReporting(report: ConversionWorker.((Int) -> Unit) -> Unit): ConversionWorker { + val worker = TestListenableWorkerBuilder( + context = app, + inputData = workDataOf( + ConversionWorker.KEY_INPUT_URI to INPUT.toString(), + ConversionWorker.KEY_DISPLAY_NAME to DISPLAY_NAME, + ConversionWorker.KEY_SIZE_BYTES to INPUT_BYTES, + ConversionWorker.KEY_CONTAINER to SPEC.container.name, + ConversionWorker.KEY_VIDEO_CODEC to SPEC.videoCodec.name, + ConversionWorker.KEY_AUDIO_CODEC to SPEC.audioCodec.name, + ConversionWorker.KEY_ENGINE_PREFERENCE to EnginePreference.FORCE_SOFTWARE.name, + ), + runAttemptCount = 0, + ).setId(JOB_ID) + .setForegroundUpdater(updater) + .build() + // Resolved when the worker reaches the engine, so assigning after `build()` is in time. + ConversionDependencies.software = { ReportingTranscoder { onProgress -> worker.report(onProgress) } } + return worker + } + + /** Stops the worker the way WorkManager does, so `isStopped` becomes true. */ + private fun ConversionWorker.stopped(reason: Int) = stop(reason) + + private fun postedNotifications(): Int = shadowOf(app.getSystemService(NotificationManager::class.java)).size() + + private companion object { + val INPUT: Uri = Uri.parse("file:///tmp/holiday.mp4") + const val DISPLAY_NAME = "holiday.mp4" + const val INPUT_BYTES = 1024L + const val PERCENT = 42 + const val TICKS = 50 + val SPEC = OutputFormat.MP4_H265.spec + val JOB_ID: UUID = UUID.fromString("00000000-0000-4000-8000-000000000021") + } +} + +/** + * Records every [ForegroundInfo] the worker publishes, and otherwise behaves as the test default. + * + * Delegating to [TestForegroundUpdater] rather than hand-rolling a `ListenableFuture`: the + * worker awaits what this returns, so a future that never completes would hang the initial + * `setForeground` rather than test anything. + */ +private class RecordingForegroundUpdater : TestForegroundUpdater() { + val infos = mutableListOf() + + override fun setForegroundAsync( + context: Context, + id: UUID, + foregroundInfo: ForegroundInfo, + ): ListenableFuture { + infos += foregroundInfo + return super.setForegroundAsync(context, id, foregroundInfo) + } +} + +/** An engine that reports whatever [report] wants reported, then writes an output. */ +private class ReportingTranscoder(private val report: ((Int) -> Unit) -> Unit) : SoftwareTranscoder { + override suspend fun run( + request: ConversionRequest, + inputPath: String, + output: File, + durationMs: Long, + onProgress: (Int) -> Unit, + ) { + report(onProgress) + output.writeBytes(ByteArray(OUTPUT_BYTES)) + } + + private companion object { + const val OUTPUT_BYTES = 512 + } +}