Init
This commit is contained in:
@@ -0,0 +1,13 @@
|
||||
# SmartCameraStreamer
|
||||
|
||||
Small Kotlin/JVM-compatible core for the Portal streaming service. It owns independent video/audio subscriber counts, starts each encoded track on its first subscriber, stops it after the last disconnects, fans out one encoded packet stream to all subscribers, waits for an IDR for new video subscribers, and exposes `/control/mode` and `/control/fixed` endpoints.
|
||||
|
||||
The host Android service supplies `EncodedTrack` implementations backed by one shared `MediaCodec` per track and a `SmartCameraController` adapter around the Portal Smart Camera binder API.
|
||||
|
||||
Build the standalone library with:
|
||||
|
||||
```sh
|
||||
"/Applications/Android Studio.app/Contents/plugins/Kotlin/kotlinc/bin/kotlinc" \
|
||||
src/main/kotlin/com/portaltv/streamer/SmartCameraStreamer.kt \
|
||||
-d build/smart-camera-streamer.jar
|
||||
```
|
||||
@@ -0,0 +1,91 @@
|
||||
package com.portaltv.streamer
|
||||
|
||||
import java.io.Closeable
|
||||
import java.io.OutputStream
|
||||
import java.net.ServerSocket
|
||||
import java.net.Socket
|
||||
import java.net.URLDecoder
|
||||
import java.util.concurrent.CopyOnWriteArraySet
|
||||
import java.util.concurrent.LinkedBlockingDeque
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
|
||||
/** Transport-independent streaming core. Media capture is supplied by the host app. */
|
||||
class SmartCameraStreamer(
|
||||
private val camera: SmartCameraController,
|
||||
private val video: EncodedTrack,
|
||||
private val audio: EncodedTrack,
|
||||
private val controlPort: Int = 8080
|
||||
) : Closeable {
|
||||
private val videoUsers = AtomicInteger()
|
||||
private val audioUsers = AtomicInteger()
|
||||
private var server: ServerSocket? = null
|
||||
|
||||
fun start() {
|
||||
if (server != null) return
|
||||
server = ServerSocket(controlPort).also { socket ->
|
||||
Thread({
|
||||
while (!socket.isClosed) runCatching { socket.accept().also { Thread { handle(it) }.start() } }
|
||||
}, "smart-camera-http").apply { isDaemon = true; start() }
|
||||
}
|
||||
}
|
||||
|
||||
fun publishVideo(packet: ByteArray, keyFrame: Boolean, ptsUs: Long) = video.publish(packet, keyFrame, ptsUs)
|
||||
fun publishAudio(packet: ByteArray, ptsUs: Long) = audio.publish(packet, false, ptsUs)
|
||||
|
||||
private fun handle(socket: Socket) {
|
||||
socket.use {
|
||||
val reader = it.getInputStream().bufferedReader()
|
||||
val request = reader.readLine() ?: return
|
||||
val path = request.split(' ').getOrNull(1) ?: return
|
||||
while (reader.readLine()?.isNotEmpty() == true) Unit
|
||||
when {
|
||||
path.startsWith("/video.h264") -> stream(it, video, videoUsers, "video/h264", true)
|
||||
path.startsWith("/audio.aac") -> stream(it, audio, audioUsers, "audio/aac", false)
|
||||
path.startsWith("/control/") -> control(it, path)
|
||||
else -> reply(it, 404, "not found")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun stream(socket: Socket, track: EncodedTrack, users: AtomicInteger, type: String, key: Boolean) {
|
||||
val out = socket.getOutputStream()
|
||||
out.write("HTTP/1.1 200 OK\r\nContent-Type: $type\r\nCache-Control: no-store\r\nConnection: keep-alive\r\n\r\n".toByteArray()); out.flush()
|
||||
val q = track.subscribe(key); if (users.incrementAndGet() == 1) track.start()
|
||||
var waitingForKey = key
|
||||
try { while (!socket.isClosed) q.take().also { packet -> if (!waitingForKey || packet.keyFrame) { waitingForKey = false; out.write(packet.data); out.flush() } } }
|
||||
catch (_: Exception) { }
|
||||
finally { track.unsubscribe(q); if (users.decrementAndGet() == 0) track.stop() }
|
||||
}
|
||||
|
||||
private fun control(socket: Socket, path: String) {
|
||||
val query = path.substringAfter('?', "").split('&').mapNotNull { p -> p.split('=', limit = 2).takeIf { it.size == 2 }?.let { URLDecoder.decode(it[0], "UTF-8") to URLDecoder.decode(it[1], "UTF-8") } }.toMap()
|
||||
val result = when {
|
||||
path.startsWith("/control/mode") -> camera.setMode(query["mode"] ?: "")
|
||||
path.startsWith("/control/fixed") -> camera.setFixed(query["x"]?.toFloatOrNull(), query["y"]?.toFloatOrNull(), query["scale"]?.toFloatOrNull())
|
||||
else -> "unknown control endpoint"
|
||||
}
|
||||
reply(socket, if (result.startsWith("ok")) 200 else 400, result)
|
||||
}
|
||||
|
||||
private fun reply(socket: Socket, code: Int, body: String) { val out = socket.getOutputStream(); val b = body.toByteArray(); out.write("HTTP/1.1 $code OK\r\nContent-Type: text/plain\r\nContent-Length: ${b.size}\r\nConnection: close\r\n\r\n".toByteArray()); out.write(b); out.flush() }
|
||||
override fun close() { server?.close(); server = null; video.stop(); audio.stop() }
|
||||
}
|
||||
|
||||
interface SmartCameraController { fun setMode(mode: String): String; fun setFixed(x: Float?, y: Float?, scale: Float?): String }
|
||||
|
||||
interface EncodedTrack {
|
||||
data class Packet(val data: ByteArray, val keyFrame: Boolean, val ptsUs: Long)
|
||||
fun subscribe(waitForKeyFrame: Boolean): LinkedBlockingDeque<Packet>
|
||||
fun unsubscribe(queue: LinkedBlockingDeque<Packet>)
|
||||
fun publish(data: ByteArray, keyFrame: Boolean, ptsUs: Long)
|
||||
fun start(); fun stop()
|
||||
}
|
||||
|
||||
class FanoutTrack(private val capacity: Int = 24) : EncodedTrack {
|
||||
private val clients = CopyOnWriteArraySet<LinkedBlockingDeque<EncodedTrack.Packet>>()
|
||||
override fun subscribe(waitForKeyFrame: Boolean) = LinkedBlockingDeque<EncodedTrack.Packet>(capacity).also { clients += it }
|
||||
override fun unsubscribe(queue: LinkedBlockingDeque<EncodedTrack.Packet>) { clients -= queue }
|
||||
override fun publish(data: ByteArray, keyFrame: Boolean, ptsUs: Long) { clients.forEach { if (!it.offer(EncodedTrack.Packet(data, keyFrame, ptsUs))) clients -= it } }
|
||||
override fun start() {}
|
||||
override fun stop() { clients.forEach { it.clear() } }
|
||||
}
|
||||
Reference in New Issue
Block a user