From c0d3d759001556ca7930741536fd571294487804 Mon Sep 17 00:00:00 2001 From: ThibaultBee <37510686+ThibaultBee@users.noreply.github.com> Date: Fri, 25 Sep 2026 23:41:21 +0200 Subject: [PATCH 1/4] =?UTF-8?q?fix(core):=20TS:=20add=20timestamp=20parame?= =?UTF-8?q?ter=20to=20write=20methods=20for=20PAT,=20PMT,=20and=20SDT=20It?= =?UTF-8?q?=20avoids=20relying=20on=20SRT=E2=80=99s=20internal=20timing=20?= =?UTF-8?q?mechanisms,=20which=20could=20cause=20ffplay=20to=20appear=20to?= =?UTF-8?q?=20drop=20frames.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../endpoints/composites/muxers/ts/TsMuxer.kt | 39 ++++++++++--------- .../composites/muxers/ts/packets/Pat.kt | 4 +- .../composites/muxers/ts/packets/Pmt.kt | 4 +- .../composites/muxers/ts/packets/Psi.kt | 4 +- .../composites/muxers/ts/packets/Sdt.kt | 4 +- 5 files changed, 29 insertions(+), 26 deletions(-) diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/TsMuxer.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/TsMuxer.kt index 4706786d7..3025608d9 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/TsMuxer.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/TsMuxer.kt @@ -182,16 +182,19 @@ class TsMuxer : IMuxerInternal { private fun generateStreams( frame: Frame, pes: Pes ) { - retransmitPsi(pes.stream.isVideo and frame.isKeyFrame) + retransmitPsi(pes.stream.isVideo and frame.isKeyFrame, frame.ptsInUs) pes.write(frame) } /** - * Manages table retransmission + * Manages table retransmission. + * If [sendPsiOnce] is true, tables (SDT, PAT, PMT) are sent only once at the beginning of the stream. + * Otherwise, they are retransmitted periodically and on every video key frame. * * @param forcePat Force to remit a PAT. Set to true on video key frame. + * @param timestamp Timestamp in µs to associate with the PSI packets (typically the current frame PTS). */ - private fun retransmitPsi(forcePat: Boolean) { + private fun retransmitPsi(forcePat: Boolean, timestamp: Long = 0L) { var sendSdt = false var sendPat = false @@ -208,41 +211,41 @@ class TsMuxer : IMuxerInternal { } if (sendSdt) { - sendSdt() + sendSdt(timestamp) } if (sendPat) { - sendPat() - sendPmts() + sendPat(timestamp) + sendPmts(timestamp) } } - private fun upgradePat() { + private fun upgradePat(timestamp: Long = 0L) { pat.versionNumber = (pat.versionNumber + 1).toByte() - sendPat() + sendPat(timestamp) } - private fun sendPat() { - pat.write() + private fun sendPat(timestamp: Long = 0L) { + pat.write(timestamp) } - private fun sendPmt(service: Service) { - service.pmt?.write() ?: throw UnsupportedOperationException("PMT must not be null") + private fun sendPmt(service: Service, timestamp: Long = 0L) { + service.pmt?.write(timestamp) ?: throw UnsupportedOperationException("PMT must not be null") } - private fun sendPmts() { + private fun sendPmts(timestamp: Long = 0L) { tsServices.filter { it.pmt != null }.forEach { - it.pmt?.write() ?: throw UnsupportedOperationException("PMT must not be null") + it.pmt?.write(timestamp) ?: throw UnsupportedOperationException("PMT must not be null") } } - private fun upgradeSdt() { + private fun upgradeSdt(timestamp: Long = 0L) { sdt.versionNumber = (sdt.versionNumber + 1).toByte() - sendSdt() + sendSdt(timestamp) } - private fun sendSdt() { - sdt.write() + private fun sendSdt(timestamp: Long = 0L) { + sdt.write(timestamp) } /** diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pat.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pat.kt index 034093978..9abcdf3ec 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pat.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pat.kt @@ -54,9 +54,9 @@ class Pat( override val size: Int get() = bitSize / Byte.SIZE_BITS - fun write() { + fun write(timestamp: Long) { if (services.any { it.pmt != null }) { - write(toByteBuffer()) + write(toByteBuffer(), timestamp) } } diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pmt.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pmt.kt index 5365c987a..69eaef5c7 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pmt.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Pmt.kt @@ -62,9 +62,9 @@ class Pmt( MediaFormat.MIMETYPE_VIDEO_HEVC ) * streams.filter { it.config.mimeType == MediaFormat.MIMETYPE_VIDEO_HEVC }.size - fun write() { + fun write(timestamp: Long) { if (service.pcrPid != null) { - write(toByteBuffer()) + write(toByteBuffer(), timestamp) } } diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Psi.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Psi.kt index e6e537fe6..c43e10454 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Psi.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Psi.kt @@ -40,8 +40,8 @@ open class Psi( const val PSI_HEADER_SIZE = 9 // contains pointer_field } - protected fun write(buffer: ByteBuffer) { - write(payload = toByteBuffer(buffer)) + protected fun write(buffer: ByteBuffer, timestamp: Long) { + write(payload = toByteBuffer(buffer), timestamp = timestamp) } fun toByteBuffer(payload: ByteBuffer): ByteBuffer { diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Sdt.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Sdt.kt index d22a0a304..de32efc9d 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Sdt.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/endpoints/composites/muxers/ts/packets/Sdt.kt @@ -66,9 +66,9 @@ class Sdt( return nBits } - fun write() { + fun write(timestamp: Long) { if (services.isNotEmpty()) { - write(toByteBuffer()) + write(toByteBuffer(), timestamp) } } From 78b6bfaac3a7368c5c51381428aff4adc562057d Mon Sep 17 00:00:00 2001 From: ThibaultBee <37510686+ThibaultBee@users.noreply.github.com> Date: Mon, 14 Sep 2026 21:59:58 +0200 Subject: [PATCH 2/4] fix(*): avoid to stall the encoder on network pressure --- .../outputs/encoding/EncodingPipelineOutput.kt | 17 +++++++++++++++-- .../ext/rtmp/elements/endpoints/RtmpEndpoint.kt | 6 +++++- .../endpoints/composites/sinks/SrtSink.kt | 8 +------- gradle/libs.versions.toml | 2 +- 4 files changed, 22 insertions(+), 11 deletions(-) diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/pipelines/outputs/encoding/EncodingPipelineOutput.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/pipelines/outputs/encoding/EncodingPipelineOutput.kt index e01806fb6..5d48d6509 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/pipelines/outputs/encoding/EncodingPipelineOutput.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/pipelines/outputs/encoding/EncodingPipelineOutput.kt @@ -265,9 +265,13 @@ internal class EncodingPipelineOutput( onInternalError(t) } + /** + * Using a size will make the [Channel] release encoder frame instead of stalling it. + */ override val outputChannel = - Channel(Channel.UNLIMITED, onUndeliveredElement = { + Channel(AUDIO_ENCODER_CHANNEL_SIZE, onUndeliveredElement = { it.close() + Logger.w(TAG, "Audio frame dropped") }) } @@ -276,9 +280,13 @@ internal class EncodingPipelineOutput( onInternalError(t) } + /** + * Using a size will make the [Channel] release encoder frame instead of stalling it. + */ override val outputChannel = - Channel(Channel.UNLIMITED, onUndeliveredElement = { + Channel(VIDEO_ENCODER_CHANNEL_SIZE, onUndeliveredElement = { it.close() + Logger.w(TAG, "Video frame dropped") }) } @@ -965,5 +973,10 @@ internal class EncodingPipelineOutput( companion object { private const val TAG = "EncodingPipelineOutput" + + private const val AUDIO_ENCODER_CHANNEL_SIZE = + 4 // This is an arbitrary value. It depends on the encoder, but 4 is a reasonable value based on real-world implementations. + private const val VIDEO_ENCODER_CHANNEL_SIZE = + 4 // This is an arbitrary value. It depends on the encoder, but 4 is a reasonable value based on real-world implementations. } } diff --git a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt index 0eecfebdd..302e9c53d 100644 --- a/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt +++ b/extensions/rtmp/src/main/java/io/github/thibaultbee/streampack/ext/rtmp/elements/endpoints/RtmpEndpoint.kt @@ -81,7 +81,8 @@ class RtmpEndpoint internal constructor( private var videoPayloadSendDroppedSize = 0L private val flvTagChannel = ChannelWithCloseableData( - 10 /* Arbitrary buffer size. TODO: add a parameter to set it */, BufferOverflow.DROP_OLDEST, + FLV_TAG_CHANNEL_SIZE, + BufferOverflow.DROP_OLDEST, onUndeliveredElement = { flvTag -> synchronized(metricsLock) { val payloadSize = flvTag.data.getSize(AmfVersion.AMF0) @@ -304,6 +305,9 @@ class RtmpEndpoint internal constructor( private const val INVALID_TIMESTAMP = -1L + private const val FLV_TAG_CHANNEL_SIZE = + 8 // Arbitrary buffer size. It's Audio output size + Video output size in the encoding output + init { System.setProperty("kotlinx.io.pool.size.bytes", "4194304") // 4MB } diff --git a/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt b/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt index b77f3acea..f3624a4fb 100644 --- a/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt +++ b/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt @@ -141,13 +141,7 @@ class SrtSink(private val coroutineDispatcher: CoroutineDispatcher) : AbstractSi ) return -1 } - } catch (t: Throwable) { - Logger.w(TAG, "Failed to get connection time: $t") - return -1 - } - - try { - return socket.send(packet.buffer, buildMsgCtrl(packet)) + return socket.trySend(packet.buffer, buildMsgCtrl(packet)) } catch (t: Throwable) { isOnError = true if (completionException != null) { diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index abbd7c988..f33814572 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -27,7 +27,7 @@ material = "1.14.0" mockk = "1.14.11" robolectric = "4.16.1" komuxer = "0.4.0" -srtdroid = "1.9.5" +srtdroid = "1.10.1" junitKtx = "1.3.0" compose = "1.12.0" From 4573053a8be63fbf4b7989d9f326cfc094487673 Mon Sep 17 00:00:00 2001 From: ThibaultBee <37510686+ThibaultBee@users.noreply.github.com> Date: Fri, 25 Sep 2026 20:38:53 +0200 Subject: [PATCH 3/4] chore(processor): remove unused surfaceTexture reference to prevent potential issues --- .../core/elements/processing/video/DefaultSurfaceProcessor.kt | 1 - 1 file changed, 1 deletion(-) diff --git a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/processing/video/DefaultSurfaceProcessor.kt b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/processing/video/DefaultSurfaceProcessor.kt index 399949775..a9d6bc8e8 100644 --- a/core/src/main/java/io/github/thibaultbee/streampack/core/elements/processing/video/DefaultSurfaceProcessor.kt +++ b/core/src/main/java/io/github/thibaultbee/streampack/core/elements/processing/video/DefaultSurfaceProcessor.kt @@ -299,7 +299,6 @@ private class DefaultSurfaceProcessor( if (isReleaseRequested.get()) { return } - surfaceTexture val timeConverter = surfaceInputsToTimeConverterMap[surfaceTexture] ?: return surfaceTexture.updateTexImage() From 56f76d7d9fa6ccc8902275cc070353fbb06671c4 Mon Sep 17 00:00:00 2001 From: ThibaultBee <37510686+ThibaultBee@users.noreply.github.com> Date: Sat, 26 Sep 2026 21:56:06 +0200 Subject: [PATCH 4/4] fix(srt): use trySend to avoid blocking the encoder --- .../ext/srt/elements/endpoints/composites/sinks/SrtSink.kt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt b/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt index f3624a4fb..aafc54ba9 100644 --- a/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt +++ b/extensions/srt/src/main/java/io/github/thibaultbee/streampack/ext/srt/elements/endpoints/composites/sinks/SrtSink.kt @@ -144,9 +144,9 @@ class SrtSink(private val coroutineDispatcher: CoroutineDispatcher) : AbstractSi return socket.trySend(packet.buffer, buildMsgCtrl(packet)) } catch (t: Throwable) { isOnError = true - if (completionException != null) { + completionException?.let { // Socket already closed - throw ClosedException(completionException!!) + throw ClosedException(it) } close() throw ClosedException(t)