Skip to content
Merged
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 @@ -113,8 +113,7 @@ class Mpeg2TsMuxerRecordController : AsyncBaseRecordController() {
}

private fun setTrackConfig(videoEnabled: Boolean, audioEnabled: Boolean) {
Pid.reset()
service.clearTracks()
service.clear()
if (audioEnabled) service.addTrack(getAudioCodec().toCodec())
if (videoEnabled) service.addTrack(getVideoCodec().toCodec())
service.generatePmt()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ class SrtStreamClient(
): StreamBaseClient() {

/**
* Set latency in micro. By default 120_000.
* Set latency in millis. By default 120.
*/
fun setLatency(latency: Int) {
srtClient.setLatency(latency)
Expand Down
2 changes: 1 addition & 1 deletion rtmp/src/main/java/com/pedro/rtmp/amf/v0/AmfData.kt
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ abstract class AmfData {
}

fun getMarkType(type: Int): AmfType {
return AmfType.entries.find { it.mark.toInt() == type } ?: AmfType.STRING
return AmfType.entries.find { it.mark.toInt() == type } ?: throw IOException("Unimplemented AMF data type: $type")
}
}

Expand Down
2 changes: 1 addition & 1 deletion rtmp/src/main/java/com/pedro/rtmp/amf/v0/AmfDate.kt
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,6 @@ class AmfDate(var date: Double = TimeUtils.getCurrentTimeMillis().toDouble()): A
override fun getSize(): Int = 10

override fun toString(): String {
return "AmfUnsupported"
return "AmfDate value: $date"
}
}
10 changes: 10 additions & 0 deletions rtmp/src/main/java/com/pedro/rtmp/amf/v0/AmfEcmaArray.kt
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,16 @@ class AmfEcmaArray(private val properties: LinkedHashMap<AmfString, AmfData> = L
length = properties.size
}

override fun setProperty(name: String, data: AmfData) {
super.setProperty(name, data)
length = properties.size
}

override fun setProperty(name: String, data: Any) {
super.setProperty(name, data)
length = properties.size
}

@Throws(IOException::class)
override fun readBody(input: InputStream) {
//get number of items as UInt32
Expand Down
8 changes: 4 additions & 4 deletions rtmp/src/main/java/com/pedro/rtmp/amf/v0/AmfObject.kt
Original file line number Diff line number Diff line change
Expand Up @@ -121,20 +121,20 @@ open class AmfObject(private val properties: LinkedHashMap<AmfString, AmfData> =
properties.clear()
bodySize = 0
val objectEnd = AmfObjectEnd()
val markInputStream: InputStream = if (input.markSupported()) input else BufferedInputStream(input)
val markInputStream = if (input.markSupported()) input else BufferedInputStream(input)
while (!objectEnd.found) {
markInputStream.mark(objectEnd.getSize())
objectEnd.readBody(input)
objectEnd.readBody(markInputStream)
if (objectEnd.found) {
bodySize += objectEnd.getSize()
} else {
markInputStream.reset()

val key = AmfString()
key.readBody(input)
key.readBody(markInputStream)
bodySize += key.getSize()

val value = getAmfData(input)
val value = getAmfData(markInputStream)
bodySize += value.getSize() + 1

properties[key] = value
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,15 +25,9 @@ import com.pedro.common.AudioUtils
*/
class AacAudioSpecificConfig(private val type: Int, private val sampleRate: Int, private val channels: Int) {

val size = 9
val size = 2

fun write(buffer: ByteArray, offset: Int) {
writeConfig(buffer, offset)
val adts = AudioUtils.createAdtsHeader(type, buffer.size, sampleRate, channels)
adts.get(buffer, offset + 2, adts.capacity())
}

private fun writeConfig(buffer: ByteArray, offset: Int) {
val frequency = AudioUtils.getFrequency(sampleRate)
buffer[offset] = ((type shl 3) or (frequency shr 1)).toByte()
buffer[offset + 1] = (frequency shl 7 and 0x80).plus(channels shl 3 and 0x78).toByte()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,15 +37,15 @@ class OpusAudioSpecificConfig(private val sampleRate: Int, private val channels:
buffer[offset + 8] = 0x01 //version 1
buffer[offset + 9] = channels.toByte()
val preSkip = 3840 //this is the recommended value by the RFC
buffer[offset + 10] = (preSkip shr 8).toByte()
buffer[offset + 11] = preSkip.toByte()
buffer[offset + 12] = (sampleRate shr 24).toByte()
buffer[offset + 13] = (sampleRate shr 16).toByte()
buffer[offset + 14] = (sampleRate shr 8).toByte()
buffer[offset + 15] = sampleRate.toByte()
buffer[offset + 10] = preSkip.toByte()
buffer[offset + 11] = (preSkip shr 8).toByte()
buffer[offset + 12] = sampleRate.toByte()
buffer[offset + 13] = (sampleRate shr 8).toByte()
buffer[offset + 14] = (sampleRate shr 16).toByte()
buffer[offset + 15] = (sampleRate shr 24).toByte()
val outputGain = 0
buffer[offset + 16] = (outputGain shr 8).toByte()
buffer[offset + 17] = outputGain.toByte()
buffer[offset + 16] = outputGain.toByte()
buffer[offset + 17] = (outputGain shr 8).toByte()
val mappingFamily = 0
buffer[offset + 18] = mappingFamily.toByte()
}
Expand Down
37 changes: 18 additions & 19 deletions rtmp/src/main/java/com/pedro/rtmp/rtmp/CommandsManager.kt
Original file line number Diff line number Diff line change
Expand Up @@ -58,12 +58,12 @@ abstract class CommandsManager {
var password: String? = null
var onAuth = false
var startTs = 0L
var readChunkSize = RtmpConfig.DEFAULT_CHUNK_SIZE
val config = RtmpConfig()
var audioDisabled = false
var videoDisabled = false
var customAmfObject: Map<String, Any> = emptyMap()
private var bytesRead = 0
private var acknowledgementSequence = 0
private var lastAcknowledgementSequence = 0

protected var width = 640
protected var height = 480
Expand Down Expand Up @@ -97,12 +97,12 @@ abstract class CommandsManager {
@Throws(IOException::class)
suspend fun sendChunkSize(socket: RtmpSocket) {
writeSync.withLock {
if (RtmpConfig.writeChunkSize != RtmpConfig.DEFAULT_CHUNK_SIZE) {
val chunkSize = SetChunkSize(RtmpConfig.writeChunkSize)
if (config.writeChunkSize != RtmpConfig.DEFAULT_CHUNK_SIZE) {
val chunkSize = SetChunkSize(config.writeChunkSize)
chunkSize.header.timeStamp = getCurrentTimestamp()
chunkSize.header.messageStreamId = streamId
chunkSize.writeHeader(socket)
chunkSize.writeBody(socket)
chunkSize.writeBody(socket, config.writeChunkSize)
socket.flush()
Log.i(TAG, "send $chunkSize")
} else {
Expand All @@ -129,7 +129,7 @@ abstract class CommandsManager {

@Throws(IOException::class)
suspend fun readMessageResponse(socket: RtmpSocket): RtmpMessage {
val message = RtmpMessage.getRtmpMessage(socket, readChunkSize, sessionHistory)
val message = RtmpMessage.getRtmpMessage(socket, config.readChunkSize, sessionHistory)
sessionHistory.setReadHeader(message.header)
Log.i(TAG, "read $message")
bytesRead += message.header.getPacketLength()
Expand All @@ -155,9 +155,9 @@ abstract class CommandsManager {
@Throws(IOException::class)
suspend fun sendWindowAcknowledgementSize(socket: RtmpSocket) {
writeSync.withLock {
val windowAcknowledgementSize = WindowAcknowledgementSize(RtmpConfig.acknowledgementWindowSize, getCurrentTimestamp())
val windowAcknowledgementSize = WindowAcknowledgementSize(config.acknowledgementWindowSize, getCurrentTimestamp())
windowAcknowledgementSize.writeHeader(socket)
windowAcknowledgementSize.writeBody(socket)
windowAcknowledgementSize.writeBody(socket, config.writeChunkSize)
socket.flush()
}
}
Expand All @@ -166,7 +166,7 @@ abstract class CommandsManager {
writeSync.withLock {
val pong = UserControl(Type.PONG_REPLY, event)
pong.writeHeader(socket)
pong.writeBody(socket)
pong.writeBody(socket, config.writeChunkSize)
socket.flush()
Log.i(TAG, "send pong")
}
Expand All @@ -176,7 +176,7 @@ abstract class CommandsManager {
writeSync.withLock {
val ping = UserControl(Type.PING_REQUEST, Event(TimeUtils.getCurrentTimeSeconds()))
ping.writeHeader(socket)
ping.writeBody(socket)
ping.writeBody(socket, config.writeChunkSize)
socket.flush()
Log.i(TAG, "send ping")
}
Expand All @@ -192,12 +192,11 @@ abstract class CommandsManager {

suspend fun checkAndSendAcknowledgement(socket: RtmpSocket) {
writeSync.withLock {
if (bytesRead >= RtmpConfig.acknowledgementWindowSize) {
acknowledgementSequence += bytesRead
bytesRead -= RtmpConfig.acknowledgementWindowSize
val acknowledgement = Acknowledgement(acknowledgementSequence)
if (bytesRead - lastAcknowledgementSequence >= config.acknowledgementWindowSize) {
lastAcknowledgementSequence = bytesRead
val acknowledgement = Acknowledgement(bytesRead)
acknowledgement.writeHeader(socket)
acknowledgement.writeBody(socket)
acknowledgement.writeBody(socket, config.writeChunkSize)
socket.flush()
Log.i(TAG, "send $acknowledgement")
}
Expand All @@ -209,7 +208,7 @@ abstract class CommandsManager {
writeSync.withLock {
val video = Video(flvPacket, streamId)
video.writeHeader(socket)
video.writeBody(socket)
video.writeBody(socket, config.writeChunkSize)
socket.flush(true)
return video.header.getPacketLength() //get packet size with header included to calculate bps
}
Expand All @@ -220,7 +219,7 @@ abstract class CommandsManager {
writeSync.withLock {
val audio = Audio(flvPacket, streamId)
audio.writeHeader(socket)
audio.writeBody(socket)
audio.writeBody(socket, config.writeChunkSize)
socket.flush(true)
return audio.header.getPacketLength() //get packet size with header included to calculate bps
}
Expand All @@ -237,9 +236,9 @@ abstract class CommandsManager {
timestamp = 0
streamId = 0
commandId = 0
readChunkSize = RtmpConfig.DEFAULT_CHUNK_SIZE
config.readChunkSize = RtmpConfig.DEFAULT_CHUNK_SIZE
sessionHistory.reset()
acknowledgementSequence = 0
lastAcknowledgementSequence = 0
bytesRead = 0
}
}
14 changes: 7 additions & 7 deletions rtmp/src/main/java/com/pedro/rtmp/rtmp/CommandsManagerAmf0.kt
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ class CommandsManagerAmf0: CommandsManager() {
connect.addData(connectInfo)

connect.writeHeader(socket)
connect.writeBody(socket)
connect.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "connect")
Log.i(TAG, "send $connect")
}
Expand All @@ -75,7 +75,7 @@ class CommandsManagerAmf0: CommandsManager() {
releaseStream.addData(AmfString(streamName))

releaseStream.writeHeader(socket)
releaseStream.writeBody(socket)
releaseStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "releaseStream")
Log.i(TAG, "send $releaseStream")

Expand All @@ -85,7 +85,7 @@ class CommandsManagerAmf0: CommandsManager() {
fcPublish.addData(AmfString(streamName))

fcPublish.writeHeader(socket)
fcPublish.writeBody(socket)
fcPublish.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "FCPublish")
Log.i(TAG, "send $fcPublish")

Expand All @@ -94,7 +94,7 @@ class CommandsManagerAmf0: CommandsManager() {
createStream.addData(AmfNull())

createStream.writeHeader(socket)
createStream.writeBody(socket)
createStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "createStream")
Log.i(TAG, "send $createStream")
}
Expand Down Expand Up @@ -134,7 +134,7 @@ class CommandsManagerAmf0: CommandsManager() {
metadata.addData(amfEcmaArray)

metadata.writeHeader(socket)
metadata.writeBody(socket)
metadata.writeBody(socket, config.writeChunkSize)
Log.i(TAG, "send $metadata")
}

Expand All @@ -147,7 +147,7 @@ class CommandsManagerAmf0: CommandsManager() {
publish.addData(AmfString("live"))

publish.writeHeader(socket)
publish.writeBody(socket)
publish.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, name)
Log.i(TAG, "send $publish")
}
Expand All @@ -158,7 +158,7 @@ class CommandsManagerAmf0: CommandsManager() {
closeStream.addData(AmfNull())

closeStream.writeHeader(socket)
closeStream.writeBody(socket)
closeStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, name)
Log.i(TAG, "send $closeStream")
}
Expand Down
14 changes: 7 additions & 7 deletions rtmp/src/main/java/com/pedro/rtmp/rtmp/CommandsManagerAmf3.kt
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ class CommandsManagerAmf3: CommandsManager() {
connect.addData(connectInfo)

connect.writeHeader(socket)
connect.writeBody(socket)
connect.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "connect")
Log.i(TAG, "send $connect")
}
Expand All @@ -70,7 +70,7 @@ class CommandsManagerAmf3: CommandsManager() {
releaseStream.addData(Amf3String(streamName))

releaseStream.writeHeader(socket)
releaseStream.writeBody(socket)
releaseStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "releaseStream")
Log.i(TAG, "send $releaseStream")

Expand All @@ -80,7 +80,7 @@ class CommandsManagerAmf3: CommandsManager() {
fcPublish.addData(Amf3String(streamName))

fcPublish.writeHeader(socket)
fcPublish.writeBody(socket)
fcPublish.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "FCPublish")
Log.i(TAG, "send $fcPublish")

Expand All @@ -89,7 +89,7 @@ class CommandsManagerAmf3: CommandsManager() {
createStream.addData(Amf3Null())

createStream.writeHeader(socket)
createStream.writeBody(socket)
createStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, "createStream")
Log.i(TAG, "send $createStream")
}
Expand Down Expand Up @@ -121,7 +121,7 @@ class CommandsManagerAmf3: CommandsManager() {
metadata.addData(amfEcmaArray)

metadata.writeHeader(socket)
metadata.writeBody(socket)
metadata.writeBody(socket, config.writeChunkSize)
Log.i(TAG, "send $metadata")
}

Expand All @@ -134,7 +134,7 @@ class CommandsManagerAmf3: CommandsManager() {
publish.addData(Amf3String("live"))

publish.writeHeader(socket)
publish.writeBody(socket)
publish.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, name)
Log.i(TAG, "send $publish")
}
Expand All @@ -145,7 +145,7 @@ class CommandsManagerAmf3: CommandsManager() {
closeStream.addData(Amf3Null())

closeStream.writeHeader(socket)
closeStream.writeBody(socket)
closeStream.writeBody(socket, config.writeChunkSize)
sessionHistory.setPacket(commandId, name)
Log.i(TAG, "send $closeStream")
}
Expand Down
Loading
Loading