Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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)
}

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

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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,9 +66,9 @@ class Sdt(
return nBits
}

fun write() {
fun write(timestamp: Long) {
if (services.isNotEmpty()) {
write(toByteBuffer())
write(toByteBuffer(), timestamp)
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -299,7 +299,6 @@ private class DefaultSurfaceProcessor(
if (isReleaseRequested.get()) {
return
}
surfaceTexture
val timeConverter = surfaceInputsToTimeConverterMap[surfaceTexture] ?: return

surfaceTexture.updateTexImage()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Frame>(Channel.UNLIMITED, onUndeliveredElement = {
Channel<Frame>(AUDIO_ENCODER_CHANNEL_SIZE, onUndeliveredElement = {
it.close()
Logger.w(TAG, "Audio frame dropped")
})
}

Expand All @@ -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<Frame>(Channel.UNLIMITED, onUndeliveredElement = {
Channel<Frame>(VIDEO_ENCODER_CHANNEL_SIZE, onUndeliveredElement = {
it.close()
Logger.w(TAG, "Video frame dropped")
})
}

Expand Down Expand Up @@ -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.
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,8 @@ class RtmpEndpoint internal constructor(
private var videoPayloadSendDroppedSize = 0L

private val flvTagChannel = ChannelWithCloseableData<FLVTag>(
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)
Expand Down Expand Up @@ -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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -141,18 +141,12 @@ 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) {
completionException?.let {
// Socket already closed
throw ClosedException(completionException!!)
throw ClosedException(it)
}
close()
throw ClosedException(t)
Expand Down
2 changes: 1 addition & 1 deletion gradle/libs.versions.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
Loading