diff --git a/app/src/main/java/com/pedro/streamer/rotation/CameraFragment.kt b/app/src/main/java/com/pedro/streamer/rotation/CameraFragment.kt index d28431966..06afde8ad 100644 --- a/app/src/main/java/com/pedro/streamer/rotation/CameraFragment.kt +++ b/app/src/main/java/com/pedro/streamer/rotation/CameraFragment.kt @@ -39,7 +39,7 @@ import com.pedro.extrasources.CameraXSource import com.pedro.library.base.StreamBase import com.pedro.library.base.recording.RecordController import com.pedro.library.generic.GenericStream -import com.pedro.library.util.BitrateAdapter +import com.pedro.library.util.QueueAwareBitrateAdapter import com.pedro.streamer.R import com.pedro.streamer.utils.PathUtils import com.pedro.streamer.utils.toast @@ -95,10 +95,8 @@ class CameraFragment: Fragment(), ConnectChecker { private val aBitrate = 128 * 1000 private var recordPath = "" //Bitrate adapter used to change the bitrate on fly depend of the bandwidth. - private val bitrateAdapter = BitrateAdapter { - genericStream.setVideoBitrateOnFly(it) - }.apply { - setMaxBitrate(vBitrate + aBitrate) + private val bitrateAdapter = QueueAwareBitrateAdapter(maxBitrate = vBitrate + aBitrate) { + genericStream.setVideoBitrateOnFly(it - aBitrate) } @SuppressLint("ClickableViewAccessibility") @@ -203,6 +201,7 @@ class CameraFragment: Fragment(), ConnectChecker { } override fun onConnectionStarted(url: String) { + bitrateAdapter.reset() } override fun onConnectionSuccess() { @@ -223,7 +222,7 @@ class CameraFragment: Fragment(), ConnectChecker { override fun onStreamingStats(report: StreamingStatsReport) { onMainThreadHandler { - bitrateAdapter.adaptBitrate(report.smoothedBitrate, genericStream.getStreamClient().hasCongestion()) + bitrateAdapter.onStreamingStats(report) if (report.throughput != Throughput.UNKNOWN) { txtBitrate.text = String.format( Locale.getDefault(), diff --git a/common/src/main/java/com/pedro/common/BitrateManager.kt b/common/src/main/java/com/pedro/common/BitrateManager.kt index 0afd4d822..081c04843 100644 --- a/common/src/main/java/com/pedro/common/BitrateManager.kt +++ b/common/src/main/java/com/pedro/common/BitrateManager.kt @@ -48,5 +48,6 @@ open class BitrateManager(private val bitrateChecker: BitrateChecker) { fun reset() { bitrate = 0 bitrateOld = 0 + timeStamp = TimeUtils.getCurrentTimeMillis() } } \ No newline at end of file diff --git a/common/src/main/java/com/pedro/common/Extensions.kt b/common/src/main/java/com/pedro/common/Extensions.kt index 9078e4a3a..761e2b60b 100644 --- a/common/src/main/java/com/pedro/common/Extensions.kt +++ b/common/src/main/java/com/pedro/common/Extensions.kt @@ -26,9 +26,11 @@ import android.media.MediaFormat import android.os.Build import android.os.Handler import android.os.Looper +import android.util.Log import android.view.Surface import androidx.annotation.RequiresApi import com.pedro.common.frame.MediaFrame +import kotlinx.coroutines.CoroutineDispatcher import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.withContext import kotlinx.io.IOException @@ -47,6 +49,7 @@ import java.util.concurrent.LinkedBlockingQueue import java.util.concurrent.ThreadPoolExecutor import java.util.concurrent.TimeUnit import kotlin.coroutines.Continuation +import kotlin.coroutines.CoroutineContext /** * Created by pedro on 3/11/23. @@ -366,3 +369,13 @@ fun ByteBuffer.clone(data: ByteArray): ByteBuffer { source.get(data, 0, length) return ByteBuffer.wrap(data, 0, length).slice() } + +@JvmOverloads +fun getSuspendContext(dispatcher: CoroutineDispatcher = Dispatchers.IO) = object: Continuation { + override val context: CoroutineContext + get() = dispatcher + + override fun resumeWith(result: Result) { + result.exceptionOrNull()?.let { Log.e("getSuspendContext", "Error", it) } + } +} \ No newline at end of file diff --git a/common/src/test/java/com/pedro/common/BitrateManagerTest.kt b/common/src/test/java/com/pedro/common/BitrateManagerTest.kt index defa7ffcf..8eea93982 100644 --- a/common/src/test/java/com/pedro/common/BitrateManagerTest.kt +++ b/common/src/test/java/com/pedro/common/BitrateManagerTest.kt @@ -23,6 +23,7 @@ import kotlinx.coroutines.test.resetMain import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.setMain import org.junit.After +import org.junit.Assert.assertEquals import org.junit.Assert.assertTrue import org.junit.Before import org.junit.Rule @@ -90,4 +91,36 @@ class BitrateManagerTest { val marginError = 20 assertTrue(expectedResult - marginError <= resultValue.firstValue && resultValue.firstValue <= expectedResult + marginError) } + + @Test + fun `GIVEN an idle instance WHEN reset and measure a second THEN report the real bitrate`() = runTest { + val bitrateManager = BitrateManager(connectChecker) + //the instance is created when the client is built, the stream may start much later + fakeTime += 300_000 + bitrateManager.reset() + + fakeTime += 1000 + bitrateManager.calculateBitrate(3_000_000L) + + val resultValue = argumentCaptor() + verify(connectChecker, times(1)).onNewBitrate(resultValue.capture()) + assertEquals(3_000_000L, resultValue.firstValue) + } + + @Test + fun `GIVEN a measured bitrate WHEN reset THEN start a new window instead of averaging the pause`() = runTest { + val bitrateManager = BitrateManager(connectChecker) + fakeTime += 1000 + bitrateManager.calculateBitrate(1_000_000L) + + //a reconnection: the sender is stopped for a while and started again + fakeTime += 120_000 + bitrateManager.reset() + fakeTime += 1000 + bitrateManager.calculateBitrate(2_000_000L) + + val resultValue = argumentCaptor() + verify(connectChecker, times(2)).onNewBitrate(resultValue.capture()) + assertEquals(2_000_000L, resultValue.secondValue) + } } diff --git a/encoder/src/main/java/com/pedro/encoder/input/video/Camera2ApiManager.kt b/encoder/src/main/java/com/pedro/encoder/input/video/Camera2ApiManager.kt index c80166cba..c1c3c1296 100644 --- a/encoder/src/main/java/com/pedro/encoder/input/video/Camera2ApiManager.kt +++ b/encoder/src/main/java/com/pedro/encoder/input/video/Camera2ApiManager.kt @@ -403,6 +403,13 @@ class Camera2ApiManager(context: Context) { } } + private fun clearRegions(): Array? { + val sensor = cameraCharacteristics?.secureGet( + CameraCharacteristics.SENSOR_INFO_ACTIVE_ARRAY_SIZE) ?: return null + return arrayOf(MeteringRectangle(0, 0, sensor.width(), sensor.height(), + MeteringRectangle.METERING_WEIGHT_DONT_CARE)) + } + /** * @param mode value from CameraCharacteristics.CONTROL_AWB_MODE_* */ @@ -413,8 +420,7 @@ class Camera2ApiManager(context: Context) { if (!modes.contains(mode)) return false val maxRegionsAwb = characteristics.secureGet(CameraCharacteristics.CONTROL_MAX_REGIONS_AWB) ?: 0 if (maxRegionsAwb > 0) { - val clearRect = MeteringRectangle(0, 0, 0, 0, MeteringRectangle.METERING_WEIGHT_DONT_CARE) - builderInputSurface.set(CaptureRequest.CONTROL_AWB_REGIONS, arrayOf(clearRect)) + clearRegions()?.let { builderInputSurface.set(CaptureRequest.CONTROL_AWB_REGIONS, it) } } builderInputSurface.set(CaptureRequest.CONTROL_AWB_MODE, mode) isAutoWhiteBalanceEnabled = applyRequest(builderInputSurface) @@ -462,8 +468,7 @@ class Camera2ApiManager(context: Context) { if (!modes.contains(CaptureRequest.CONTROL_AE_MODE_ON)) return false val maxRegionsAe = characteristics.secureGet(CameraCharacteristics.CONTROL_MAX_REGIONS_AE) ?: 0 if (maxRegionsAe > 0) { - val clearRect = MeteringRectangle(0, 0, 0, 0, MeteringRectangle.METERING_WEIGHT_DONT_CARE) - builderInputSurface.set(CaptureRequest.CONTROL_AE_REGIONS, arrayOf(clearRect)) + clearRegions()?.let { builderInputSurface.set(CaptureRequest.CONTROL_AE_REGIONS, it) } } builderInputSurface.set(CaptureRequest.CONTROL_AE_MODE, CaptureRequest.CONTROL_AE_MODE_ON) isAutoExposureEnabled = applyRequest(builderInputSurface) @@ -784,8 +789,7 @@ class Camera2ApiManager(context: Context) { builderInputSurface.setTag("") val maxRegionsAf = characteristics.secureGet(CameraCharacteristics.CONTROL_MAX_REGIONS_AF) ?: 0 if (maxRegionsAf > 0) { - val clearRect = MeteringRectangle(0, 0, 0, 0, MeteringRectangle.METERING_WEIGHT_DONT_CARE) - builderInputSurface.set(CaptureRequest.CONTROL_AF_REGIONS, arrayOf(clearRect)) + clearRegions()?.let { builderInputSurface.set(CaptureRequest.CONTROL_AF_REGIONS, it) } } builderInputSurface.set(CaptureRequest.CONTROL_AF_TRIGGER, CameraMetadata.CONTROL_AF_TRIGGER_CANCEL) builderInputSurface.set(CaptureRequest.CONTROL_AF_MODE, CaptureRequest.CONTROL_AF_MODE_OFF) diff --git a/library/src/main/java/com/pedro/library/base/Camera1Base.java b/library/src/main/java/com/pedro/library/base/Camera1Base.java index bf742cfcd..55095cb2e 100644 --- a/library/src/main/java/com/pedro/library/base/Camera1Base.java +++ b/library/src/main/java/com/pedro/library/base/Camera1Base.java @@ -293,8 +293,10 @@ public boolean prepareVideo(int width, int height, int fps, int bitrate, int iFr } FormatVideoEncoder formatVideoEncoder = glInterface == null ? FormatVideoEncoder.YUV420Dynamical : FormatVideoEncoder.SURFACE; - return videoEncoder.prepareVideoEncoder(width, height, fps, bitrate, rotation, iFrameInterval, + boolean result = videoEncoder.prepareVideoEncoder(width, height, fps, bitrate, rotation, iFrameInterval, formatVideoEncoder, profile, level); + forceFpsLimit(true); + return result; } /** diff --git a/library/src/main/java/com/pedro/library/base/Camera2Base.java b/library/src/main/java/com/pedro/library/base/Camera2Base.java index 609f1142d..c66dfc4ac 100644 --- a/library/src/main/java/com/pedro/library/base/Camera2Base.java +++ b/library/src/main/java/com/pedro/library/base/Camera2Base.java @@ -19,6 +19,8 @@ import android.content.Context; import android.graphics.Point; import android.hardware.camera2.CameraCharacteristics; +import android.hardware.camera2.CaptureRequest; +import android.hardware.camera2.TotalCaptureResult; import android.media.MediaCodec; import android.media.MediaFormat; import android.media.MediaRecorder; @@ -67,6 +69,8 @@ import java.util.Arrays; import java.util.List; +import kotlin.Unit; + /** * Wrapper to stream with camera2 api and microphone. Support stream with SurfaceView, TextureView, * OpenGlView(Custom SurfaceView that use OpenGl) and Context(background mode). All views use @@ -151,6 +155,29 @@ public void setCustomAudioEffect(CustomAudioEffect customAudioEffect) { microphoneManager.setCustomAudioEffect(customAudioEffect); } + public interface RequestListener { + void onRequest(CaptureRequest.Builder builder); + } + + public interface CaptureResultListener { + void onCaptureResult(TotalCaptureResult result); + } + + public boolean setCustomRequest(RequestListener listener) { + return cameraManager.setCustomRequest((builder) -> { + if (listener != null) listener.onRequest(builder); + return Unit.INSTANCE; + }); + } + + public void setCustomOnCaptureCompletedCallback(CaptureResultListener listener) { + cameraManager.setCustomOnCaptureCompletedCallback(listener == null ? null : + (session, request, result) -> { + listener.onCaptureResult(result); + return Unit.INSTANCE; + }); + } + /** * @param callback get fps while record or stream */ @@ -364,8 +391,10 @@ public boolean prepareVideo( iFrameInterval, FormatVideoEncoder.SURFACE, profile, level); if (!result) return false; } - return videoEncoder.prepareVideoEncoder(width, height, fps, bitrate, rotation, + boolean result = videoEncoder.prepareVideoEncoder(width, height, fps, bitrate, rotation, iFrameInterval, FormatVideoEncoder.SURFACE, profile, level); + forceFpsLimit(true); + return result; } public boolean prepareVideo( diff --git a/library/src/main/java/com/pedro/library/base/DisplayBase.java b/library/src/main/java/com/pedro/library/base/DisplayBase.java index 955ca2127..cabae172a 100644 --- a/library/src/main/java/com/pedro/library/base/DisplayBase.java +++ b/library/src/main/java/com/pedro/library/base/DisplayBase.java @@ -171,6 +171,7 @@ public boolean prepareVideo(int width, int height, int fps, int bitrate, int rot glStreamInterface.setIsPortrait(isPortrait); } } + forceFpsLimit(true); return videoInitialized; } diff --git a/library/src/main/java/com/pedro/library/base/FromFileBase.java b/library/src/main/java/com/pedro/library/base/FromFileBase.java index f6cca639e..ac742c7c0 100644 --- a/library/src/main/java/com/pedro/library/base/FromFileBase.java +++ b/library/src/main/java/com/pedro/library/base/FromFileBase.java @@ -214,6 +214,7 @@ private boolean finishPrepareVideo(int bitRate, int rotation, int profile, int if (!result) return false; result = videoDecoder.prepareVideo(videoEncoder.getInputSurface()); videoEnabled = result; + forceFpsLimit(true); return result; } diff --git a/library/src/main/java/com/pedro/library/base/StreamBase.kt b/library/src/main/java/com/pedro/library/base/StreamBase.kt index ff1332b08..84f6dc8de 100644 --- a/library/src/main/java/com/pedro/library/base/StreamBase.kt +++ b/library/src/main/java/com/pedro/library/base/StreamBase.kt @@ -540,10 +540,10 @@ abstract class StreamBase( audioSource.stop() glInterface.removeMediaCodecSurface() glInterface.removeMediaCodecRecordSurface() - if (!isOnPreview) glInterface.stop() videoEncoder.stop() videoEncoderRecord.stop() audioEncoder.stop() + if (!isOnPreview) glInterface.stop() if (!isRecording) recordController.resetFormats() } diff --git a/library/src/main/java/com/pedro/library/util/BitrateAdapter.java b/library/src/main/java/com/pedro/library/util/BitrateAdapter.java index a1f34d940..7736adbcc 100644 --- a/library/src/main/java/com/pedro/library/util/BitrateAdapter.java +++ b/library/src/main/java/com/pedro/library/util/BitrateAdapter.java @@ -21,6 +21,8 @@ */ public class BitrateAdapter { + private static final float MAX_BITRATE_TOLERANCE = 0.95f; + public interface Listener { void onBitrateAdapted(int bitrate); } @@ -74,7 +76,8 @@ public void adaptBitrate(long actualBitrate, boolean hasCongestion) { private int getBitrateAdapted(int bitrate) { if (bitrate >= maxBitrate) { //You have high speed and max bitrate. Keep max speed oldBitrate = maxBitrate; - } else if (bitrate <= oldBitrate * 0.9f) { //You have low speed and bitrate too high. Reduce bitrate by 10%. + } else if (bitrate <= oldBitrate * 0.9f || isStuckAtMaxBitrate(bitrate)) { + //You have low speed and bitrate too high. Reduce bitrate by 10%. oldBitrate = (int) (bitrate * decreaseRange); } else { //You have high speed and bitrate too low. Increase bitrate by 10%. oldBitrate = (int) (bitrate * increaseRange); @@ -83,8 +86,12 @@ private int getBitrateAdapted(int bitrate) { return oldBitrate; } + private boolean isStuckAtMaxBitrate(int bitrate) { + return oldBitrate >= maxBitrate && bitrate < maxBitrate * MAX_BITRATE_TOLERANCE; + } + private int getBitrateAdapted(int bitrate, boolean hasCongestion) { - if (bitrate >= maxBitrate) { //You have high speed and max bitrate. Keep max speed + if (bitrate >= maxBitrate && !hasCongestion) { //You have high speed and max bitrate. Keep max speed oldBitrate = maxBitrate; } else if (hasCongestion) { //You have low speed and bitrate too high. Reduce bitrate by 10%. oldBitrate = (int) (bitrate * decreaseRange); diff --git a/library/src/main/java/com/pedro/library/util/QueueAwareBitrateAdapter.kt b/library/src/main/java/com/pedro/library/util/QueueAwareBitrateAdapter.kt new file mode 100644 index 000000000..6d03c365d --- /dev/null +++ b/library/src/main/java/com/pedro/library/util/QueueAwareBitrateAdapter.kt @@ -0,0 +1,129 @@ +/* + * Copyright (C) 2026 pedroSG94. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.pedro.library.util + +import com.pedro.common.StreamingStatsReport +import com.pedro.common.Throughput + +/** + * Alternative to [BitrateAdapter] driven by the send queue instead of by the measured bitrate. + * + * [BitrateAdapter] probes upwards blindly and only reduces when the measured bitrate falls far + * enough below the configured one, so a link slightly slower than the target makes it oscillate + * above and below the real capacity. This one reads the queue: when frames start piling up the + * link is the limit, so it records the bitrate the link actually delivered and stays below it + * instead of climbing back to the maximum. + * + * Feed it from onStreamingStats and apply [Listener.onBitrateAdapted] to the video encoder. + * [maxBitrate] is the whole wire budget, video plus audio, because that is what the report + * measures, so subtract the audio bitrate before applying it to video. + */ + +class QueueAwareBitrateAdapter( + private val maxBitrate: Int, + minBitrate: Int, + /** + * Seconds to check if the network improved doing set bitrate to max and decrease again if the network still is bad. + * Use 0 or less if you never want check if the network improved. 300s by default + */ + private val ceilingTtl: Int, + private val listener: Listener +) { + + constructor(maxBitrate: Int, listener: Listener): + this(maxBitrate, maxBitrate / 10, DEFAULT_CEILING_TTL, listener) + + constructor(maxBitrate: Int, minBitrate: Int, listener: Listener): + this(maxBitrate, minBitrate, DEFAULT_CEILING_TTL, listener) + + fun interface Listener { + fun onBitrateAdapted(bitrate: Int) + } + + companion object { + const val DEFAULT_CEILING_TTL = 300 + private const val QUEUE_ALERT_FRACTION = 0.15f //seconds of video allowed in queue + private const val BACKOFF = 0.90f + private const val PROBE = 1.05f + private const val HOLD_SECONDS = 4 + private const val CEILING_MARGIN = 0.97f + //a transient must not define the link, so the capacity is averaged over this window + private const val CAPACITY_WINDOW = 15 + //and even a real drop cannot halve the ceiling twice in a row + private const val MAX_CEILING_DROP = 0.5f + //each probe closes part of the distance to the ceiling, so recovering from a deep drop + //does not take minutes + private const val GAP_CLOSE = 0.25f + } + + private val floor = minBitrate.coerceIn(1, maxBitrate) + private val recent = ArrayDeque() + private var hasMeasured = false + private var target = maxBitrate + private var ceiling = maxBitrate + private var good = 0 + private var age = 0 + + fun onStreamingStats(report: StreamingStatsReport) { + if (report.smoothedBitrate > 0) hasMeasured = true + if (report.queueBytesOut > 0 && hasMeasured) { + recent.addLast(report.smoothedBitrate) + if (recent.size > CAPACITY_WINDOW) recent.removeFirst() + } + val alertBytes = (target / 8) * QUEUE_ALERT_FRACTION + val congested = report.throughput == Throughput.INSUFFICIENT || report.queueBytesOut > alertBytes + if (congested) { + val measured = if (recent.isNotEmpty()) recent.average().toLong() + else report.smoothedBitrate.takeIf { it > 0 } + if (measured != null) { + val dropped = minOf(ceiling.toLong(), measured).toInt() + ceiling = maxOf(dropped, (ceiling * MAX_CEILING_DROP).toInt()) + target = (ceiling * BACKOFF).toInt().coerceAtLeast(floor) + good = 0 + age = 0 + listener.onBitrateAdapted(target) + } + } else { + good++ + if (good >= HOLD_SECONDS) { + good = 0 + val cap = if (ceiling >= maxBitrate) maxBitrate + else (ceiling * CEILING_MARGIN).toInt().coerceAtLeast(floor) + val gapStep = target + ((cap - target) * GAP_CLOSE).toInt() + val next = minOf(maxOf(gapStep, (target * PROBE).toInt()), cap) + if (next != target) { + target = next + listener.onBitrateAdapted(target) + } + } + } + if (ceilingTtl > 0 && ++age > ceilingTtl) { + ceiling = maxBitrate + age = 0 + } + } + + fun reset() { + target = maxBitrate + ceiling = maxBitrate + good = 0 + age = 0 + recent.clear() + hasMeasured = false + listener.onBitrateAdapted(target) + } +} diff --git a/library/src/test/java/com/pedro/library/util/BitrateAdapterTest.kt b/library/src/test/java/com/pedro/library/util/BitrateAdapterTest.kt new file mode 100644 index 000000000..8219be222 --- /dev/null +++ b/library/src/test/java/com/pedro/library/util/BitrateAdapterTest.kt @@ -0,0 +1,89 @@ +/* + * Copyright (C) 2024 pedroSG94. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.pedro.library.util + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * maxBitrate is 3200000 of video plus 64000 of audio, the configuration reported in issue #2177. + * adaptBitrate only produces a value every 5 samples. + */ +class BitrateAdapterTest { + + private val maxBitrate = 3264000 + + private fun adapt(samples: List): List { + val results = mutableListOf() + val adapter = BitrateAdapter { results.add(it) } + adapter.setMaxBitrate(maxBitrate) + samples.forEach { adapter.adaptBitrate(it) } + return results + } + + private fun adapt(samples: List, hasCongestion: Boolean): List { + val results = mutableListOf() + val adapter = BitrateAdapter { results.add(it) } + adapter.setMaxBitrate(maxBitrate) + samples.forEach { adapter.adaptBitrate(it, hasCongestion) } + return results + } + + @Test + fun `GIVEN a link faster than max WHEN adapt THEN keep max bitrate`() { + val results = adapt(List(5) { 3400000L }) + assertEquals(listOf(3264000), results) + } + + @Test + fun `GIVEN a link a bit slower than max WHEN adapt THEN go below the link instead of pinning at max`() { + //3150000 is fast enough to stay above oldBitrate * 0.9, so the decrease branch was never + //taken and the bitrate stayed pinned at maxBitrate over a link that cannot carry it + val results = adapt(List(5) { 3150000L }) + assertEquals(listOf(2441249), results) + assertTrue(results.first() < 3150000) + } + + @Test + fun `GIVEN a link a bit slower than max WHEN adapt many times THEN never settle above the link`() { + val link = 3150000L + val results = mutableListOf() + val adapter = BitrateAdapter { results.add(it) } + adapter.setMaxBitrate(maxBitrate) + var configured = maxBitrate + repeat(10 * 5) { + //the sender can only push what the link carries + adapter.adaptBitrate(minOf(configured.toLong(), link)) + configured = results.lastOrNull() ?: configured + } + assertEquals(10, results.size) + assertTrue("settled above the link: $results", results.count { it > link } < results.size) + } + + @Test + fun `GIVEN congestion WHEN measured bitrate reaches max THEN reduce anyway`() { + val results = adapt(List(5) { 4000000L }, hasCongestion = true) + assertEquals(listOf(3100000), results) + } + + @Test + fun `GIVEN no congestion WHEN measured bitrate reaches max THEN keep max bitrate`() { + val results = adapt(List(5) { 4000000L }, hasCongestion = false) + assertEquals(listOf(3264000), results) + } +} diff --git a/library/src/test/java/com/pedro/library/util/QueueAwareBitrateAdapterTest.kt b/library/src/test/java/com/pedro/library/util/QueueAwareBitrateAdapterTest.kt new file mode 100644 index 000000000..bfa821c1d --- /dev/null +++ b/library/src/test/java/com/pedro/library/util/QueueAwareBitrateAdapterTest.kt @@ -0,0 +1,221 @@ +/* + * Copyright (C) 2024 pedroSG94. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.pedro.library.util + +import com.pedro.common.StreamingStatsReport +import com.pedro.common.Throughput +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * maxBitrate is 3200000 of video plus 64000 of audio, the configuration reported in issue #2177. + */ +class QueueAwareBitrateAdapterTest { + + private val maxBitrate = 3264000 + + private fun report(smoothedBitrate: Long, queueBytesOut: Long, throughput: Throughput) = + StreamingStatsReport( + bytesOutPerSecond = 0, + queueBytesOut = queueBytesOut, + totalBytesOut = 0, + queueCongestionPercent = 0f, + throughput = throughput, + bitrate = 0, + smoothedBitrate = smoothedBitrate, + ) + + private fun healthy(bitrate: Int) = report(bitrate.toLong(), 0, Throughput.SUFFICIENT) + + @Test + fun `GIVEN a healthy link WHEN adapt for a long time THEN keep the whole configured bitrate`() { + //the encoder starts at maxBitrate, an adapter that never touches it is the correct result + var configured = maxBitrate + val adapter = QueueAwareBitrateAdapter(maxBitrate) { configured = it } + repeat(300) { adapter.onStreamingStats(healthy(maxBitrate)) } + assertEquals(maxBitrate, configured) + } + + @Test + fun `GIVEN the queue piling up WHEN adapt THEN drop below the bitrate the link delivered`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + adapter.onStreamingStats(report(3000000, 500000, Throughput.INSUFFICIENT)) + assertEquals(2700000, last) + } + + @Test + fun `GIVEN no measurement yet WHEN the queue piles up THEN do not report a zero bitrate`() { + val emitted = mutableListOf() + val adapter = QueueAwareBitrateAdapter(maxBitrate) { emitted.add(it) } + //BitrateManager reports 0 until its first one second window closes + adapter.onStreamingStats(report(0, 900000, Throughput.INSUFFICIENT)) + assertEquals(emptyList(), emitted) + } + + @Test + fun `GIVEN a single terrible second WHEN it is the only measurement THEN do not collapse`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //a second delivering 50 kbps used to define the link and drop the bitrate to the floor + adapter.onStreamingStats(report(50000, 900000, Throughput.INSUFFICIENT)) + assertEquals(1468800, last) + assertTrue("collapsed to $last", last > maxBitrate / 4) + } + + @Test + fun `GIVEN backlogged seconds WHEN one goes bad THEN average them instead of trusting the worst`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //seconds with frames waiting are real link measurements, they build the capacity window + repeat(10) { adapter.onStreamingStats(report(3000000, 500000, Throughput.INSUFFICIENT)) } + adapter.onStreamingStats(report(50000, 900000, Throughput.INSUFFICIENT)) + //the average absorbs the transient instead of letting it define the link + assertTrue("collapsed to $last", last > 2000000) + } + + @Test + fun `GIVEN a throttled session WHEN reset THEN go back to max and tell the encoder`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + adapter.onStreamingStats(report(1000000, 900000, Throughput.INSUFFICIENT)) + assertTrue(last < maxBitrate) + + adapter.reset() + assertEquals(maxBitrate, last) + } + + @Test + fun `GIVEN a link slower than max WHEN adapt in a closed loop THEN stay under it almost always`() { + val link = 3000000 + var configured = maxBitrate + val adapter = QueueAwareBitrateAdapter(maxBitrate) { configured = it } + var secondsOverTheLink = 0 + repeat(600) { + val congested = configured > link + if (congested) secondsOverTheLink++ + //the link only carries what it carries, and the queue grows while we push more + adapter.onStreamingStats( + if (congested) report(link.toLong(), 500000, Throughput.INSUFFICIENT) + else healthy(configured) + ) + } + //it re-tests the link once every CEILING_TTL, so it goes over briefly by design + assertTrue("over the link $secondsOverTheLink seconds of 600", secondsOverTheLink < 60) + } + + @Test + fun `GIVEN a bursty link WHEN the queue drains in a burst THEN do not read the burst as capacity`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //a shaped link stalls and then drains the backlog at twice the rate. Taking the best second + //would read 6 Mbps as capacity on a link that only carries 2.5 + val pattern = listOf(2500000L, 0L, 6000000L, 2500000L, 1000000L, 6000000L, 2500000L) + repeat(3) { + pattern.forEach { adapter.onStreamingStats(report(it, 500000, Throughput.INSUFFICIENT)) } + } + assertTrue("aimed at $last, above what the link carries", last < 2500000) + } + + @Test + fun `GIVEN an empty queue WHEN adapt THEN do not use it as a capacity measurement`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //ten healthy seconds at max, the link was never the limit so they say nothing about capacity + repeat(10) { adapter.onStreamingStats(healthy(maxBitrate)) } + //now the link backs up and only delivers 1 Mbps. Had the healthy seconds been used as + //measurements the window would average near maxBitrate and barely reduce anything + adapter.onStreamingStats(report(1000000, 500000, Throughput.INSUFFICIENT)) + assertEquals(1468800, last) + } + + @Test + fun `GIVEN a link that collapses WHEN it probes back up THEN never go under the floor`() { + val emitted = mutableListOf() + val adapter = QueueAwareBitrateAdapter(maxBitrate) { emitted.add(it) } + //a link that keeps failing drives the ceiling down; the probe branch used to cap the + //target at the collapsed ceiling, ignoring the floor and reaching zero + repeat(40) { + adapter.onStreamingStats(report(20000, 900000, Throughput.INSUFFICIENT)) + repeat(5) { adapter.onStreamingStats(healthy(20000)) } + } + val lowest = emitted.min() + assertTrue("emitted $lowest, under the floor", lowest >= maxBitrate / 10) + } + + @Test + fun `GIVEN a link that stalls WHEN frames are waiting THEN count the stall as a measurement`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //one real value first, so the adapter knows BitrateManager is producing measurements + adapter.onStreamingStats(report(2000000, 500000, Throughput.INSUFFICIENT)) + val afterSlowLink = last + //now the link stops delivering entirely while frames pile up + repeat(10) { adapter.onStreamingStats(report(0, 900000, Throughput.INSUFFICIENT)) } + assertTrue("stalls did not lower the estimate: $afterSlowLink -> $last", last < afterSlowLink) + } + + @Test + fun `GIVEN no measurement yet WHEN the queue piles up THEN ignore the zero as a capacity value`() { + var last = 0 + val adapter = QueueAwareBitrateAdapter(maxBitrate) { last = it } + //BitrateManager reports 0 until its first window closes; with a queue already growing that + //zero must not be read as "the link delivers nothing" + adapter.onStreamingStats(report(0, 900000, Throughput.INSUFFICIENT)) + assertEquals(0, last) + } + + @Test + fun `GIVEN a ttl of zero WHEN the link stays the same THEN never climb over it again`() { + val link = 2800000 + var configured = maxBitrate + val adapter = QueueAwareBitrateAdapter(maxBitrate, maxBitrate / 10, 0) { configured = it } + var secondsOverTheLink = 0 + repeat(1800) { + val over = configured > link + if (over) secondsOverTheLink++ + adapter.onStreamingStats( + if (over) report(link.toLong(), 500000, Throughput.INSUFFICIENT) + else healthy(configured) + ) + } + //without re-testing it settles under the link and stays there + assertTrue("over the link $secondsOverTheLink seconds of 1800", secondsOverTheLink < 30) + assertTrue("settled at $configured, over the link", configured <= link) + } + + @Test + fun `GIVEN the default ttl WHEN the link stays the same THEN re-test it now and then`() { + val link = 2800000 + var configured = maxBitrate + val adapter = QueueAwareBitrateAdapter(maxBitrate) { configured = it } + var overshoots = 0 + var wasOver = false + repeat(1800) { + val over = configured > link + if (over && !wasOver) overshoots++ + wasOver = over + adapter.onStreamingStats( + if (over) report(link.toLong(), 500000, Throughput.INSUFFICIENT) + else healthy(configured) + ) + } + //1800 seconds at the default ttl of 300 leaves a handful of re-tests, not one per minute + assertTrue("re-tested $overshoots times in 1800s", overshoots in 1..8) + } +}