From 4495711fdecf550370beff9e8c3af72d6bff616a Mon Sep 17 00:00:00 2001 From: YashasVM Date: Wed, 5 Aug 2026 22:12:12 +0530 Subject: [PATCH 1/2] test: align contracts with stream optimization --- tests/test_repo_contract.py | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/tests/test_repo_contract.py b/tests/test_repo_contract.py index 65cab08..1aa6741 100644 --- a/tests/test_repo_contract.py +++ b/tests/test_repo_contract.py @@ -32,8 +32,10 @@ def test_android_project_declares_camera_media_codec_srt_discovery_boundaries() assert "startPhoneServerIfAllowed" in app assert "MediaCodecAudioEncoder" in app stream_config = read("android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt") - assert "PreferHevc" in stream_config - assert "50_000_000" in stream_config + assert "Default1080p30" in stream_config + assert "codecPreference = CodecPreference.ForceAvc" in stream_config + assert "MIN_BITRATE_MBPS = 8" in stream_config + assert "MAX_BITRATE_MBPS = 50" in stream_config assert "OPENSTREAM_PHONE/1" in discovery assert "DISCOVERY_PORT = 51515" in discovery assert "dev.openstream.phone" in discovery @@ -119,8 +121,11 @@ def test_obs_sources_are_named_camera_slots_with_advanced_transport() -> None: assert "OBS_GROUP_CHECKABLE, advanced_group" in source assert "listener_port" in source assert "SRT latency (ms)" in source - assert '"bitrate_mbps", 50' in source - assert '"bitrate_mbps", "Expected bitrate (Mbps)", 8, 120, 1' in source + assert "kDefaultBitrateMbps = 12" in source + assert "kMinBitrateMbps = 8" in source + assert "kMaxBitrateMbps = 50" in source + assert '"bitrate_mbps"' in source + assert '"Expected bitrate (Mbps)"' in source def test_obs_discovery_beacons_advertise_slots_not_raw_listener_only() -> None: From 65e3259a3cc6c0c6a9d847954eb5dcd617454a7e Mon Sep 17 00:00:00 2001 From: YashasVM Date: Wed, 5 Aug 2026 22:45:05 +0530 Subject: [PATCH 2/2] Restore legacy Wi-Fi streaming stability --- README.md | 18 +- android/app/src/main/cpp/openstream_srt.cpp | 127 ++++++++++++-- .../java/dev/openstream/app/MainActivity.kt | 120 ++++++++++--- .../app/camera/Camera2Controller.kt | 15 ++ .../app/control/CameraControlServer.kt | 49 +++++- .../app/discovery/ObsDiscoveryClient.kt | 93 ++++++---- .../app/discovery/PhoneDiscoveryAdvertiser.kt | 52 ++++-- .../app/encoder/MediaCodecAudioEncoder.kt | 30 +++- .../app/encoder/MediaCodecVideoEncoder.kt | 77 ++++++-- .../openstream/app/stream/ConnectionTarget.kt | 3 +- .../openstream/app/stream/SrtStreamClient.kt | 24 ++- .../dev/openstream/app/stream/StreamConfig.kt | 26 +-- .../openstream/app/stream/StreamConfigTest.kt | 25 +++ docs/architecture.md | 25 ++- docs/protocol.md | 18 +- docs/testing.md | 24 ++- obs-plugin/CMakeLists.txt | 12 ++ obs-plugin/src/async-control-client.cpp | 33 +++- obs-plugin/src/async-control-client.hpp | 12 +- obs-plugin/src/media-clock.hpp | 38 ++++ obs-plugin/src/openstream-dock.cpp | 7 +- obs-plugin/src/openstream-source.cpp | 165 +++++++++++++++--- obs-plugin/tests/test_contracts.cpp | 44 +++++ tests/test_repo_contract.py | 5 +- website/src/main.jsx | 4 +- 25 files changed, 857 insertions(+), 189 deletions(-) create mode 100644 android/app/src/test/java/dev/openstream/app/stream/StreamConfigTest.kt create mode 100644 obs-plugin/src/media-clock.hpp create mode 100644 obs-plugin/tests/test_contracts.cpp diff --git a/README.md b/README.md index 7604a0a..e03fbfd 100644 --- a/README.md +++ b/README.md @@ -21,6 +21,10 @@ > [!IMPORTANT] > OpenStream V1.0.1 is a fixed patch release for local Wi-Fi camera workflows. Device-specific camera behavior and network quality can still vary; please report bugs in [GitHub Issues](https://github.com/YashasVM/OpenStream/issues). +The supported V1.0.1 path is the legacy Android app plus the legacy OBS camera +source/dock over local Wi-Fi. USB transport and USB camera-source support are +not included in this release. + ## Quick Downloads | Step | Download | Install on | @@ -38,7 +42,7 @@ Need the non-technical walkthrough with screenshots? Start with [`docs/set-up.md OpenStream V1.0.1 sends your Android phone camera directly into OBS Studio over local Wi-Fi. It uses Camera2, MediaCodec video/audio encoding, MPEG-TS muxing, SRT transport, two-way LAN discovery, source-slot pairing, and a native OBS source plugin. ```text -Phone camera -> HEVC/H.264 + AAC -> SRT over Wi-Fi -> OpenStream camera slot in OBS +Phone camera -> hardware AVC/H.264 + AAC -> SRT over Wi-Fi -> OpenStream camera slot in OBS ``` ### V1.0.1 Release @@ -97,8 +101,8 @@ Use OBS as usual. Open **Docks → OpenStream Camera Control** for the camera co | Feature | Details | |---|---| -| **Full HD Streaming** | Streams a 1080p camera feed at up to 60 fps over SRT. | -| **Hardware Encoding** | Uses MediaCodec HEVC/H.265 with H.264 fallback. | +| **Full HD Streaming** | Streams a sustainable 1080p30 camera feed over SRT, with a 720p30 fallback profile. | +| **Hardware Encoding** | Uses an explicit hardware AVC/H.264 surface encoder; unsupported hardware is reported instead of silently using software video encoding. | | **Audio Streaming** | Sends microphone audio with the video stream as AAC. | | **Multi-Lens Switching** | Supports rear, ultrawide, telephoto, and front cameras when available. | | **Pinch-to-Zoom** | Smooth digital zoom with a live zoom indicator. | @@ -147,7 +151,7 @@ Use OBS as usual. Open **Docks → OpenStream Camera Control** for the camera co ```text Android phone Camera2 preview/capture - MediaCodec HEVC/H.264 video + MediaCodec hardware AVC/H.264 video MediaCodec AAC audio MPEG-TS muxer libsrt sender @@ -168,9 +172,9 @@ Discovery uses UDP port `51515`. Camera remote controls use the phone HTTP contr | Parameter | Default | Notes | |---|---|---| | Resolution | `1920x1080` | Full HD capture target | -| Frame rate | `60 fps` | 30-60 fps depending on device support | -| Video codec | HEVC/H.265 | H.264 fallback | -| Bitrate | `50 Mbps` | High-quality local Wi-Fi default; lower if the link drops frames | +| Frame rate | `30 fps` | Bounded default to reduce heat and receiver backlog | +| Video codec | AVC/H.264 | Explicit hardware surface path | +| Bitrate | `12 Mbps` | Valid tuning range is `8-50 Mbps`; lower to 720p30 if a phone still runs hot | | SRT latency | `120 ms` | 80-200 ms useful range | | Discovery port | `51515/udp` | LAN discovery beacon | | SRT port | `9000` | Media stream | diff --git a/android/app/src/main/cpp/openstream_srt.cpp b/android/app/src/main/cpp/openstream_srt.cpp index 010f3ee..3d57f06 100644 --- a/android/app/src/main/cpp/openstream_srt.cpp +++ b/android/app/src/main/cpp/openstream_srt.cpp @@ -4,14 +4,18 @@ #include #include +#include #include +#include #include #include #include +#include #include #include #include #include +#include #include #include @@ -27,6 +31,8 @@ namespace { constexpr const char *kTag = "OpenStreamSRT"; +constexpr int kMinSrtLatencyMs = 80; +constexpr int kMaxSrtLatencyMs = 200; constexpr int kMediaCodecBufferFlagKeyFrame = 1; constexpr int kMediaCodecBufferFlagCodecConfig = 2; constexpr int kAudioSampleRate = 48000; @@ -467,7 +473,7 @@ std::optional parseSrtUrl(const std::string &url) { std::from_chars(value.data(), value.data() + value.size(), latency); if (latencyResult.ec == std::errc{} && latencyResult.ptr == value.data() + value.size()) { - parsed.latencyMs = std::clamp(latency, 80, 200); + parsed.latencyMs = std::clamp(latency, kMinSrtLatencyMs, kMaxSrtLatencyMs); } } if (amp == std::string::npos) { @@ -482,6 +488,13 @@ std::optional parseSrtUrl(const std::string &url) { class NativeSender { public: + NativeSender() : sendWorker_(&NativeSender::runSendWorker, this) {} + + ~NativeSender() { + stopSendWorker(); + disconnect(); + } + bool connect(const std::string &url) { #if OPENSTREAM_HAVE_LIBSRT disconnect(); @@ -506,11 +519,19 @@ class NativeSender { int yes = 1; int transportType = SRTT_LIVE; int payloadSize = 188 * 7; - int sendTimeoutMs = 500; + // Never let a congested Wi-Fi link stall MediaCodec callbacks. SRT's + // live late-packet drop recovers latency instead of growing it. + int sendTimeoutMs = 120; + int tooLatePacketDrop = 1; + int connectTimeoutMs = 2000; + int peerIdleTimeoutMs = 4000; srt_setsockopt(socket, 0, SRTO_TRANSTYPE, &transportType, sizeof transportType); srt_setsockopt(socket, 0, SRTO_SENDER, &yes, sizeof yes); srt_setsockopt(socket, 0, SRTO_PAYLOADSIZE, &payloadSize, sizeof payloadSize); srt_setsockopt(socket, 0, SRTO_SNDTIMEO, &sendTimeoutMs, sizeof sendTimeoutMs); + srt_setsockopt(socket, 0, SRTO_TLPKTDROP, &tooLatePacketDrop, sizeof tooLatePacketDrop); + srt_setsockopt(socket, 0, SRTO_CONNTIMEO, &connectTimeoutMs, sizeof connectTimeoutMs); + srt_setsockopt(socket, 0, SRTO_PEERIDLETIMEO, &peerIdleTimeoutMs, sizeof peerIdleTimeoutMs); const int latency = parsed->latencyMs; srt_setsockopt(socket, 0, SRTO_LATENCY, &latency, sizeof latency); srt_setsockopt(socket, 0, SRTO_PEERLATENCY, &latency, sizeof latency); @@ -572,13 +593,17 @@ class NativeSender { int yes = 1; int transportType = SRTT_LIVE; int payloadSize = 188 * 7; - int sendTimeoutMs = 500; + int sendTimeoutMs = 120; + int tooLatePacketDrop = 1; + int peerIdleTimeoutMs = 4000; const int latency = parsed->latencyMs; srt_setsockopt(listenerSocket, 0, SRTO_TRANSTYPE, &transportType, sizeof transportType); srt_setsockopt(listenerSocket, 0, SRTO_SENDER, &yes, sizeof yes); srt_setsockopt(listenerSocket, 0, SRTO_REUSEADDR, &yes, sizeof yes); srt_setsockopt(listenerSocket, 0, SRTO_PAYLOADSIZE, &payloadSize, sizeof payloadSize); srt_setsockopt(listenerSocket, 0, SRTO_SNDTIMEO, &sendTimeoutMs, sizeof sendTimeoutMs); + srt_setsockopt(listenerSocket, 0, SRTO_TLPKTDROP, &tooLatePacketDrop, sizeof tooLatePacketDrop); + srt_setsockopt(listenerSocket, 0, SRTO_PEERIDLETIMEO, &peerIdleTimeoutMs, sizeof peerIdleTimeoutMs); srt_setsockopt(listenerSocket, 0, SRTO_LATENCY, &latency, sizeof latency); srt_setsockopt(listenerSocket, 0, SRTO_PEERLATENCY, &latency, sizeof latency); @@ -609,6 +634,8 @@ class NativeSender { return false; } srt_setsockopt(acceptedSocket, 0, SRTO_SNDTIMEO, &sendTimeoutMs, sizeof sendTimeoutMs); + srt_setsockopt(acceptedSocket, 0, SRTO_TLPKTDROP, &tooLatePacketDrop, sizeof tooLatePacketDrop); + srt_setsockopt(acceptedSocket, 0, SRTO_PEERIDLETIMEO, &peerIdleTimeoutMs, sizeof peerIdleTimeoutMs); setSocket(acceptedSocket); logInfo("OBS connected to Android SRT listener"); return true; @@ -619,10 +646,10 @@ class NativeSender { #endif } - bool send(const std::vector &bytes) { + bool sendNow(const std::vector &bytes) { #if OPENSTREAM_HAVE_LIBSRT - // Closing an SRT socket or tearing down libsrt while another codec thread - // is inside srt_sendmsg is unsafe. Serialize the whole send with teardown, + // Closing an SRT socket or tearing down libsrt while the send worker is + // inside srt_sendmsg is unsafe. Serialize the whole send with teardown, // not just the socket-handle lookup. std::lock_guard ioLock(ioMutex_); const SRTSOCKET socket = currentSocket(); @@ -655,8 +682,45 @@ class NativeSender { #endif } + bool send(std::vector bytes) { +#if OPENSTREAM_HAVE_LIBSRT + if (!healthy_.load()) return false; + const size_t byteCount = bytes.size(); + { + std::lock_guard lock(sendQueueMutex_); + // This is a byte-bounded queue, not a packet-count guess: one large + // keyframe cannot turn a short receiver stall into unbounded latency. + if (sendWorkerStopping_ || + byteCount > kMaximumSendQueueBytes || + sendQueueBytes_ > kMaximumSendQueueBytes - byteCount) { + healthy_ = false; + sendQueue_.clear(); + sendQueueBytes_ = 0; + __android_log_print( + ANDROID_LOG_WARN, + kTag, + "SRT send queue saturated; dropping the session to avoid latency growth"); + return false; + } + sendQueue_.push_back(std::move(bytes)); + sendQueueBytes_ += byteCount; + } + sendQueueWake_.notify_one(); + return true; +#else + (void)bytes; + return false; +#endif + } + void disconnect() { #if OPENSTREAM_HAVE_LIBSRT + healthy_ = false; + { + std::lock_guard lock(sendQueueMutex_); + sendQueue_.clear(); + sendQueueBytes_ = 0; + } std::lock_guard ioLock(ioMutex_); const SRTSOCKET listenerSocket = takeListenerSocket(); const SRTSOCKET socket = takeSocket(); @@ -671,6 +735,39 @@ class NativeSender { } private: + void runSendWorker() { + for (;;) { + std::vector bytes; + { + std::unique_lock lock(sendQueueMutex_); + sendQueueWake_.wait(lock, [this] { + return sendWorkerStopping_ || !sendQueue_.empty(); + }); + if (sendWorkerStopping_) return; + bytes = std::move(sendQueue_.front()); + sendQueue_.pop_front(); + sendQueueBytes_ -= bytes.size(); + } + if (!sendNow(bytes)) { + healthy_ = false; + std::lock_guard lock(sendQueueMutex_); + sendQueue_.clear(); + sendQueueBytes_ = 0; + } + } + } + + void stopSendWorker() { + { + std::lock_guard lock(sendQueueMutex_); + sendWorkerStopping_ = true; + sendQueue_.clear(); + sendQueueBytes_ = 0; + } + sendQueueWake_.notify_one(); + if (sendWorker_.joinable()) sendWorker_.join(); + } + #if OPENSTREAM_HAVE_LIBSRT SRTSOCKET currentSocket() const { std::lock_guard lock(socketMutex_); @@ -680,6 +777,7 @@ class NativeSender { void setSocket(SRTSOCKET socket) { std::lock_guard lock(socketMutex_); socket_ = socket; + healthy_ = true; } SRTSOCKET takeSocket() { @@ -695,6 +793,7 @@ class NativeSender { std::lock_guard lock(socketMutex_); if (socket_ == socket) { socket_ = SRT_INVALID_SOCK; + healthy_ = false; ownsSocket = true; } } @@ -734,6 +833,14 @@ class NativeSender { SRTSOCKET socket_ = SRT_INVALID_SOCK; SRTSOCKET listener_socket_ = SRT_INVALID_SOCK; #endif + static constexpr size_t kMaximumSendQueueBytes = 768 * 1024; + std::atomic healthy_{false}; + std::mutex sendQueueMutex_; + std::condition_variable sendQueueWake_; + std::deque> sendQueue_; + size_t sendQueueBytes_ = 0; + bool sendWorkerStopping_ = false; + std::thread sendWorker_; }; struct StreamState { @@ -863,9 +970,9 @@ Java_dev_openstream_app_stream_SrtNativeBridge_sendVideo( annexB = std::move(withConfig); } - const std::vector ts = + std::vector ts = g_state.muxer->muxAccessUnit(annexB, static_cast(presentation_time_us), keyFrame); - const bool sent = g_state.sender.send(ts); + const bool sent = g_state.sender.send(std::move(ts)); if (!sent) { g_state.connected = false; } @@ -905,9 +1012,9 @@ Java_dev_openstream_app_stream_SrtNativeBridge_sendAudio( return JNI_TRUE; } - const std::vector ts = g_state.muxer->muxAudioAccessUnit( + std::vector ts = g_state.muxer->muxAudioAccessUnit( bytes, g_state.audioCodecConfig, static_cast(presentation_time_us)); - const bool sent = g_state.sender.send(ts); + const bool sent = g_state.sender.send(std::move(ts)); if (!sent) { g_state.connected = false; } diff --git a/android/app/src/main/java/dev/openstream/app/MainActivity.kt b/android/app/src/main/java/dev/openstream/app/MainActivity.kt index e1eae12..88ca123 100644 --- a/android/app/src/main/java/dev/openstream/app/MainActivity.kt +++ b/android/app/src/main/java/dev/openstream/app/MainActivity.kt @@ -71,7 +71,7 @@ class MainActivity : Activity() { private lateinit var obsDiscoveryClient: ObsDiscoveryClient private lateinit var controlServer: CameraControlServer - private val streamConfig = StreamConfig.Default1080p60 + private val streamConfig = StreamConfig.Default1080p30 private val mainHandler = Handler(Looper.getMainLooper()) @Volatile private var activeTargetName: String? = null @Volatile private var phoneServerRunning = false @@ -79,6 +79,8 @@ class MainActivity : Activity() { @Volatile private var reservedBy: String? = null @Volatile private var reservedSlotLabel: String? = null @Volatile private var listenerThread: Thread? = null + @Volatile private var callerConnectThread: Thread? = null + @Volatile private var callerGeneration = 0L @Volatile private var pendingListenerStart = false @Volatile private var listenerGeneration = 0L @Volatile private var activityStarted = false @@ -97,6 +99,7 @@ class MainActivity : Activity() { private var pendingConnectAfterSettings = false private var currentDevices: List = emptyList() private var activeStreamBitrate: Int = streamConfig.bitrate + private val callerLifecycleLock = Any() private val statsTicker = object : Runnable { override fun run() { @@ -144,13 +147,16 @@ class MainActivity : Activity() { channelCount = streamConfig.audioChannelCount, bitrate = streamConfig.audioBitrate, onEncodedAccessUnit = { accessUnit -> - streamClient.sendAudioAccessUnit(accessUnit) + if (!streamClient.sendAudioAccessUnit(accessUnit)) { + handleMediaTransportFailure() + } }, ) camera = Camera2Controller( context = this, previewSurfaceProvider = { cameraPreview.holder.surface }, lensProvider = { currentLens }, + targetFps = streamConfig.fps, ) controlServer = CameraControlServer( cameraProvider = { camera }, @@ -332,16 +338,36 @@ class MainActivity : Activity() { onEncodedAccessUnit = { accessUnit -> val sent = streamClient.sendVideoAccessUnit(accessUnit) if (!sent) { - phoneConnected = false - runOnUiThread { renderStreamStats(forceFailure = true) } + handleMediaTransportFailure() } }, ) } + private fun handleMediaTransportFailure() { + phoneConnected = false + mainHandler.post { + if (phoneServerRunning) { + if (activeTargetName != null) { + statusText.text = "Connection lost" + statusText.setTextColor(getColor(R.color.os_warning)) + statusDetail.text = "Waiting for OBS to reconnect" + } + return@post + } + if (activeTargetName == null && callerConnectThread == null) return@post + stopStream(updateStatus = false) + startPreviewIfAllowed() + startPhoneServerIfAllowed() + statusText.text = "Connection lost" + statusText.setTextColor(getColor(R.color.os_warning)) + statusDetail.text = getString(R.string.status_waiting) + } + } + private fun useStreamBitrate(bitrateMbps: Int?) { val nextBitrate = (bitrateMbps ?: streamConfig.bitrateMbps) - .coerceIn(1, 200) * 1_000_000 + .coerceIn(StreamConfig.MIN_BITRATE_MBPS, StreamConfig.MAX_BITRATE_MBPS) * 1_000_000 if (activeStreamBitrate == nextBitrate) return activeStreamBitrate = nextBitrate if (activeTargetName == null) { @@ -598,7 +624,8 @@ class MainActivity : Activity() { val sourceInstanceId = uri.getQueryParameter("sourceInstanceId")?.trim().orEmpty() if (sourceInstanceId.isNotBlank()) { val slotLabel = uri.getQueryParameter("slotLabel")?.trim().orEmpty() - val bitrateMbps = uri.getQueryParameter("bitrateMbps")?.toIntOrNull()?.coerceIn(1, 200) + val bitrateMbps = uri.getQueryParameter("bitrateMbps")?.toIntOrNull() + ?.coerceIn(StreamConfig.MIN_BITRATE_MBPS, StreamConfig.MAX_BITRATE_MBPS) if (reserveForSource(sourceInstanceId, slotLabel, bitrateMbps)) { statusText.text = "Paired to ${slotLabel.ifBlank { "OBS slot" }}" statusDetail.text = "Waiting for OBS to go live" @@ -610,34 +637,61 @@ class MainActivity : Activity() { } private fun startStream(target: ConnectionTarget) { + stopStream(updateStatus = false) // Caller mode and listener mode share one native SRT transport. Fully // stop the listener before opening a manual caller connection. stopPhoneServer(clearReservation = true, updateStatus = false) useStreamBitrate(target.bitrateMbps) statusText.text = "Connecting…" statusDetail.text = "${currentLens.displayName} → ${target.name}" - runCatching { - streamClient.connect( - url = target.toSrtCallerUrl(), - codecMime = encoder.codecName, - width = streamConfig.width, - height = streamConfig.height, - fps = streamConfig.fps, - ) - encoder.start() - startAudioIfAllowed() - camera.startStreaming(encoder.inputSurface()) - activeTargetName = target.name - mainHandler.removeCallbacks(statsTicker) - mainHandler.post(statsTicker) - showLiveState(target.name) - }.onFailure { error -> - stopStream() - startPreviewIfAllowed() - startPhoneServerIfAllowed() - statusText.text = "Connection failed" - statusDetail.text = error.message ?: "Unknown error" + val generation = callerGeneration + 1 + callerGeneration = generation + val thread = Thread({ + try { + streamClient.connect( + url = target.toSrtCallerUrl(), + codecMime = encoder.codecName, + width = streamConfig.width, + height = streamConfig.height, + fps = streamConfig.fps, + ) + synchronized(callerLifecycleLock) { + check(callerGeneration == generation) { "SRT caller connection was cancelled" } + encoder.start() + startAudioIfAllowed() + camera.startStreaming(encoder.inputSurface()) + } + mainHandler.post { + if (callerGeneration != generation) return@post + activeTargetName = target.name + mainHandler.removeCallbacks(statsTicker) + mainHandler.post(statsTicker) + showLiveState(target.name) + } + } catch (error: Throwable) { + synchronized(callerLifecycleLock) { + if (callerGeneration == generation) { + streamClient.disconnect() + stopActiveEncoding(updateStatus = false) + } + } + mainHandler.post { + if (callerGeneration != generation) return@post + startPreviewIfAllowed() + startPhoneServerIfAllowed() + statusText.text = "Connection failed" + statusDetail.text = error.message ?: "Unknown error" + } + } finally { + if (callerConnectThread === Thread.currentThread()) { + callerConnectThread = null + } + } + }, "OpenStreamPhoneSrtCaller").apply { + isDaemon = true } + callerConnectThread = thread + thread.start() } private fun startAudioIfAllowed() { @@ -759,6 +813,8 @@ class MainActivity : Activity() { clearReservation: Boolean = true, updateStatus: Boolean = true, ) { + callerGeneration += 1 + callerConnectThread?.interrupt() pendingListenerStart = false listenerGeneration += 1 phoneServerRunning = false @@ -778,17 +834,23 @@ class MainActivity : Activity() { } else if (listenerThread === thread) { listenerThread = null } - stopActiveEncoding(updateStatus) + synchronized(callerLifecycleLock) { + stopActiveEncoding(updateStatus) + } hideLiveState() btnStop.visibility = View.GONE } private fun stopStream(updateStatus: Boolean = true) { + callerGeneration += 1 + callerConnectThread?.interrupt() activeTargetName = null mainHandler.removeCallbacks(statsTicker) phoneConnected = false streamClient.disconnect() - stopActiveEncoding(updateStatus) + synchronized(callerLifecycleLock) { + stopActiveEncoding(updateStatus) + } hideLiveState() } diff --git a/android/app/src/main/java/dev/openstream/app/camera/Camera2Controller.kt b/android/app/src/main/java/dev/openstream/app/camera/Camera2Controller.kt index 699054f..df2acba 100644 --- a/android/app/src/main/java/dev/openstream/app/camera/Camera2Controller.kt +++ b/android/app/src/main/java/dev/openstream/app/camera/Camera2Controller.kt @@ -12,6 +12,7 @@ import android.os.Build import android.os.Handler import android.os.HandlerThread import android.util.Log +import android.util.Range import android.view.Surface /** @@ -27,6 +28,7 @@ class Camera2Controller( private val context: Context, private val previewSurfaceProvider: () -> Surface, private val lensProvider: () -> CameraLens = { CameraLens.Back }, + private val targetFps: Int = 30, ) { private val cameraManager = context.getSystemService(CameraManager::class.java) private val thread = HandlerThread("OpenStreamCamera") @@ -43,6 +45,7 @@ class Camera2Controller( private var minZoomRatio = 1.0f private var sensorRect: Rect? = null private var supportsZoomRatioKey = false + private var captureFpsRange: Range? = null // Torch state private var torchEnabled = false @@ -228,6 +231,10 @@ class Camera2Controller( private fun loadZoomCapabilities(cameraId: String) { val chars = cameraManager.getCameraCharacteristics(cameraId) + captureFpsRange = chars + .get(CameraCharacteristics.CONTROL_AE_AVAILABLE_TARGET_FPS_RANGES) + ?.filter { range -> range.upper <= targetFps } + ?.maxWithOrNull(compareBy> { it.upper }.thenBy { it.lower }) sensorRect = chars.get(CameraCharacteristics.SENSOR_INFO_ACTIVE_ARRAY_SIZE) if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.R) { @@ -272,6 +279,7 @@ class Camera2Controller( set(CaptureRequest.CONTROL_AF_MODE, CaptureRequest.CONTROL_AF_MODE_CONTINUOUS_VIDEO) set(CaptureRequest.CONTROL_AE_MODE, CaptureRequest.CONTROL_AE_MODE_ON) set(CaptureRequest.CONTROL_AWB_MODE, CaptureRequest.CONTROL_AWB_MODE_AUTO) + applyFrameRate(this) applyZoom(this) applyTorch(this) }.build() @@ -303,6 +311,12 @@ class Camera2Controller( } } + private fun applyFrameRate(builder: CaptureRequest.Builder) { + captureFpsRange?.let { range -> + builder.set(CaptureRequest.CONTROL_AE_TARGET_FPS_RANGE, range) + } + } + private fun applyTorch(builder: CaptureRequest.Builder) { if (torchEnabled) { builder.set(CaptureRequest.FLASH_MODE, CaptureRequest.FLASH_MODE_TORCH) @@ -327,6 +341,7 @@ class Camera2Controller( set(CaptureRequest.CONTROL_AF_MODE, CaptureRequest.CONTROL_AF_MODE_CONTINUOUS_VIDEO) set(CaptureRequest.CONTROL_AE_MODE, CaptureRequest.CONTROL_AE_MODE_ON) set(CaptureRequest.CONTROL_AWB_MODE, CaptureRequest.CONTROL_AWB_MODE_AUTO) + applyFrameRate(this) applyZoom(this) applyTorch(this) }.build() diff --git a/android/app/src/main/java/dev/openstream/app/control/CameraControlServer.kt b/android/app/src/main/java/dev/openstream/app/control/CameraControlServer.kt index 4d1c2a6..d2f90e0 100644 --- a/android/app/src/main/java/dev/openstream/app/control/CameraControlServer.kt +++ b/android/app/src/main/java/dev/openstream/app/control/CameraControlServer.kt @@ -3,6 +3,7 @@ package dev.openstream.app.control import android.util.Log import dev.openstream.app.camera.Camera2Controller import dev.openstream.app.camera.CameraLens +import dev.openstream.app.stream.StreamConfig import org.json.JSONObject import java.io.BufferedInputStream import java.io.OutputStreamWriter @@ -18,7 +19,7 @@ import java.util.concurrent.atomic.AtomicBoolean * - POST /zoom {"value": 2.5} * - POST /torch {"enabled": true} * - POST /lens {"lens": "Back"} - * - POST /reserve {"sourceInstanceId": "...", "bitrateMbps": 50} + * - POST /reserve {"sourceInstanceId": "...", "bitrateMbps": 12} * - POST /release {"sourceInstanceId": "..."} * - POST /identify {"label": "CAM B", "subtitle": "Close-up"} * - GET /status Returns current camera state @@ -36,10 +37,16 @@ class CameraControlServer( private val onIdentify: (String, String) -> Unit, ) { private val running = AtomicBoolean(false) - private var serverSocket: ServerSocket? = null - private var worker: Thread? = null + @Volatile private var serverSocket: ServerSocket? = null + @Volatile private var activeClient: Socket? = null + @Volatile private var worker: Thread? = null fun start() { + if (running.get()) return + if (worker?.isAlive == true) { + Log.w(TAG, "Control server worker is still stopping; delaying restart") + return + } if (!running.compareAndSet(false, true)) return worker = Thread(::run, "OpenStreamControlServer").apply { isDaemon = true @@ -50,13 +57,23 @@ class CameraControlServer( fun stop() { running.set(false) runCatching { serverSocket?.close() } + runCatching { activeClient?.close() } + val thread = worker + thread?.interrupt() + if (thread != null && thread !== Thread.currentThread()) { + runCatching { thread.join(STOP_TIMEOUT_MS) } + .onFailure { Thread.currentThread().interrupt() } + } + if (worker === thread && thread?.isAlive != true) worker = null serverSocket = null - worker = null + activeClient = null } private fun run() { + var openedSocket: ServerSocket? = null try { val socket = ServerSocket(port) + openedSocket = socket serverSocket = socket socket.soTimeout = 1000 Log.i(TAG, "Camera control server listening on port $port") @@ -70,10 +87,22 @@ class CameraControlServer( if (running.get()) Log.w(TAG, "Socket error", e) break } - handleClient(client) + activeClient = client + try { + handleClient(client) + } finally { + activeClient = null + } } } catch (e: Exception) { - Log.e(TAG, "Control server error", e) + if (running.get()) Log.e(TAG, "Control server error", e) + } finally { + runCatching { openedSocket?.close() } + if (serverSocket === openedSocket) serverSocket = null + if (worker === Thread.currentThread()) { + worker = null + } + running.set(false) } } @@ -230,7 +259,12 @@ class CameraControlServer( val sourceInstanceId = json.optString("sourceInstanceId").trim() if (sourceInstanceId.isEmpty()) return """{"error":"missing sourceInstanceId"}""" val slotLabel = json.optString("slotLabel", "") - val bitrateMbps = if (json.has("bitrateMbps")) json.optInt("bitrateMbps").coerceIn(1, 200) else null + val bitrateMbps = if (json.has("bitrateMbps")) { + json.optInt("bitrateMbps").coerceIn( + StreamConfig.MIN_BITRATE_MBPS, + StreamConfig.MAX_BITRATE_MBPS, + ) + } else null val accepted = onReserve(sourceInstanceId, slotLabel, bitrateMbps) return if (accepted) { JSONObject() @@ -268,5 +302,6 @@ class CameraControlServer( private const val MAX_HEADER_LINE_BYTES = 2_048 private const val MAX_HEADER_BYTES = 8_192 private const val MAX_BODY_BYTES = 8_192 + private const val STOP_TIMEOUT_MS = 1_000L } } diff --git a/android/app/src/main/java/dev/openstream/app/discovery/ObsDiscoveryClient.kt b/android/app/src/main/java/dev/openstream/app/discovery/ObsDiscoveryClient.kt index f4511e9..59fe3be 100644 --- a/android/app/src/main/java/dev/openstream/app/discovery/ObsDiscoveryClient.kt +++ b/android/app/src/main/java/dev/openstream/app/discovery/ObsDiscoveryClient.kt @@ -4,9 +4,10 @@ import android.content.Context import android.net.wifi.WifiManager import android.os.Handler import android.os.Looper +import android.util.Log +import dev.openstream.app.stream.StreamConfig import org.json.JSONObject import java.net.DatagramPacket -import java.net.DatagramSocket import java.net.InetSocketAddress import java.net.InetAddress import java.net.MulticastSocket @@ -23,11 +24,16 @@ class ObsDiscoveryClient( private val running = AtomicBoolean(false) private val mainHandler = Handler(Looper.getMainLooper()) private val devices = linkedMapOf() - private var socket: MulticastSocket? = null - private var worker: Thread? = null + @Volatile private var socket: MulticastSocket? = null + @Volatile private var worker: Thread? = null private var multicastLock: WifiManager.MulticastLock? = null fun start() { + if (running.get()) return + if (worker?.isAlive == true) { + Log.w(TAG, "Discovery worker is still stopping; delaying restart") + return + } if (!running.compareAndSet(false, true)) return acquireMulticastLock() worker = Thread(::receiveLoop, "OpenStreamDiscovery").apply { @@ -37,10 +43,16 @@ class ObsDiscoveryClient( } fun stop() { - if (!running.compareAndSet(true, false)) return + running.set(false) socket?.close() socket = null - worker = null + val thread = worker + thread?.interrupt() + if (thread != null && thread !== Thread.currentThread()) { + runCatching { thread.join(STOP_TIMEOUT_MS) } + .onFailure { Thread.currentThread().interrupt() } + } + if (worker === thread && thread?.isAlive != true) worker = null synchronized(devices) { devices.clear() } @@ -66,40 +78,50 @@ class ObsDiscoveryClient( } private fun receiveLoop() { - val udp = MulticastSocket(null).apply { - reuseAddress = true - soTimeout = 500 - bind(InetSocketAddress(DISCOVERY_PORT)) - } - socket = udp - joinDiscoveryMulticast(udp) - val buffer = ByteArray(4096) - - while (running.get()) { - try { - val packet = DatagramPacket(buffer, buffer.size) - udp.receive(packet) - val payload = String(packet.data, packet.offset, packet.length, StandardCharsets.UTF_8) - val host = packet.address.hostAddress ?: continue - val device = ObsDiscoveryProtocol.parseBeacon(payload, host, nowMs()) ?: continue - synchronized(devices) { - devices[device.instanceId.ifBlank { "${device.host}:${device.port}" }] = device - } - pruneExpired() - publishDevices() - } catch (_: SocketTimeoutException) { - if (pruneExpired()) { + var udp: MulticastSocket? = null + try { + val multicast = MulticastSocket(null).apply { + reuseAddress = true + soTimeout = 500 + bind(InetSocketAddress(DISCOVERY_PORT)) + } + udp = multicast + socket = multicast + joinDiscoveryMulticast(multicast) + val buffer = ByteArray(4096) + + while (running.get()) { + try { + val packet = DatagramPacket(buffer, buffer.size) + multicast.receive(packet) + val payload = String(packet.data, packet.offset, packet.length, StandardCharsets.UTF_8) + val host = packet.address.hostAddress ?: continue + val device = ObsDiscoveryProtocol.parseBeacon(payload, host, nowMs()) ?: continue + synchronized(devices) { + devices[device.instanceId.ifBlank { "${device.host}:${device.port}" }] = device + } + pruneExpired() publishDevices() - } - } catch (_: Exception) { - if (running.get()) { + } catch (_: SocketTimeoutException) { if (pruneExpired()) { publishDevices() } + } catch (_: Exception) { + if (running.get()) { + if (pruneExpired()) { + publishDevices() + } + } } } + } catch (e: Exception) { + if (running.get()) Log.w(TAG, "Discovery worker failed", e) + } finally { + runCatching { udp?.close() } + if (socket === udp) socket = null + if (worker === Thread.currentThread()) worker = null + running.set(false) } - udp.close() } private fun joinDiscoveryMulticast(socket: MulticastSocket) { @@ -135,9 +157,11 @@ class ObsDiscoveryClient( } companion object { + private const val TAG = "OpenStreamDiscovery" const val DISCOVERY_PORT = 51515 const val DISCOVERY_MULTICAST_ADDRESS = "239.255.42.99" const val DEVICE_TTL_MS = 5_000L + private const val STOP_TIMEOUT_MS = 1_000L } } @@ -161,7 +185,10 @@ object ObsDiscoveryProtocol { host = host, port = port, latencyMs = json.optInt("latencyMs", 120).coerceIn(80, 200), - bitrateMbps = json.optInt("bitrateMbps", 12).coerceIn(1, 200), + bitrateMbps = json.optInt("bitrateMbps", StreamConfig.Default1080p30.bitrateMbps).coerceIn( + StreamConfig.MIN_BITRATE_MBPS, + StreamConfig.MAX_BITRATE_MBPS, + ), instanceId = json.optString("instanceId", "$packetHost:$port"), sourceInstanceId = json.optString("sourceInstanceId", json.optString("instanceId", "$packetHost:$port")), slotId = json.optString("slotId", json.optString("instanceId", "$packetHost:$port")), diff --git a/android/app/src/main/java/dev/openstream/app/discovery/PhoneDiscoveryAdvertiser.kt b/android/app/src/main/java/dev/openstream/app/discovery/PhoneDiscoveryAdvertiser.kt index dd69598..4a79e35 100644 --- a/android/app/src/main/java/dev/openstream/app/discovery/PhoneDiscoveryAdvertiser.kt +++ b/android/app/src/main/java/dev/openstream/app/discovery/PhoneDiscoveryAdvertiser.kt @@ -21,10 +21,12 @@ class PhoneDiscoveryAdvertiser( private val reservedByProvider: () -> String? = { null }, ) { private val running = AtomicBoolean(false) - private var worker: Thread? = null + @Volatile private var worker: Thread? = null private val instanceId = UUID.randomUUID().toString() fun start() { + if (running.get()) return + if (worker?.isAlive == true) return if (!running.compareAndSet(false, true)) return worker = Thread(::run, "OpenStreamPhoneAdvertiser").apply { isDaemon = true @@ -34,27 +36,46 @@ class PhoneDiscoveryAdvertiser( fun stop() { running.set(false) - worker = null + val thread = worker + thread?.interrupt() + if (thread != null && thread !== Thread.currentThread()) { + runCatching { thread.join(STOP_TIMEOUT_MS) } + .onFailure { Thread.currentThread().interrupt() } + } + if (worker === thread && thread?.isAlive != true) worker = null } private fun run() { - val socket = DatagramSocket().apply { - broadcast = true + val socket = runCatching { + DatagramSocket().apply { broadcast = true } + }.getOrElse { + running.set(false) + if (worker === Thread.currentThread()) worker = null + return } - val destinations = listOf( - InetAddress.getByName("255.255.255.255"), - InetAddress.getByName(DISCOVERY_MULTICAST_ADDRESS), - ) - while (running.get()) { - val bytes = beaconPayload().toByteArray(StandardCharsets.UTF_8) - destinations.forEach { destination -> - runCatching { - socket.send(DatagramPacket(bytes, bytes.size, destination, DISCOVERY_PORT)) + try { + val destinations = listOf( + InetAddress.getByName("255.255.255.255"), + InetAddress.getByName(DISCOVERY_MULTICAST_ADDRESS), + ) + while (running.get()) { + val bytes = beaconPayload().toByteArray(StandardCharsets.UTF_8) + destinations.forEach { destination -> + runCatching { + socket.send(DatagramPacket(bytes, bytes.size, destination, DISCOVERY_PORT)) + } + } + try { + Thread.sleep(1_000) + } catch (_: InterruptedException) { + break } } - Thread.sleep(1_000) + } finally { + socket.close() + if (worker === Thread.currentThread()) worker = null + running.set(false) } - socket.close() } private fun beaconPayload(): String { @@ -94,5 +115,6 @@ class PhoneDiscoveryAdvertiser( const val DISCOVERY_MULTICAST_ADDRESS = "239.255.42.99" const val PREFIX = "OPENSTREAM_PHONE/1" const val TYPE = "dev.openstream.phone" + private const val STOP_TIMEOUT_MS = 1_000L } } diff --git a/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecAudioEncoder.kt b/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecAudioEncoder.kt index b798f28..45e8168 100644 --- a/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecAudioEncoder.kt +++ b/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecAudioEncoder.kt @@ -8,6 +8,7 @@ import android.media.AudioFormat import android.media.AudioRecord import android.media.MediaCodec import android.media.MediaCodecInfo +import android.media.MediaCodecList import android.media.MediaFormat import android.media.MediaRecorder import android.os.Build @@ -24,7 +25,7 @@ class MediaCodecAudioEncoder( context: Context, private val sampleRate: Int = 48_000, private val channelCount: Int = 1, - private val bitrate: Int = 192_000, + private val bitrate: Int = 128_000, private val onEncodedAccessUnit: (EncodedAccessUnit) -> Unit, ) { private val context = context.applicationContext @@ -52,7 +53,7 @@ class MediaCodecAudioEncoder( } } - val encoder = MediaCodec.createEncoderByType(mime) + val encoder = createAudioCodec(mime) codec = encoder encoder.configure(format, null, null, MediaCodec.CONFIGURE_FLAG_ENCODE) encoder.start() @@ -66,7 +67,9 @@ class MediaCodecAudioEncoder( sampleRate, channelConfig, AudioFormat.ENCODING_PCM_16BIT ) check(minBufferSize > 0) { "AudioRecord does not support $sampleRate Hz / $channelCount ch PCM16" } - val bufferSize = max(minBufferSize * 4, bytesForDurationMs(250)) + // Keep capture buffering below 80 ms. READ_BLOCKING applies backpressure + // once this capacity is reached instead of allowing audio to trail video. + val bufferSize = max(minBufferSize * 2, bytesForDurationMs(MAX_CAPTURE_BUFFER_MS)) val recorder = createRecorder(channelConfig, bufferSize) audioRecord = recorder @@ -114,12 +117,12 @@ class MediaCodecAudioEncoder( fun stop() { running = false - captureThread?.join(500) - captureThread = null val recorder = audioRecord audioRecord = null runCatching { recorder?.stop() } runCatching { recorder?.release() } + captureThread?.join(500) + captureThread = null val encoder = codec codec = null @@ -185,6 +188,22 @@ class MediaCodecAudioEncoder( return bytes.takeIf { it.isNotEmpty() } } + private fun createAudioCodec(mime: String): MediaCodec { + val hardware = MediaCodecList(MediaCodecList.REGULAR_CODECS).codecInfos + .firstOrNull { info -> + info.isEncoder && + info.isHardwareAccelerated && + !info.isSoftwareOnly && + info.supportedTypes.any { it.equals(mime, ignoreCase = true) } + } + if (hardware != null) { + Log.i(TAG, "Using hardware audio encoder ${hardware.name} for $mime") + return MediaCodec.createByCodecName(hardware.name) + } + Log.w(TAG, "No hardware AAC encoder is available; explicitly using the platform software encoder") + return MediaCodec.createEncoderByType(mime) + } + @SuppressLint("MissingPermission") private fun createRecorder(channelConfig: Int, bufferSize: Int): AudioRecord { val sources = buildList { @@ -227,5 +246,6 @@ class MediaCodecAudioEncoder( companion object { private const val TAG = "OpenStreamAudioEncoder" private const val BYTES_PER_PCM16_SAMPLE = 2 + private const val MAX_CAPTURE_BUFFER_MS = 80 } } diff --git a/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecVideoEncoder.kt b/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecVideoEncoder.kt index 588cf43..84ab557 100644 --- a/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecVideoEncoder.kt +++ b/android/app/src/main/java/dev/openstream/app/encoder/MediaCodecVideoEncoder.kt @@ -29,6 +29,11 @@ data class EncodedAccessUnit( val flags: Int, ) +private data class EncoderSelection( + val mimeType: String, + val codecName: String, +) + class MediaCodecVideoEncoder( private val preference: CodecPreference, private val width: Int, @@ -38,7 +43,8 @@ class MediaCodecVideoEncoder( private val keyframeIntervalSeconds: Int, private val onEncodedAccessUnit: (EncodedAccessUnit) -> Unit, ) { - private var mimeType = chooseMimeType(preference) + private var selection = chooseEncoder(preference) + private var mimeType = selection.mimeType private var codec: MediaCodec? = null private val thread = HandlerThread("OpenStreamEncoder") private lateinit var handler: Handler @@ -54,8 +60,10 @@ class MediaCodecVideoEncoder( if (codec != null) { stop() } - mimeType = chooseMimeType(preference) - val encoder = MediaCodec.createEncoderByType(mimeType) + selection = chooseEncoder(preference) + mimeType = selection.mimeType + Log.i("OpenStreamEncoder", "Using hardware encoder ${selection.codecName} for $mimeType") + val encoder = MediaCodec.createByCodecName(selection.codecName) codec = encoder val format = MediaFormat.createVideoFormat(mimeType, width, height).apply { setInteger(MediaFormat.KEY_COLOR_FORMAT, MediaCodecInfo.CodecCapabilities.COLOR_FormatSurface) @@ -67,8 +75,15 @@ class MediaCodecVideoEncoder( setInteger(MediaFormat.KEY_LATENCY, 0) } } - encoder.configure(format, null, null, MediaCodec.CONFIGURE_FLAG_ENCODE) - surface = encoder.createInputSurface() + try { + encoder.configure(format, null, null, MediaCodec.CONFIGURE_FLAG_ENCODE) + surface = encoder.createInputSurface() + } catch (error: Throwable) { + codec = null + surface = null + runCatching { encoder.release() } + throw error + } if (!thread.isAlive) { thread.start() @@ -116,7 +131,14 @@ class MediaCodecVideoEncoder( } } }, handler) - encoder.start() + try { + encoder.start() + } catch (error: Throwable) { + codec = null + surface = null + runCatching { encoder.release() } + throw error + } } fun stop() { @@ -127,14 +149,45 @@ class MediaCodecVideoEncoder( runCatching { encoder.release() } } - private fun chooseMimeType(preference: CodecPreference): String { - if (preference == CodecPreference.ForceAvc) { - return MediaFormat.MIMETYPE_VIDEO_AVC + private fun chooseEncoder(preference: CodecPreference): EncoderSelection { + val mimeTypes = when (preference) { + CodecPreference.ForceAvc -> listOf(MediaFormat.MIMETYPE_VIDEO_AVC) + CodecPreference.PreferHevc -> listOf( + MediaFormat.MIMETYPE_VIDEO_HEVC, + MediaFormat.MIMETYPE_VIDEO_AVC, + ) } - val hasHevc = MediaCodecList(MediaCodecList.REGULAR_CODECS).codecInfos.any { info -> - info.isEncoder && info.supportedTypes.any { it.equals(MediaFormat.MIMETYPE_VIDEO_HEVC, true) } + val infos = MediaCodecList(MediaCodecList.REGULAR_CODECS).codecInfos + for (mime in mimeTypes) { + val info = infos.firstOrNull { candidate -> + candidate.isEncoder && + candidate.isHardwareAccelerated && + !candidate.isSoftwareOnly && + candidate.supportedTypes.any { it.equals(mime, true) } && + runCatching { + candidate.getCapabilitiesForType(mime).colorFormats.contains( + MediaCodecInfo.CodecCapabilities.COLOR_FormatSurface, + ) + }.getOrDefault(false) + } ?: continue + if (preference == CodecPreference.PreferHevc && + mime == MediaFormat.MIMETYPE_VIDEO_AVC + ) { + Log.w( + "OpenStreamEncoder", + "Hardware HEVC is unavailable; explicitly falling back to hardware AVC", + ) + } + return EncoderSelection(mime, info.name) } - return if (hasHevc) MediaFormat.MIMETYPE_VIDEO_HEVC else MediaFormat.MIMETYPE_VIDEO_AVC + + Log.e( + "OpenStreamEncoder", + "No hardware surface video encoder is available; refusing a silent software fallback", + ) + throw IllegalStateException( + "OpenStream needs a hardware AVC/HEVC surface encoder on this phone", + ) } private fun codecConfigFrom(format: MediaFormat): ByteArray? { diff --git a/android/app/src/main/java/dev/openstream/app/stream/ConnectionTarget.kt b/android/app/src/main/java/dev/openstream/app/stream/ConnectionTarget.kt index 6e90e71..e1ade02 100644 --- a/android/app/src/main/java/dev/openstream/app/stream/ConnectionTarget.kt +++ b/android/app/src/main/java/dev/openstream/app/stream/ConnectionTarget.kt @@ -34,7 +34,8 @@ data class ConnectionTarget( if (host.isBlank()) return null val port = uri.getQueryParameter("port")?.toIntOrNull()?.coerceIn(1, 65535) ?: DEFAULT_PORT val latencyMs = uri.getQueryParameter("latency")?.toIntOrNull()?.coerceIn(80, 200) ?: DEFAULT_LATENCY_MS - val bitrateMbps = uri.getQueryParameter("bitrateMbps")?.toIntOrNull()?.coerceIn(1, 200) + val bitrateMbps = uri.getQueryParameter("bitrateMbps")?.toIntOrNull() + ?.coerceIn(StreamConfig.MIN_BITRATE_MBPS, StreamConfig.MAX_BITRATE_MBPS) val name = uri.getQueryParameter("name")?.ifBlank { DEFAULT_NAME } ?: DEFAULT_NAME return ConnectionTarget( name = name, diff --git a/android/app/src/main/java/dev/openstream/app/stream/SrtStreamClient.kt b/android/app/src/main/java/dev/openstream/app/stream/SrtStreamClient.kt index d690e9c..5999bd5 100644 --- a/android/app/src/main/java/dev/openstream/app/stream/SrtStreamClient.kt +++ b/android/app/src/main/java/dev/openstream/app/stream/SrtStreamClient.kt @@ -53,7 +53,9 @@ class SrtStreamClient { } fun sendVideoAccessUnit(accessUnit: EncodedAccessUnit): Boolean { - if (!connected) return false + val generation = synchronized(stateLock) { + if (!connected) null else sessionGeneration.get() + } ?: return false val sent = SrtNativeBridge.sendVideo(accessUnit.data, accessUnit.presentationTimeUs, accessUnit.flags) val isCodecConfig = (accessUnit.flags and BUFFER_FLAG_CODEC_CONFIG) != 0 if (sent) { @@ -69,13 +71,21 @@ class SrtStreamClient { } } else { sendFailures.incrementAndGet() + markSendFailure(generation) } return sent } fun sendAudioAccessUnit(accessUnit: EncodedAccessUnit): Boolean { - if (!connected) return false - return SrtNativeBridge.sendAudio(accessUnit.data, accessUnit.presentationTimeUs, accessUnit.flags) + val generation = synchronized(stateLock) { + if (!connected) null else sessionGeneration.get() + } ?: return false + val sent = SrtNativeBridge.sendAudio(accessUnit.data, accessUnit.presentationTimeUs, accessUnit.flags) + if (!sent) { + sendFailures.incrementAndGet() + markSendFailure(generation) + } + return sent } fun disconnect() { @@ -119,6 +129,14 @@ class SrtStreamClient { lastPresentationTimeUs.set(0) } + private fun markSendFailure(generation: Long) { + synchronized(stateLock) { + if (sessionGeneration.get() == generation) { + connected = false + } + } + } + companion object { private const val BUFFER_FLAG_KEY_FRAME = 1 private const val BUFFER_FLAG_CODEC_CONFIG = 2 diff --git a/android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt b/android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt index 0306123..0d2cdf0 100644 --- a/android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt +++ b/android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt @@ -18,30 +18,36 @@ data class StreamConfig( get() = bitrate / 1_000_000 companion object { - val Default1080p60 = StreamConfig( + const val MIN_BITRATE_MBPS = 8 + const val MAX_BITRATE_MBPS = 50 + + // The sustainable default is 1080p30 AVC. 60 fps and HEVC remain + // useful device-specific choices, but they should not be forced on a + // phone before its thermal and hardware limits are known. + val Default1080p30 = StreamConfig( width = 1920, height = 1080, - fps = 60, - bitrate = 50_000_000, + fps = 30, + bitrate = 12_000_000, keyframeIntervalSeconds = 1, latencyMs = 120, - codecPreference = CodecPreference.PreferHevc, + codecPreference = CodecPreference.ForceAvc, audioSampleRate = 48_000, audioChannelCount = 1, - audioBitrate = 192_000, + audioBitrate = 128_000, ) - val Fallback720p60 = StreamConfig( + val Fallback720p30 = StreamConfig( width = 1280, height = 720, - fps = 60, - bitrate = 25_000_000, + fps = 30, + bitrate = 8_000_000, keyframeIntervalSeconds = 1, latencyMs = 120, - codecPreference = CodecPreference.PreferHevc, + codecPreference = CodecPreference.ForceAvc, audioSampleRate = 48_000, audioChannelCount = 1, - audioBitrate = 192_000, + audioBitrate = 128_000, ) } } diff --git a/android/app/src/test/java/dev/openstream/app/stream/StreamConfigTest.kt b/android/app/src/test/java/dev/openstream/app/stream/StreamConfigTest.kt new file mode 100644 index 0000000..df815be --- /dev/null +++ b/android/app/src/test/java/dev/openstream/app/stream/StreamConfigTest.kt @@ -0,0 +1,25 @@ +package dev.openstream.app.stream + +import dev.openstream.app.encoder.CodecPreference +import org.junit.Assert.assertEquals +import org.junit.Test + +class StreamConfigTest { + @Test + fun defaultProfileIsBoundedForSustainablePhoneStreaming() { + val config = StreamConfig.Default1080p30 + + assertEquals(30, config.fps) + assertEquals(12, config.bitrateMbps) + assertEquals(CodecPreference.ForceAvc, config.codecPreference) + assertEquals(1, config.audioChannelCount) + assertEquals(128_000, config.audioBitrate) + } + + @Test + fun bitrateBoundsMatchPairingAndManualTargetValidation() { + assertEquals(8, StreamConfig.MIN_BITRATE_MBPS) + assertEquals(50, StreamConfig.MAX_BITRATE_MBPS) + } + +} diff --git a/docs/architecture.md b/docs/architecture.md index 4db8372..b0fe11d 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -4,7 +4,7 @@ ```text Android Camera2 - -> MediaCodec hardware HEVC/H.264 video encode + -> MediaCodec hardware AVC/H.264 video encode -> MediaCodec AAC audio encode -> MPEG-TS muxer (video + audio) -> SRT caller over local Wi-Fi @@ -13,8 +13,12 @@ Android Camera2 -> OBS final stream/record encode ``` -The V1 prototype uses hardware video and audio encoding on the phone. This is -the best path for stable, high-quality wireless streaming over commodity Wi-Fi. +The V1 path uses hardware video and audio encoding on the phone. AVC/H.264 at +1080p30 is the sustainable default; the app does not silently switch video to +software encoding when a hardware surface encoder is unavailable. + +USB transport is intentionally outside this release. The maintained path is +the legacy Android app and OBS source/dock over local Wi-Fi and SRT. ## Discovery & Connection @@ -67,13 +71,13 @@ transport can be tested without blocking the reliable product path. ## Encoding Defaults -- Preferred video codec: HEVC/H.265. -- Fallback video codec: AVC/H.264. +- Default video codec: AVC/H.264. +- HEVC/H.265 remains an explicit device/profile option only; it is not forced by the default path. - Audio codec: AAC. - Default resolution: `1920x1080`. - Fallback resolution: `1280x720`. -- Default frame rate: `60 fps`. -- Bitrate presets: `8-120 Mbps`; the default high-quality path uses `50 Mbps`. +- Default frame rate: `30 fps`. +- Bitrate range: `8-50 Mbps`; the default path uses `12 Mbps`. - Keyframe interval: `1 second`. ## Video Output Path @@ -91,8 +95,11 @@ extracted from FFmpeg frame properties and forwarded to OBS. ## Sync Model V1 does not promise PTP or genlock. Each Android frame/audio buffer is -timestamped with a monotonic clock. The Windows receiver and OBS source use -those timestamps to estimate drift and preserve A/V alignment. +timestamped with a monotonic clock. The Windows receiver maps the source PTS +into one OBS monotonic session clock, preserving the audio/video offset instead +of replacing it with receiver-arrival time. Frames that arrive more than +`250 ms` late are dropped and logged so virtual-camera load cannot create an +unbounded backlog. ## Telemetry diff --git a/docs/protocol.md b/docs/protocol.md index 016e1b9..c45c810 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -43,11 +43,11 @@ transport stream. ### Video Payload -- Codec: `video/hevc` (preferred) or `video/avc` (fallback) +- Codec: `video/avc` by default; `video/hevc` only for an explicit hardware profile - Source: MediaCodec hardware encoder - Resolution: 1920x1080 default -- Frame rate: 60 fps default -- Bitrate: 50 Mbps default +- Frame rate: 30 fps default +- Bitrate: 12 Mbps default (bounded to 8-50 Mbps) - Keyframe interval: 1 second - No B-frame dependency in target encoder profile @@ -76,7 +76,7 @@ on port `51515`: **Beacon format:** ```text -OPENSTREAM/1 {"type":"dev.openstream.listener","version":1,"name":"OpenStream","instanceId":"...","sourceInstanceId":"...","slotId":"...","slotLabel":"CAM A","pairingUrl":"openstream://connect?...","host":"","listenerPort":9000,"latencyMs":120,"bitrateMbps":50,"busy":false} +OPENSTREAM/1 {"type":"dev.openstream.listener","version":1,"name":"OpenStream","instanceId":"...","sourceInstanceId":"...","slotId":"...","slotLabel":"CAM A","pairingUrl":"openstream://connect?...","host":"","listenerPort":9000,"latencyMs":120,"bitrateMbps":12,"busy":false} ``` | Field | Type | Description | @@ -102,7 +102,7 @@ The Android app advertises itself on the same multicast group: **Beacon format:** ```text -OPENSTREAM_PHONE/1 {"type":"dev.openstream.phone","version":1,"name":"","instanceId":"...","host":"","listenerPort":9000,"controlPort":9001,"latencyMs":120,"codec":"video/hevc","width":1920,"height":1080,"fps":60,"bitrateMbps":50,"busy":false,"reservedBy":""} +OPENSTREAM_PHONE/1 {"type":"dev.openstream.phone","version":1,"name":"","instanceId":"...","host":"","listenerPort":9000,"controlPort":9001,"latencyMs":120,"codec":"video/avc","width":1920,"height":1080,"fps":30,"bitrateMbps":12,"busy":false,"reservedBy":""} ``` | Field | Type | Description | @@ -157,7 +157,7 @@ Content-Type: application/json "sourceInstanceId": "openstream-...", "slotId": "slot-a", "slotLabel": "CAM A", - "bitrateMbps": 50 + "bitrateMbps": 12 } ``` @@ -272,11 +272,11 @@ Telemetry is separate from media. The Android app samples: ```json { "deviceName": "Google Pixel", - "codec": "video/hevc", + "codec": "video/avc", "width": 1920, "height": 1080, - "fps": 60, - "bitrate": 20000000, + "fps": 30, + "bitrate": 12000000, "batteryPercent": 87, "wifiRssi": -48, "temperatureCelsius": null, diff --git a/docs/testing.md b/docs/testing.md index fb100d7..c2eca63 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -28,7 +28,10 @@ implementation. - OBS receives video as one source. - OBS receives mono AAC audio at 48 kHz on the source's mixer channel. - The OBS source shows only the phone camera feed, never the Android screen. -- 1080p60 works on devices that advertise hardware support. +- Starting OBS Virtual Camera while the source is live does not create an + ever-growing delay; video remains within the configured low-latency window. +- A clap or tone remains aligned between the OBS video and mono audio paths + after Virtual Camera starts and after a reconnect. - Telemetry updates at least once per second. ## Android/OBS compatibility matrix @@ -80,3 +83,22 @@ Run 1080p30 for 30 minutes and log: - Bitrate If temperature exceeds the warning threshold, the app should recommend lowering bitrate or switching to 720p. + +## Performance measurement record + +The pre-fix defaults and the candidate defaults are measured from the checked-in +profiles and bounded buffers below. Runtime temperature and end-to-end latency +still require a physical phone/OBS run before release sign-off. + +| Item | Before | Candidate | Change | +|---|---:|---:|---:| +| Default video rate | 60 fps | 30 fps | 50% fewer encoded frames | +| Default video bitrate | 50 Mbps | 12 Mbps | 76% lower target bitrate | +| Audio capture buffer | 250 ms | 80 ms | 68% less capture backlog | +| App send path | Inline blocking callback | 768 KiB hard cap | Drops the session instead of growing latency | +| OBS stale-frame threshold | Arrival-time timestamps | 250 ms source-clock backlog | Old frames are dropped and logged | + +For the device run, record the 30-minute temperature delta, average measured +end-to-end latency, maximum audio/video offset, reconnect time, and dropped +frame count for the old and candidate builds under the same Wi-Fi and OBS +Virtual Camera workload. diff --git a/obs-plugin/CMakeLists.txt b/obs-plugin/CMakeLists.txt index 27dc516..377144d 100644 --- a/obs-plugin/CMakeLists.txt +++ b/obs-plugin/CMakeLists.txt @@ -1,5 +1,7 @@ cmake_minimum_required(VERSION 3.24) project(openstream_obs_plugin VERSION 1.0.1 LANGUAGES CXX) +include(CTest) +find_package(Threads REQUIRED) set(OPENSTREAM_VERSION "${PROJECT_VERSION}" CACHE STRING "Version embedded in the OBS plugin") # Releases are built against OBS Studio 32.2.1 x64. Its bundled FFmpeg ABI is @@ -22,6 +24,16 @@ add_library(openstream-obs MODULE src/openstream-source.cpp ) +if (BUILD_TESTING) + add_executable(openstream-obs-contract-tests + tests/test_contracts.cpp + src/async-control-client.cpp + ) + target_include_directories(openstream-obs-contract-tests PRIVATE src) + target_link_libraries(openstream-obs-contract-tests PRIVATE Threads::Threads) + add_test(NAME openstream-obs-contract-tests COMMAND openstream-obs-contract-tests) +endif() + find_package(Qt6 6.2 COMPONENTS Network Widgets REQUIRED) find_path(OBS_FRONTEND_INCLUDE_DIR NAMES obs-frontend-api.h HINTS ${OBS_ROOT} PATH_SUFFIXES frontend/api REQUIRED) diff --git a/obs-plugin/src/async-control-client.cpp b/obs-plugin/src/async-control-client.cpp index a541cc7..1f55127 100644 --- a/obs-plugin/src/async-control-client.cpp +++ b/obs-plugin/src/async-control-client.cpp @@ -4,14 +4,15 @@ AsyncControlClient::AsyncControlClient() : worker_(&AsyncControlClient::run, thi AsyncControlClient::~AsyncControlClient() { stop(); } -void AsyncControlClient::post(std::function command) { - if (!command) return; +bool AsyncControlClient::post(std::function command) { + if (!command) return false; { std::lock_guard lock(mutex_); - if (stopping_) return; + if (stopping_ || commands_.size() >= kQueueCapacity) return false; commands_.push(std::move(command)); } wake_.notify_one(); + return true; } void AsyncControlClient::stop() { @@ -28,15 +29,33 @@ void AsyncControlClient::stop() { if (worker_.joinable()) worker_.join(); } +bool AsyncControlClient::post_urgent(std::function command) { + if (!command) return false; + { + std::lock_guard lock(mutex_); + if (stopping_ || urgent_command_) return false; + urgent_command_ = std::move(command); + } + wake_.notify_one(); + return true; +} + void AsyncControlClient::run() { for (;;) { std::function command; { std::unique_lock lock(mutex_); - wake_.wait(lock, [this] { return stopping_ || !commands_.empty(); }); - if (stopping_) return; - command = std::move(commands_.front()); - commands_.pop(); + wake_.wait(lock, [this] { + return stopping_ || urgent_command_ || !commands_.empty(); + }); + if (urgent_command_) { + command = std::move(urgent_command_); + } else if (stopping_) { + return; + } else { + command = std::move(commands_.front()); + commands_.pop(); + } } command(); } diff --git a/obs-plugin/src/async-control-client.hpp b/obs-plugin/src/async-control-client.hpp index 0a511e2..1931daa 100644 --- a/obs-plugin/src/async-control-client.hpp +++ b/obs-plugin/src/async-control-client.hpp @@ -1,5 +1,6 @@ #pragma once +#include #include #include #include @@ -16,14 +17,23 @@ class AsyncControlClient { AsyncControlClient(const AsyncControlClient &) = delete; AsyncControlClient &operator=(const AsyncControlClient &) = delete; - void post(std::function command); + bool post(std::function command); + // Reservation release is lifecycle-critical. Keep one separate command so + // teardown can still release a phone when the transient control queue is full. + bool post_urgent(std::function command); void stop(); private: + // Camera controls are transient. Rejecting new work when this small queue + // is full prevents a disconnected phone from turning UI clicks into stale + // network requests. + static constexpr std::size_t kQueueCapacity = 16; + void run(); std::mutex mutex_; std::condition_variable wake_; std::queue> commands_; + std::function urgent_command_; bool stopping_ = false; std::thread worker_; }; diff --git a/obs-plugin/src/media-clock.hpp b/obs-plugin/src/media-clock.hpp new file mode 100644 index 0000000..1b5d940 --- /dev/null +++ b/obs-plugin/src/media-clock.hpp @@ -0,0 +1,38 @@ +#pragma once + +#include +#include +#include + +// Maps the phone's shared media clock into OBS's monotonic clock. The mapping +// preserves the audio/video offset carried by MPEG-TS instead of replacing it +// with the time at which a frame happened to arrive on the receiver thread. +// One instance belongs to one receiver session, so reconnects start cleanly. +class MediaClock { + public: + std::optional map(int64_t source_ns, uint64_t obs_now_ns) { + if (source_ns < 0) return std::nullopt; + if (!source_origin_ns_) { + source_origin_ns_ = source_ns; + obs_origin_ns_ = obs_now_ns; + } + + const int64_t delta = source_ns - *source_origin_ns_; + if (delta < 0) { + const uint64_t magnitude = + static_cast(-(delta + 1)) + 1u; + if (magnitude > obs_origin_ns_) return std::nullopt; + return obs_origin_ns_ - magnitude; + } + + const uint64_t offset = static_cast(delta); + if (offset > (std::numeric_limits::max)() - obs_origin_ns_) { + return std::nullopt; + } + return obs_origin_ns_ + offset; + } + + private: + std::optional source_origin_ns_; + uint64_t obs_origin_ns_ = 0; +}; diff --git a/obs-plugin/src/openstream-dock.cpp b/obs-plugin/src/openstream-dock.cpp index cdb6d64..93f48cb 100644 --- a/obs-plugin/src/openstream-dock.cpp +++ b/obs-plugin/src/openstream-dock.cpp @@ -228,7 +228,7 @@ OpenStreamDock *g_dock = nullptr; void openstream_register_dock() { if (g_dock) return; g_dock = new OpenStreamDock(); - // OBS 30+ owns the dock wrapper; the plugin retains/deletes the widget on unload. + // OBS 30+ owns the dock wrapper and the widget after registration. if (!obs_frontend_add_dock_by_id("openstream-camera-control", "OpenStream Camera Control", g_dock)) { delete g_dock; @@ -238,7 +238,8 @@ void openstream_register_dock() { void openstream_unregister_dock() { if (!g_dock) return; - obs_frontend_remove_dock("openstream-camera-control"); - delete g_dock; + // OBS owns the QDockWidget wrapper and destroys the child widget when the + // dock is removed. Clear our non-owning pointer before triggering teardown. g_dock = nullptr; + obs_frontend_remove_dock("openstream-camera-control"); } diff --git a/obs-plugin/src/openstream-source.cpp b/obs-plugin/src/openstream-source.cpp index 34c3073..da605d3 100644 --- a/obs-plugin/src/openstream-source.cpp +++ b/obs-plugin/src/openstream-source.cpp @@ -17,6 +17,7 @@ #include #include "async-control-client.hpp" +#include "media-clock.hpp" #include "openstream-control-api.hpp" #include @@ -90,6 +91,10 @@ using SwsContextPtr = std::unique_ptr; constexpr int kDiscoveryPort = 51515; constexpr int kDefaultListenerPort = 9000; +constexpr int kDefaultBitrateMbps = 12; +constexpr int kMinBitrateMbps = 8; +constexpr int kMaxBitrateMbps = 50; +constexpr uint64_t kMaximumMediaBacklogNs = 250'000'000; constexpr auto kReconnectReservationWindow = std::chrono::seconds(45); constexpr const char *kOpenStreamSourceName = "OpenStream V8"; constexpr const char *kDiscoveryMulticastAddress = "239.255.42.99"; @@ -442,7 +447,7 @@ struct PhoneDevice { int width = 1920; int height = 1080; int fps = 30; - int bitrate_mbps = 50; + int bitrate_mbps = kDefaultBitrateMbps; bool busy = false; std::string reserved_by; std::chrono::steady_clock::time_point last_seen = std::chrono::steady_clock::now(); @@ -613,7 +618,10 @@ class PhoneDiscoveryReceiver { device.width = std::clamp(json_int_value(json, "width").value_or(1920), 16, 8192); device.height = std::clamp(json_int_value(json, "height").value_or(1080), 16, 8192); device.fps = std::clamp(json_int_value(json, "fps").value_or(30), 1, 240); - device.bitrate_mbps = std::clamp(json_int_value(json, "bitrateMbps").value_or(12), 1, 200); + device.bitrate_mbps = std::clamp( + json_int_value(json, "bitrateMbps").value_or(kDefaultBitrateMbps), + kMinBitrateMbps, + kMaxBitrateMbps); device.busy = json_bool_value(json, "busy").value_or(false); device.reserved_by = json_string_value(json, "reservedBy").value_or(""); device.last_seen = std::chrono::steady_clock::now(); @@ -745,7 +753,7 @@ class DiscoveryAdvertiser { std::thread worker_; int listener_port_ = 9000; int latency_ms_ = 120; - int bitrate_mbps_ = 50; + int bitrate_mbps_ = kDefaultBitrateMbps; std::string source_name_ = kOpenStreamSourceName; std::string instance_id_; std::string slot_id_; @@ -768,7 +776,7 @@ struct OpenStreamSource { std::string selected_phone_id = PhoneDiscoveryReceiver::kAutoPhoneId; int listener_port = 0; int latency_ms = 120; - int bitrate_mbps = 50; + int bitrate_mbps = kDefaultBitrateMbps; bool listener_enabled = true; std::atomic listener_running = false; std::atomic phone_connected = false; @@ -788,6 +796,8 @@ struct OpenStreamSource { std::string active_selected_phone_id = PhoneDiscoveryReceiver::kAutoPhoneId; std::optional active_phone; uint64_t frames_output = 0; + uint64_t stale_video_frames = 0; + uint64_t stale_audio_frames = 0; double last_cam_zoom = 1.0; std::shared_ptr camera_controls = std::make_shared(); @@ -805,12 +815,17 @@ bool queue_control_command(OpenStreamSource *ctx, const std::string &path, const auto client = ctx->camera_controls; const std::string host = phone->host; const int port = phone->control_port; - client->post([host, port, path, body] { + const bool queued = client->post([host, port, path, body] { if (!send_control_command(host, port, path, body)) { blog(LOG_WARNING, "[OpenStream] Camera command %s failed", path.c_str()); } }); - return true; + if (!queued) { + blog(LOG_WARNING, + "[OpenStream] Camera command %s dropped: control queue is full or stopping", + path.c_str()); + } + return queued; } void set_slot_status(OpenStreamSource *ctx, std::string status) { @@ -976,6 +991,24 @@ void release_phone(OpenStreamSource *ctx, const PhoneDevice &phone) { send_control_command(phone.host, phone.control_port, "/release", body.str()); } +void queue_release_phone(OpenStreamSource *ctx, const PhoneDevice &phone) { + const auto client = ctx->camera_controls; + const std::string host = phone.host; + const int port = phone.control_port; + const std::string source_instance_id = ctx->instance_id; + const bool queued = client->post_urgent([host, port, source_instance_id] { + const std::string body = "{\"sourceInstanceId\":\"" + + json_escape(source_instance_id) + "\"}"; + if (!send_control_command(host, port, "/release", body)) { + blog(LOG_WARNING, "[OpenStream] Camera reservation release failed"); + } + }); + if (!queued) { + blog(LOG_WARNING, + "[OpenStream] Camera reservation release dropped: control queue is full or stopping"); + } +} + std::string av_error(int error) { char buffer[AV_ERROR_MAX_STRING_SIZE] = {}; av_strerror(error, buffer, sizeof(buffer)); @@ -1033,7 +1066,6 @@ const char *openstream_get_name(void *) { } void openstream_stop_worker(OpenStreamSource *ctx) { - const auto phone = control_phone(ctx); ctx->stop_requested = true; ctx->listener_running = false; ctx->phone_connected = false; @@ -1041,8 +1073,17 @@ void openstream_stop_worker(OpenStreamSource *ctx) { if (ctx->worker.joinable()) { ctx->worker.join(); } - if (phone.has_value()) { - release_phone(ctx, *phone); + std::optional reserved_phone; + { + std::lock_guard lock(ctx->settings_mutex); + // A live worker releases its reservation on its own stop path. Only queue + // a release for a reservation left behind by an interrupted/retry path. + reserved_phone = ctx->active_phone; + } + if (reserved_phone.has_value()) { + // Source stop/reconfiguration can be called by an OBS/Qt callback. Queue + // the release so that callback never performs control-socket I/O. + queue_release_phone(ctx, *reserved_phone); } set_active_phone(ctx, std::nullopt); set_slot_status(ctx, "Offline"); @@ -1157,10 +1198,25 @@ bool open_audio_decoder(AVFormatContext *format_ctx, return true; } +std::optional source_timestamp_ns(const AVFrame *frame, + const AVStream *stream) { + if (!frame || !stream) return std::nullopt; + const int64_t timestamp = frame->best_effort_timestamp != AV_NOPTS_VALUE + ? frame->best_effort_timestamp + : frame->pts; + if (timestamp == AV_NOPTS_VALUE) return std::nullopt; + return av_rescale_q(timestamp, stream->time_base, AVRational{1, 1'000'000'000}); +} + +bool stream_timestamp_is_stale(uint64_t timestamp_ns) { + const uint64_t now_ns = os_gettime_ns(); + return timestamp_ns < now_ns && now_ns - timestamp_ns > kMaximumMediaBacklogNs; +} + bool output_decoded_frame(OpenStreamSource *ctx, - AVStream *stream, AVCodecContext *decoder_ctx, AVFrame *decoded_frame, + uint64_t timestamp_ns, SwsContextPtr *sws_ctx, std::vector *bgra_buffer) { if (decoded_frame->width <= 0 || decoded_frame->height <= 0) { @@ -1182,7 +1238,7 @@ bool output_decoded_frame(OpenStreamSource *ctx, yuv_frame.format = can_output_i420 ? VIDEO_FORMAT_I420 : VIDEO_FORMAT_NV12; yuv_frame.width = static_cast(width); yuv_frame.height = static_cast(height); - yuv_frame.timestamp = os_gettime_ns(); + yuv_frame.timestamp = timestamp_ns; yuv_frame.range = obs_range_for_frame(decoded_frame, source_format); yuv_frame.trc = static_cast(obs_trc_for_frame(decoded_frame)); yuv_frame.flip = false; @@ -1261,13 +1317,11 @@ bool output_decoded_frame(OpenStreamSource *ctx, return false; } - (void)stream; - struct obs_source_frame obs_frame = {}; obs_frame.format = VIDEO_FORMAT_BGRA; obs_frame.width = static_cast(width); obs_frame.height = static_cast(height); - obs_frame.timestamp = os_gettime_ns(); + obs_frame.timestamp = timestamp_ns; obs_frame.data[0] = bgra_buffer->data(); obs_frame.linesize[0] = static_cast(linesize); obs_frame.flip = false; @@ -1342,6 +1396,10 @@ void decode_packets(OpenStreamSource *ctx, SwsContextPtr sws_ctx(nullptr); std::vector bgra_buffer; AVStream *video_stream = format_ctx->streams[video_stream_index]; + AVStream *audio_stream = audio_stream_index >= 0 + ? format_ctx->streams[audio_stream_index] + : nullptr; + MediaClock media_clock; uint64_t audio_frames_output = 0; const auto drain_video = [&]() -> int { @@ -1356,8 +1414,34 @@ void decode_packets(OpenStreamSource *ctx, av_error(result).c_str()); return result; } - output_decoded_frame( - ctx, video_stream, video_decoder_ctx, frame.get(), &sws_ctx, &bgra_buffer); + const auto source_ns = source_timestamp_ns(frame.get(), video_stream); + const auto timestamp_ns = source_ns + ? media_clock.map(*source_ns, os_gettime_ns()) + : std::nullopt; + if (!timestamp_ns) { + blog(LOG_WARNING, + "[OpenStream] Dropping video frame without a usable source timestamp (media gap surfaced)"); + av_frame_unref(frame.get()); + continue; + } + if (stream_timestamp_is_stale(*timestamp_ns)) { + const uint64_t dropped = ++ctx->stale_video_frames; + if (dropped == 1 || dropped % 60 == 0) { + blog(LOG_WARNING, + "[OpenStream] Dropped stale video frame(s): %" PRIu64 + " (receiver backlog exceeded %u ms; media gap surfaced)", + dropped, + static_cast(kMaximumMediaBacklogNs / 1'000'000)); + } + av_frame_unref(frame.get()); + continue; + } + output_decoded_frame(ctx, + video_decoder_ctx, + frame.get(), + *timestamp_ns, + &sws_ctx, + &bgra_buffer); av_frame_unref(frame.get()); } return AVERROR_EXIT; @@ -1389,10 +1473,32 @@ void decode_packets(OpenStreamSource *ctx, return AVERROR(EINVAL); } + const auto source_ns = source_timestamp_ns(audio_frame.get(), audio_stream); + const auto timestamp_ns = source_ns + ? media_clock.map(*source_ns, os_gettime_ns()) + : std::nullopt; + if (!timestamp_ns) { + blog(LOG_WARNING, + "[OpenStream] Dropping audio frame without a usable source timestamp (media gap surfaced)"); + av_frame_unref(audio_frame.get()); + continue; + } + if (stream_timestamp_is_stale(*timestamp_ns)) { + const uint64_t dropped = ++ctx->stale_audio_frames; + if (dropped == 1 || dropped % 100 == 0) { + blog(LOG_WARNING, + "[OpenStream] Dropped stale audio frame(s): %" PRIu64 + " (audio gap surfaced)", + dropped); + } + av_frame_unref(audio_frame.get()); + continue; + } + struct obs_source_audio obs_audio = {}; obs_audio.samples_per_sec = sample_rate; obs_audio.frames = static_cast(audio_frame->nb_samples); - obs_audio.timestamp = os_gettime_ns(); + obs_audio.timestamp = *timestamp_ns; obs_audio.format = obs_format; obs_audio.speakers = speakers; for (int i = 0; i < MAX_AV_PLANES && audio_frame->data[i]; i++) { @@ -1519,8 +1625,8 @@ void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::str AVDictionary *options = nullptr; av_dict_set(&options, "fflags", "nobuffer", 0); av_dict_set(&options, "flags", "low_delay", 0); - av_dict_set(&options, "probesize", "1048576", 0); - av_dict_set(&options, "analyzeduration", "1000000", 0); + av_dict_set(&options, "probesize", "262144", 0); + av_dict_set(&options, "analyzeduration", "250000", 0); blog(LOG_INFO, "[OpenStream] Opening Android stream at %s", @@ -1556,6 +1662,8 @@ void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::str ctx->phone_connected = true; set_slot_status(ctx, "Live"); ctx->frames_output = 0; + ctx->stale_video_frames = 0; + ctx->stale_audio_frames = 0; int video_stream_index = -1; CodecContextPtr video_decoder_ctx; @@ -1614,7 +1722,7 @@ void openstream_start_worker(OpenStreamSource *ctx) { std::string srt_url; int listener_port = kDefaultListenerPort; int latency_ms = 120; - int bitrate_mbps = 50; + int bitrate_mbps = kDefaultBitrateMbps; std::string source_name; std::string instance_id; std::string slot_id; @@ -1690,8 +1798,12 @@ void openstream_update(void *data, obs_data_t *settings) { obs_data_set_int(settings, "listener_port", requested_port); } ctx->listener_port = requested_port; - ctx->latency_ms = static_cast(obs_data_get_int(settings, "latency_ms")); - ctx->bitrate_mbps = static_cast(obs_data_get_int(settings, "bitrate_mbps")); + ctx->latency_ms = std::clamp( + static_cast(obs_data_get_int(settings, "latency_ms")), 80, 200); + ctx->bitrate_mbps = std::clamp( + static_cast(obs_data_get_int(settings, "bitrate_mbps")), + kMinBitrateMbps, + kMaxBitrateMbps); const char *slot_id = obs_data_get_string(settings, "slot_id"); if (slot_id && slot_id[0] != '\0') { ctx->slot_id = slot_id; @@ -1815,7 +1927,7 @@ void openstream_defaults(obs_data_t *settings) { obs_data_set_default_bool(settings, "show_advanced", false); obs_data_set_default_int(settings, "listener_port", kDefaultListenerPort); obs_data_set_default_int(settings, "latency_ms", 120); - obs_data_set_default_int(settings, "bitrate_mbps", 50); + obs_data_set_default_int(settings, "bitrate_mbps", kDefaultBitrateMbps); obs_data_set_default_double(settings, "cam_zoom", 1.0); } @@ -1964,7 +2076,12 @@ obs_properties_t *openstream_properties(void *data) { obs_property_int_set_suffix(latency, " ms"); obs_property_set_long_description(latency, "Higher values are more stable on Wi-Fi; lower values reduce delay."); obs_property_t *bitrate = - obs_properties_add_int_slider(advanced_group, "bitrate_mbps", "Expected bitrate (Mbps)", 8, 120, 1); + obs_properties_add_int_slider(advanced_group, + "bitrate_mbps", + "Expected bitrate (Mbps)", + kMinBitrateMbps, + kMaxBitrateMbps, + 1); obs_property_int_set_suffix(bitrate, " Mbps"); obs_property_set_long_description(bitrate, "Used in discovery so the phone can tune stream quality for this slot."); obs_properties_add_group(props, "show_advanced", "3. Network & Pairing (Advanced)", OBS_GROUP_CHECKABLE, advanced_group); diff --git a/obs-plugin/tests/test_contracts.cpp b/obs-plugin/tests/test_contracts.cpp new file mode 100644 index 0000000..25a89fa --- /dev/null +++ b/obs-plugin/tests/test_contracts.cpp @@ -0,0 +1,44 @@ +#include "../src/async-control-client.hpp" +#include "../src/media-clock.hpp" + +#include +#include + +namespace { +void check(bool condition) { + if (!condition) std::abort(); +} +} // namespace + +int main() { + { + MediaClock clock; + const uint64_t origin = 10'000'000'000ULL; + check(clock.map(1'000'000, origin).value() == origin); + check(clock.map(1'033'333, origin + 33'333).value() == origin + 33'333); + check(!clock.map(-1, origin).has_value()); + } + + { + AsyncControlClient client; + std::promise started; + std::promise release; + const auto release_signal = release.get_future().share(); + check(client.post([&] { + started.set_value(); + release_signal.wait(); + })); + started.get_future().wait(); + + for (int i = 0; i < 16; ++i) { + check(client.post([] {})); + } + check(!client.post([] {})); + release.set_value(); + bool urgent_ran = false; + check(client.post_urgent([&] { urgent_ran = true; })); + client.stop(); + check(urgent_ran); + check(!client.post([] {})); + } +} diff --git a/tests/test_repo_contract.py b/tests/test_repo_contract.py index 1aa6741..a953534 100644 --- a/tests/test_repo_contract.py +++ b/tests/test_repo_contract.py @@ -11,7 +11,7 @@ def read(path: str) -> str: def test_architecture_documents_practical_v1_transport() -> None: architecture = read("docs/architecture.md") - assert "MediaCodec hardware HEVC/H.264 video encode" in architecture + assert "MediaCodec hardware AVC/H.264 video encode" in architecture assert "MediaCodec AAC audio encode" in architecture assert "SRT caller" in architecture assert "UDP discovery" in architecture @@ -31,6 +31,9 @@ def test_android_project_declares_camera_media_codec_srt_discovery_boundaries() assert "startPreviewIfAllowed" in app assert "startPhoneServerIfAllowed" in app assert "MediaCodecAudioEncoder" in app + camera = read("android/app/src/main/java/dev/openstream/app/camera/Camera2Controller.kt") + assert "CONTROL_AE_TARGET_FPS_RANGE" in camera + assert "targetFps" in camera stream_config = read("android/app/src/main/java/dev/openstream/app/stream/StreamConfig.kt") assert "Default1080p30" in stream_config assert "codecPreference = CodecPreference.ForceAvc" in stream_config diff --git a/website/src/main.jsx b/website/src/main.jsx index a7a5b77..2b8bf9e 100644 --- a/website/src/main.jsx +++ b/website/src/main.jsx @@ -54,7 +54,7 @@ function ObsPreview() {
OBS Studio — □ ×
File   Edit   View   Docks   Profile   Scene Collection
-
OPENSTREAM / CAM ALIVE · 60 FPS
+
OPENSTREAM / CAM ALIVE · 30 FPS
Scenes
Camera
Overlay
Sources
OpenStream
Audio
Audio Mixer
▮▮▮▮▮▯▯
OpenStream AAC
OpenStream Camera Control ● Connected
RearFrontTorchIdentify
Zoom 1.8×
@@ -147,7 +147,7 @@ function Setup() { function Pipeline() { return (
-

Local network.
Explicit pipes.

Your media stays on the LAN. Camera2 and MediaCodec handle capture, MPEG-TS carries HEVC/H.264 plus AAC, and SRT delivers it to the native OBS source.

+

Local network.
Explicit pipes.

Your media stays on the LAN. Camera2 and MediaCodec handle capture, MPEG-TS carries hardware AVC/H.264 plus AAC, and SRT delivers it to the native OBS source.

ANDROID CAMERACamera2 + MediaCodecSRT STREAMMPEG-TS · port 9000OBS SOURCEFFmpeg decode + mixer
Media
SRT :9000
Discovery
UDP :51515
Control
HTTP :9001
Default bitrate
16 Mbps