Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 5 additions & 6 deletions app/src/main/java/com/pedro/streamer/rotation/CameraFragment.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -203,6 +201,7 @@ class CameraFragment: Fragment(), ConnectChecker {
}

override fun onConnectionStarted(url: String) {
bitrateAdapter.reset()
}

override fun onConnectionSuccess() {
Expand All @@ -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(),
Expand Down
1 change: 1 addition & 0 deletions common/src/main/java/com/pedro/common/BitrateManager.kt
Original file line number Diff line number Diff line change
Expand Up @@ -48,5 +48,6 @@ open class BitrateManager(private val bitrateChecker: BitrateChecker) {
fun reset() {
bitrate = 0
bitrateOld = 0
timeStamp = TimeUtils.getCurrentTimeMillis()
}
}
13 changes: 13 additions & 0 deletions common/src/main/java/com/pedro/common/Extensions.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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<Any?> {
override val context: CoroutineContext
get() = dispatcher

override fun resumeWith(result: Result<Any?>) {
result.exceptionOrNull()?.let { Log.e("getSuspendContext", "Error", it) }
}
}
33 changes: 33 additions & 0 deletions common/src/test/java/com/pedro/common/BitrateManagerTest.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<Long>()
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<Long>()
verify(connectChecker, times(2)).onNewBitrate(resultValue.capture())
assertEquals(2_000_000L, resultValue.secondValue)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -364,8 +364,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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,7 @@ public boolean prepareVideo(int width, int height, int fps, int bitrate, int rot
glStreamInterface.setIsPortrait(isPortrait);
}
}
forceFpsLimit(true);
return videoInitialized;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down
11 changes: 9 additions & 2 deletions library/src/main/java/com/pedro/library/util/BitrateAdapter.java
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
*/
public class BitrateAdapter {

private static final float MAX_BITRATE_TOLERANCE = 0.95f;

public interface Listener {
void onBitrateAdapted(int bitrate);
}
Expand Down Expand Up @@ -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);
Expand All @@ -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);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
/*
* 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,
private val listener: Listener
) {

constructor(maxBitrate: Int, listener: Listener): this(maxBitrate, maxBitrate / 10, listener)

fun interface Listener {
fun onBitrateAdapted(bitrate: Int)
}

private companion object {
const val QUEUE_ALERT_FRACTION = 0.15f //seconds of video allowed in queue
const val BACKOFF = 0.90f
const val PROBE = 1.05f
const val HOLD_SECONDS = 4
const val CEILING_MARGIN = 0.97f
const val CEILING_TTL = 60
}

private val floor = minBitrate.coerceIn(1, maxBitrate)
private var target = maxBitrate
private var ceiling = maxBitrate
private var good = 0
private var age = 0

fun onStreamingStats(report: StreamingStatsReport) {
val alertBytes = (target / 8) * QUEUE_ALERT_FRACTION
val congested = report.throughput == Throughput.INSUFFICIENT || report.queueBytesOut > alertBytes
if (congested) {
//BitrateManager reports 0 until its first window closes, that is not a measurement
if (report.smoothedBitrate > 0) {
ceiling = minOf(ceiling.toLong(), report.smoothedBitrate).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()
val next = minOf((target * PROBE).toInt(), cap)
if (next != target) {
target = next
listener.onBitrateAdapted(target)
}
}
}
if (++age > CEILING_TTL) {
ceiling = maxBitrate
age = 0
}
}

fun reset() {
target = maxBitrate
ceiling = maxBitrate
good = 0
age = 0
listener.onBitrateAdapted(target)
}
}
89 changes: 89 additions & 0 deletions library/src/test/java/com/pedro/library/util/BitrateAdapterTest.kt
Original file line number Diff line number Diff line change
@@ -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<Long>): List<Int> {
val results = mutableListOf<Int>()
val adapter = BitrateAdapter { results.add(it) }
adapter.setMaxBitrate(maxBitrate)
samples.forEach { adapter.adaptBitrate(it) }
return results
}

private fun adapt(samples: List<Long>, hasCongestion: Boolean): List<Int> {
val results = mutableListOf<Int>()
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<Int>()
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)
}
}
Loading
Loading