diff --git a/README.md b/README.md index 7604a0a..86f3875 100644 --- a/README.md +++ b/README.md @@ -12,7 +12,7 @@ [![OBS](https://img.shields.io/badge/OBS-Studio%20Plugin-purple?style=flat-square&labelColor=1a1a2e)](https://obsproject.com) [![Website](https://img.shields.io/badge/website-openstream.pages.dev-00D4AA?style=flat-square&labelColor=1a1a2e)](https://openstream.pages.dev) -**Open-source** | **Low-latency** | **Hardware-accelerated** | **Local Wi-Fi** +**Open-source** | **Low-latency** | **Hardware-accelerated** | **USB or Local Wi-Fi** --- @@ -35,7 +35,7 @@ Need the non-technical walkthrough with screenshots? Start with [`docs/set-up.md ## What is OpenStream? -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. +OpenStream V1.0.1 sends your Android phone camera directly into OBS Studio over USB tethering or 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 diff --git a/android/app/src/main/cpp/openstream_srt.cpp b/android/app/src/main/cpp/openstream_srt.cpp index 010f3ee..e732caa 100644 --- a/android/app/src/main/cpp/openstream_srt.cpp +++ b/android/app/src/main/cpp/openstream_srt.cpp @@ -27,6 +27,8 @@ namespace { constexpr const char *kTag = "OpenStreamSRT"; +constexpr int kMinSrtLatencyMs = 30; +constexpr int kMaxSrtLatencyMs = 200; constexpr int kMediaCodecBufferFlagKeyFrame = 1; constexpr int kMediaCodecBufferFlagCodecConfig = 2; constexpr int kAudioSampleRate = 48000; @@ -467,7 +469,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) { 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..483f09f 100644 --- a/android/app/src/main/java/dev/openstream/app/MainActivity.kt +++ b/android/app/src/main/java/dev/openstream/app/MainActivity.kt @@ -35,6 +35,7 @@ import dev.openstream.app.encoder.MediaCodecAudioEncoder import dev.openstream.app.encoder.MediaCodecVideoEncoder import dev.openstream.app.stream.ConnectionTarget import dev.openstream.app.stream.StreamConfig +import dev.openstream.app.stream.TransportMode import dev.openstream.app.stream.SrtStreamClient import dev.openstream.app.telemetry.TelemetrySampler @@ -92,6 +93,7 @@ class MainActivity : Activity() { private var zoomHideRunnable: Runnable? = null private var liveDotAnimator: ObjectAnimator? = null private var currentPort: Int = ConnectionTarget.DEFAULT_PORT + private var currentTransportMode: TransportMode = TransportMode.Wifi private var releaseReservationRunnable: Runnable? = null private var lensRestartRunnable: Runnable? = null private var pendingConnectAfterSettings = false @@ -120,12 +122,14 @@ class MainActivity : Activity() { .getInt(SettingsActivity.KEY_LISTENING_PORT, ConnectionTarget.DEFAULT_PORT) .takeIf { it in 1024..65535 } ?: ConnectionTarget.DEFAULT_PORT + currentTransportMode = selectedTransportMode() streamClient = SrtStreamClient() telemetry = TelemetrySampler(this) phoneAdvertiser = PhoneDiscoveryAdvertiser( context = this, - config = streamConfig, + config = activeStreamConfig(), + transport = currentTransportMode, port = currentPort, busyProvider = { phoneConnected || reservedBy != null }, reservedByProvider = { reservedBy }, @@ -202,8 +206,7 @@ class MainActivity : Activity() { override fun onStart() { super.onStart() activityStarted = true - phoneAdvertiser.start() - obsDiscoveryClient.start() + startDiscovery() controlServer.start() startPreviewIfAllowed() startPhoneServerIfAllowed() @@ -211,11 +214,12 @@ class MainActivity : Activity() { override fun onResume() { super.onResume() - // Reload listening port from settings if it changed + // The listener must be restarted to apply a selected SRT latency. val settingsPrefs = getSharedPreferences(SettingsActivity.PREFS_NAME, MODE_PRIVATE) val savedPort = settingsPrefs.getInt(SettingsActivity.KEY_LISTENING_PORT, currentPort) - if (savedPort != currentPort && savedPort in 1024..65535) { - changePort(savedPort) + val savedTransport = selectedTransportMode(settingsPrefs) + if ((savedPort != currentPort && savedPort in 1024..65535) || savedTransport != currentTransportMode) { + changeListenerSettings(savedPort, savedTransport) } if (pendingConnectAfterSettings) { pendingConnectAfterSettings = false @@ -228,8 +232,7 @@ class MainActivity : Activity() { cancelLensRestart() camera.stop() stopPhoneServer(clearReservation = false, updateStatus = false) - obsDiscoveryClient.stop() - phoneAdvertiser.stop() + stopDiscovery() controlServer.stop() stopLiveDotAnimation() super.onStop() @@ -668,13 +671,13 @@ class MainActivity : Activity() { activeTargetName = null statusText.text = getString(R.string.status_ready) statusText.setTextColor(getColor(R.color.os_text_primary)) - statusDetail.text = getString(R.string.status_waiting) + statusDetail.text = listenerStatusDetail() btnStop.visibility = View.GONE val thread = Thread({ try { while (isListenerActive(generation)) { - val listenUrl = "srt://0.0.0.0:${currentPort}?mode=listener&latency=${streamConfig.latencyMs}" + val listenUrl = "srt://0.0.0.0:${currentPort}?mode=listener&latency=${activeStreamConfig().latencyMs}" val listenResult = runCatching { streamClient.listen( url = listenUrl, @@ -729,7 +732,7 @@ class MainActivity : Activity() { hideLiveState() statusText.text = getString(R.string.status_ready) statusDetail.text = reservedSlotLabel?.let { "Holding $it for reconnect" } - ?: getString(R.string.status_waiting) + ?: listenerStatusDetail() } } } @@ -952,12 +955,14 @@ class MainActivity : Activity() { val host = settingsPrefs.getString(SettingsActivity.KEY_OBS_HOST, ConnectionTarget.DEFAULT_HOST) ?.trim().orEmpty().ifBlank { ConnectionTarget.DEFAULT_HOST } val port = settingsPrefs.getInt(SettingsActivity.KEY_OBS_PORT, ConnectionTarget.DEFAULT_PORT) - val latencyMs = settingsPrefs.getInt(SettingsActivity.KEY_LATENCY, ConnectionTarget.DEFAULT_LATENCY_MS) + val latencyMs = selectedTransportMode(settingsPrefs).latencyMs( + settingsPrefs.getInt(SettingsActivity.KEY_LATENCY, ConnectionTarget.DEFAULT_LATENCY_MS), + ) return ConnectionTarget( name = ConnectionTarget.DEFAULT_NAME, host = host, port = port.coerceIn(1, 65535), - latencyMs = latencyMs.coerceIn(80, 200), + latencyMs = latencyMs, ) } @@ -996,25 +1001,50 @@ class MainActivity : Activity() { // ─────────────────────────── Port selector ─────────────────────────── - private fun changePort(newPort: Int) { + private fun changeListenerSettings(newPort: Int, newTransportMode: TransportMode) { val clamped = newPort.coerceIn(1024, 65535) - if (clamped == currentPort) return + if (clamped == currentPort && newTransportMode == currentTransportMode) return currentPort = clamped - // Restart the phone server on the new port + currentTransportMode = newTransportMode stopPhoneServer(clearReservation = false) - // Re-create the advertiser with new port - phoneAdvertiser.stop() + stopDiscovery() phoneAdvertiser = PhoneDiscoveryAdvertiser( context = this, - config = streamConfig, + config = activeStreamConfig(), + transport = currentTransportMode, port = currentPort, busyProvider = { phoneConnected || reservedBy != null }, reservedByProvider = { reservedBy }, ) - phoneAdvertiser.start() + startDiscovery() startPhoneServerIfAllowed() } + private fun activeStreamConfig(): StreamConfig = streamConfig.copy( + latencyMs = currentTransportMode.latencyMs(streamConfig.latencyMs), + ) + + private fun selectedTransportMode( + prefs: android.content.SharedPreferences = getSharedPreferences(SettingsActivity.PREFS_NAME, MODE_PRIVATE), + ): TransportMode = TransportMode.fromPreference(prefs.getString(SettingsActivity.KEY_TRANSPORT, null)) + + private fun startDiscovery() { + phoneAdvertiser.start() + if (currentTransportMode == TransportMode.Wifi) { + obsDiscoveryClient.start() + } + } + + private fun stopDiscovery() { + obsDiscoveryClient.stop() + phoneAdvertiser.stop() + } + + private fun listenerStatusDetail(): String = when (currentTransportMode) { + TransportMode.Wifi -> getString(R.string.status_waiting) + TransportMode.UsbTether -> "Waiting for USB tethered network on SRT port $currentPort" + } + // ─────────────────────────── Preview aspect ratio fix ─────────────────────────── private fun adjustPreviewAspectRatio(surfaceWidth: Int, surfaceHeight: Int) { diff --git a/android/app/src/main/java/dev/openstream/app/SettingsActivity.kt b/android/app/src/main/java/dev/openstream/app/SettingsActivity.kt index 4062b8e..7ddac83 100644 --- a/android/app/src/main/java/dev/openstream/app/SettingsActivity.kt +++ b/android/app/src/main/java/dev/openstream/app/SettingsActivity.kt @@ -6,9 +6,11 @@ import android.os.Build import android.os.Bundle import android.view.WindowInsets import android.widget.EditText +import android.widget.RadioButton import android.widget.TextView import android.widget.Toast import dev.openstream.app.stream.ConnectionTarget +import dev.openstream.app.stream.TransportMode class SettingsActivity : Activity() { @@ -16,6 +18,8 @@ class SettingsActivity : Activity() { private lateinit var inputObsPort: EditText private lateinit var inputLatency: EditText private lateinit var inputListeningPort: EditText + private lateinit var transportWifi: RadioButton + private lateinit var transportUsb: RadioButton private lateinit var btnSave: TextView private lateinit var btnSaveAndConnect: TextView private lateinit var btnBack: TextView @@ -29,6 +33,8 @@ class SettingsActivity : Activity() { inputObsPort = findViewById(R.id.settingsObsPort) inputLatency = findViewById(R.id.settingsLatency) inputListeningPort = findViewById(R.id.settingsListeningPort) + transportWifi = findViewById(R.id.settingsTransportWifi) + transportUsb = findViewById(R.id.settingsTransportUsb) btnSave = findViewById(R.id.btnSaveSettings) btnSaveAndConnect = findViewById(R.id.btnSaveAndConnect) btnBack = findViewById(R.id.btnBackSettings) @@ -51,6 +57,10 @@ class SettingsActivity : Activity() { if (latency != ConnectionTarget.DEFAULT_LATENCY_MS) inputLatency.setText(latency.toString()) val listenPort = prefs.getInt(KEY_LISTENING_PORT, ConnectionTarget.DEFAULT_PORT) inputListeningPort.setText(listenPort.toString()) + when (TransportMode.fromPreference(prefs.getString(KEY_TRANSPORT, null))) { + TransportMode.Wifi -> transportWifi.isChecked = true + TransportMode.UsbTether -> transportUsb.isChecked = true + } } private fun saveSettings(connectAfterSave: Boolean) { @@ -67,10 +77,11 @@ class SettingsActivity : Activity() { validRange = 1..65535, label = "OBS port", ) ?: return + val mode = if (transportUsb.isChecked) TransportMode.UsbTether else TransportMode.Wifi val latency = validatedNumber( input = inputLatency, defaultValue = ConnectionTarget.DEFAULT_LATENCY_MS, - validRange = 80..200, + validRange = TransportMode.WIFI_MIN_LATENCY_MS..TransportMode.WIFI_MAX_LATENCY_MS, label = "Latency", ) ?: return val listenPort = validatedNumber( @@ -85,6 +96,7 @@ class SettingsActivity : Activity() { .putInt(KEY_OBS_PORT, port) .putInt(KEY_LATENCY, latency) .putInt(KEY_LISTENING_PORT, listenPort) + .putString(KEY_TRANSPORT, mode.preferenceValue) .apply() Toast.makeText(this, "Settings saved", Toast.LENGTH_SHORT).show() @@ -138,6 +150,7 @@ class SettingsActivity : Activity() { const val KEY_OBS_PORT = "obs_port" const val KEY_LATENCY = "latency_ms" const val KEY_LISTENING_PORT = "listening_port" + const val KEY_TRANSPORT = "transport" const val EXTRA_CONNECT_AFTER_SAVE = "connect_after_save" } } 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..ec880b8 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 @@ -14,6 +14,7 @@ import java.net.NetworkInterface import java.net.SocketTimeoutException import java.nio.charset.StandardCharsets import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicLong class ObsDiscoveryClient( private val context: Context, @@ -21,6 +22,8 @@ class ObsDiscoveryClient( private val nowMs: () -> Long = { System.currentTimeMillis() }, ) { private val running = AtomicBoolean(false) + private val generation = AtomicLong() + private val socketLock = Any() private val mainHandler = Handler(Looper.getMainLooper()) private val devices = linkedMapOf() private var socket: MulticastSocket? = null @@ -29,8 +32,9 @@ class ObsDiscoveryClient( fun start() { if (!running.compareAndSet(false, true)) return + val runGeneration = generation.incrementAndGet() acquireMulticastLock() - worker = Thread(::receiveLoop, "OpenStreamDiscovery").apply { + worker = Thread({ receiveLoop(runGeneration) }, "OpenStreamDiscovery").apply { isDaemon = true start() } @@ -38,8 +42,11 @@ class ObsDiscoveryClient( fun stop() { if (!running.compareAndSet(true, false)) return - socket?.close() - socket = null + generation.incrementAndGet() + synchronized(socketLock) { + socket?.close() + socket = null + } worker = null synchronized(devices) { devices.clear() @@ -65,17 +72,23 @@ class ObsDiscoveryClient( multicastLock = null } - private fun receiveLoop() { + private fun receiveLoop(runGeneration: Long) { val udp = MulticastSocket(null).apply { reuseAddress = true soTimeout = 500 bind(InetSocketAddress(DISCOVERY_PORT)) } - socket = udp + synchronized(socketLock) { + if (!running.get() || generation.get() != runGeneration) { + udp.close() + return + } + socket = udp + } joinDiscoveryMulticast(udp) val buffer = ByteArray(4096) - while (running.get()) { + while (running.get() && generation.get() == runGeneration) { try { val packet = DatagramPacket(buffer, buffer.size) udp.receive(packet) @@ -92,7 +105,7 @@ class ObsDiscoveryClient( publishDevices() } } catch (_: Exception) { - if (running.get()) { + if (running.get() && generation.get() == runGeneration) { if (pruneExpired()) { publishDevices() } @@ -100,6 +113,9 @@ class ObsDiscoveryClient( } } udp.close() + synchronized(socketLock) { + if (socket === udp) socket = null + } } private fun joinDiscoveryMulticast(socket: MulticastSocket) { 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..c8733a7 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 @@ -5,6 +5,7 @@ import android.net.wifi.WifiManager import android.os.Build import dev.openstream.app.encoder.advertisedMimeType import dev.openstream.app.stream.StreamConfig +import dev.openstream.app.stream.TransportMode import org.json.JSONObject import java.net.DatagramPacket import java.net.DatagramSocket @@ -16,6 +17,7 @@ import java.util.concurrent.atomic.AtomicBoolean class PhoneDiscoveryAdvertiser( private val context: Context, private val config: StreamConfig, + private val transport: TransportMode, private val port: Int, private val busyProvider: () -> Boolean, private val reservedByProvider: () -> String? = { null }, @@ -63,7 +65,8 @@ class PhoneDiscoveryAdvertiser( .put("version", 1) .put("name", "${Build.MANUFACTURER} ${Build.MODEL}".trim()) .put("instanceId", instanceId) - .put("host", localWifiAddress().orEmpty()) + .put("host", localAddress().orEmpty()) + .put("transport", transport.preferenceValue) .put("listenerPort", port) .put("latencyMs", config.latencyMs) .put("bitrateMbps", config.bitrateMbps) @@ -89,6 +92,26 @@ class PhoneDiscoveryAdvertiser( ).joinToString(".") } + private fun localAddress(): String? = when (transport) { + TransportMode.Wifi -> localWifiAddress() + TransportMode.UsbTether -> localUsbTetherAddress() + } + + private fun localUsbTetherAddress(): String? = runCatching { + java.net.NetworkInterface.getNetworkInterfaces().asSequence() + .filter { it.isUp && !it.isLoopback && isUsbTetherInterface(it.name) } + .flatMap { it.inetAddresses.asSequence() } + .filterIsInstance() + .firstOrNull { !it.isLoopbackAddress && it.isSiteLocalAddress } + ?.hostAddress + }.getOrNull() + + private fun isUsbTetherInterface(name: String): Boolean { + val normalized = name.lowercase() + return normalized.startsWith("rndis") || normalized.startsWith("usb") || + normalized.startsWith("ncm") || normalized.startsWith("enx") + } + companion object { const val DISCOVERY_PORT = 51515 const val DISCOVERY_MULTICAST_ADDRESS = "239.255.42.99" diff --git a/android/app/src/main/java/dev/openstream/app/stream/TransportMode.kt b/android/app/src/main/java/dev/openstream/app/stream/TransportMode.kt new file mode 100644 index 0000000..4f7555a --- /dev/null +++ b/android/app/src/main/java/dev/openstream/app/stream/TransportMode.kt @@ -0,0 +1,21 @@ +package dev.openstream.app.stream + +enum class TransportMode(val preferenceValue: String, val displayName: String) { + Wifi("wifi", "Wi-Fi"), + UsbTether("usb", "USB tethering"); + + fun latencyMs(wifiLatencyMs: Int): Int = when (this) { + Wifi -> wifiLatencyMs.coerceIn(WIFI_MIN_LATENCY_MS, WIFI_MAX_LATENCY_MS) + UsbTether -> USB_LATENCY_MS + } + + companion object { + const val USB_LATENCY_MS = 30 + const val WIFI_MIN_LATENCY_MS = 80 + const val WIFI_MAX_LATENCY_MS = 200 + + fun fromPreference(value: String?): TransportMode = + entries.firstOrNull { it.preferenceValue == value } + ?: if (value == "usb_adb") UsbTether else Wifi + } +} diff --git a/android/app/src/main/res/layout/activity_settings.xml b/android/app/src/main/res/layout/activity_settings.xml index 692e02e..07d167f 100644 --- a/android/app/src/main/res/layout/activity_settings.xml +++ b/android/app/src/main/res/layout/activity_settings.xml @@ -22,6 +22,44 @@ android:fontFamily="sans-serif-medium" android:layout_marginBottom="40dp" /> + + + + + + + + + + + OBS PORT LATENCY (ms) LISTENING PORT + TRANSPORT + Wi-Fi (80–200 ms SRT latency) + USB tethering / RNDIS (30 ms SRT latency) + For USB mode, enable USB tethering on the phone, then connect the OBS PC to the phone’s tethered network. ADB port forwarding cannot carry SRT. Screen dimmed. Tap to restore controls. Save & Connect Save Settings diff --git a/android/app/src/test/java/dev/openstream/app/SettingsValidatorTest.kt b/android/app/src/test/java/dev/openstream/app/SettingsValidatorTest.kt index 8d2ca75..f0bb30f 100644 --- a/android/app/src/test/java/dev/openstream/app/SettingsValidatorTest.kt +++ b/android/app/src/test/java/dev/openstream/app/SettingsValidatorTest.kt @@ -5,6 +5,7 @@ import org.junit.Assert.assertFalse import org.junit.Assert.assertNull import org.junit.Assert.assertTrue import org.junit.Test +import dev.openstream.app.stream.TransportMode class SettingsValidatorTest { @Test @@ -30,4 +31,13 @@ class SettingsValidatorTest { assertNull(SettingsValidator.parseNumber("79", 120, 80..200)) assertNull(SettingsValidator.parseNumber("invalid", 120, 80..200)) } + + @Test + fun usbTransportUsesFixedLowLatencyWhileWifiRemainsBounded() { + assertEquals(30, TransportMode.UsbTether.latencyMs(120)) + assertEquals(80, TransportMode.Wifi.latencyMs(20)) + assertEquals(200, TransportMode.Wifi.latencyMs(500)) + assertEquals(TransportMode.UsbTether, TransportMode.fromPreference("usb_adb")) + assertEquals(TransportMode.Wifi, TransportMode.fromPreference("unknown")) + } } diff --git a/docs/architecture.md b/docs/architecture.md index 4db8372..d8c0184 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -13,6 +13,11 @@ Android Camera2 -> OBS final stream/record encode ``` +The same encoded MPEG-TS/SRT session can travel over Android USB tethering +(RNDIS) instead of Wi-Fi. USB mode does not add another encoder, re-encode +video, or depend on ADB; Windows connects to the phone's tether-interface +gateway with a 30 ms SRT latency profile. + 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. @@ -63,7 +68,8 @@ transport can be tested without blocking the reliable product path. - Default SRT port: `9000`. - Default control port: `9001`. - Default latency: `120 ms`. -- Valid latency tuning range: `80-200 ms`. +- Wi-Fi latency tuning range: `80-200 ms`. +- USB tether latency: `30 ms`. ## Encoding Defaults @@ -91,8 +97,10 @@ 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 shared +MPEG-TS source timeline once into the OBS monotonic clock and preserves every +audio/video offset; frames without usable source timestamps are dropped and +logged instead of being stamped with arrival time. ## Telemetry diff --git a/docs/protocol.md b/docs/protocol.md index 016e1b9..6d59a71 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -8,7 +8,7 @@ version numbers and protocol versions are independent: `OPENSTREAM/1` and OpenStream uses three communication channels between the Android phone and OBS: -1. **Media Stream** - SRT/MPEG-TS for video + audio (phone -> OBS) +1. **Media Stream** - SRT/MPEG-TS for video + audio (phone -> OBS), over Wi-Fi or USB tethering 2. **Discovery** - UDP multicast/broadcast beacons (bidirectional) 3. **Control** - HTTP POST for remote camera commands (OBS -> phone) @@ -32,8 +32,18 @@ OBS listener: srt://0.0.0.0:9000?mode=listener&latency=120 ``` -The supported latency range is `80-200 ms`; both pairing links and the OBS -source clamp values to that range. +The supported Wi-Fi latency range is `80-200 ms`; pairing links clamp values to +that range. The explicit USB tether profile uses `30 ms`. + +### USB tether transport + +USB mode keeps the same codecs, MPEG-TS container, source PTS, and SRT session. +Android exposes its listener through the operating system's USB tether/RNDIS +interface; OBS connects to that adapter's IPv4 gateway using a `30 ms` SRT +latency. This is not ADB forwarding (ADB forwards TCP, while SRT uses UDP). +Automatic gateway detection is limited to Windows adapters whose name or +description identifies USB, RNDIS, or mobile tethering; the OBS source accepts +an explicit IPv4 phone address when a vendor driver uses another label. ### Container Format diff --git a/docs/testing.md b/docs/testing.md index fb100d7..44e5a87 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -67,6 +67,22 @@ Expected behavior: ## Developer receiver `tools/openstream_receiver.py` remains available for FFmpeg/SRT smoke tests + +## USB tether acceptance + +Before release, capture the same 1080p60 scene over Wi-Fi (120 ms) and USB +tethering (30 ms), then record p50/p95 glass-to-glass latency and maximum A/V +skew. Run USB for at least 30 minutes and unplug/replug the cable five times. + +- The OBS log must identify the USB tether URL and every disconnect/reconnect. +- Audio and video must retain MPEG-TS source timestamps; no arrival-time stamp is allowed. +- Maximum accepted A/V skew is 20 ms during steady state. +- A cable failure must affect only its owning source slot. +- Memory must remain bounded; camera-control backlog is capped at 16 commands + and rejects newest commands when full. + +Record before/after hardware measurements in the pull request. This repository +does not substitute synthetic timing numbers for a physical phone/cable run. without OBS. It is not part of the normal user workflow. ## Thermal tests diff --git a/obs-plugin/CMakeLists.txt b/obs-plugin/CMakeLists.txt index 27dc516..36c78bd 100644 --- a/obs-plugin/CMakeLists.txt +++ b/obs-plugin/CMakeLists.txt @@ -11,6 +11,13 @@ set(OPENSTREAM_FFMPEG_ABI "avformat-62;avcodec-62;avutil-60;swscale-9") set(CMAKE_CXX_STANDARD 20) set(CMAKE_CXX_STANDARD_REQUIRED ON) +include(CTest) + +if (BUILD_TESTING) + add_executable(openstream-media-clock-test ../tests/media_clock_test.cpp) + target_include_directories(openstream-media-clock-test PRIVATE src) + add_test(NAME openstream-media-clock COMMAND openstream-media-clock-test) +endif() if (NOT DEFINED OBS_ROOT) message(FATAL_ERROR "Set OBS_ROOT to the OBS Studio install, source checkout, build tree, or development SDK root.") diff --git a/obs-plugin/src/async-control-client.cpp b/obs-plugin/src/async-control-client.cpp index a541cc7..ff27bb2 100644 --- a/obs-plugin/src/async-control-client.cpp +++ b/obs-plugin/src/async-control-client.cpp @@ -4,14 +4,17 @@ 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; + // Capacity is intentionally small: UI controls are transient and must not + // accumulate behind a disconnected phone. Reject newest preserves order. + if (stopping_ || commands_.size() >= kQueueCapacity) return false; commands_.push(std::move(command)); } wake_.notify_one(); + return true; } void AsyncControlClient::stop() { diff --git a/obs-plugin/src/async-control-client.hpp b/obs-plugin/src/async-control-client.hpp index 0a511e2..28cb403 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 @@ -11,12 +12,14 @@ // wait for its network timeouts. class AsyncControlClient { public: + static constexpr size_t kQueueCapacity = 16; + AsyncControlClient(); ~AsyncControlClient(); AsyncControlClient(const AsyncControlClient &) = delete; AsyncControlClient &operator=(const AsyncControlClient &) = delete; - void post(std::function command); + bool post(std::function command); void stop(); private: diff --git a/obs-plugin/src/media-clock.hpp b/obs-plugin/src/media-clock.hpp new file mode 100644 index 0000000..c3e1dc1 --- /dev/null +++ b/obs-plugin/src/media-clock.hpp @@ -0,0 +1,29 @@ +#pragma once + +#include +#include + +// Maps the phone's shared MPEG-TS clock into OBS's monotonic clock while +// retaining the original audio/video offset. One instance belongs to one +// receiver session, so reconnects cannot leak timing state across cameras. +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 && static_cast(-delta) > obs_origin_ns_) { + return std::nullopt; + } + return delta < 0 ? obs_origin_ns_ - static_cast(-delta) + : obs_origin_ns_ + static_cast(delta); + } + + private: + std::optional source_origin_ns_; + uint64_t obs_origin_ns_ = 0; +}; diff --git a/obs-plugin/src/openstream-source.cpp b/obs-plugin/src/openstream-source.cpp index 34c3073..445751f 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 @@ -27,6 +28,7 @@ #include #include #include +#include #include #include #include @@ -90,6 +92,42 @@ using SwsContextPtr = std::unique_ptr; constexpr int kDiscoveryPort = 51515; constexpr int kDefaultListenerPort = 9000; +constexpr int kUsbLatencyMs = 30; + +std::optional usb_tether_gateway() { +#ifdef _WIN32 + ULONG size = 16 * 1024; + std::vector buffer(size); + auto *addresses = reinterpret_cast(buffer.data()); + ULONG result = GetAdaptersAddresses(AF_INET, GAA_FLAG_INCLUDE_GATEWAYS, nullptr, addresses, &size); + if (result == ERROR_BUFFER_OVERFLOW) { + buffer.resize(size); + addresses = reinterpret_cast(buffer.data()); + result = GetAdaptersAddresses(AF_INET, GAA_FLAG_INCLUDE_GATEWAYS, nullptr, addresses, &size); + } + if (result != NO_ERROR) return std::nullopt; + + for (auto *adapter = addresses; adapter; adapter = adapter->Next) { + std::wstring label = adapter->FriendlyName ? adapter->FriendlyName : L""; + if (adapter->Description) label += adapter->Description; + std::transform(label.begin(), label.end(), label.begin(), [](wchar_t c) { + return static_cast(std::towlower(c)); + }); + if (label.find(L"usb") == std::wstring::npos && + label.find(L"ndis") == std::wstring::npos && + label.find(L"mobile") == std::wstring::npos) { + continue; + } + for (auto *gateway = adapter->FirstGatewayAddress; gateway; gateway = gateway->Next) { + if (!gateway->Address.lpSockaddr || gateway->Address.lpSockaddr->sa_family != AF_INET) continue; + const auto *ipv4 = reinterpret_cast(gateway->Address.lpSockaddr); + char host[INET_ADDRSTRLEN] = {}; + if (inet_ntop(AF_INET, &ipv4->sin_addr, host, sizeof(host))) return std::string(host); + } + } +#endif + return std::nullopt; +} constexpr auto kReconnectReservationWindow = std::chrono::seconds(45); constexpr const char *kOpenStreamSourceName = "OpenStream V8"; constexpr const char *kDiscoveryMulticastAddress = "239.255.42.99"; @@ -769,6 +807,8 @@ struct OpenStreamSource { int listener_port = 0; int latency_ms = 120; int bitrate_mbps = 50; + bool usb_mode = false; + std::string usb_host; bool listener_enabled = true; std::atomic listener_running = false; std::atomic phone_connected = false; @@ -782,6 +822,8 @@ struct OpenStreamSource { int active_listener_port = 0; int active_latency_ms = 0; int active_bitrate_mbps = 0; + bool active_usb_mode = false; + std::string active_usb_host; std::string active_device_name; std::string active_slot_id; std::string active_slot_label; @@ -805,12 +847,11 @@ 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] { + return 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; } void set_slot_status(OpenStreamSource *ctx, std::string status) { @@ -1161,6 +1202,7 @@ 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 +1224,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; @@ -1267,7 +1309,7 @@ bool output_decoded_frame(OpenStreamSource *ctx, 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,8 +1384,20 @@ 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 source_timestamp_ns = [](const AVFrame *decoded, + const AVStream *stream) -> std::optional { + const int64_t timestamp = decoded->best_effort_timestamp != AV_NOPTS_VALUE + ? decoded->best_effort_timestamp + : decoded->pts; + if (!stream || timestamp == AV_NOPTS_VALUE) return std::nullopt; + return av_rescale_q(timestamp, stream->time_base, AVRational{1, 1000000000}); + }; + const auto drain_video = [&]() -> int { while (!ctx->stop_requested.load()) { const int result = avcodec_receive_frame(video_decoder_ctx, frame.get()); @@ -1356,8 +1410,19 @@ 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) { + output_decoded_frame(ctx, + video_stream, + video_decoder_ctx, + frame.get(), + *timestamp_ns, + &sws_ctx, + &bgra_buffer); + } else { + blog(LOG_WARNING, "[OpenStream] Dropping video frame without a usable source timestamp"); + } av_frame_unref(frame.get()); } return AVERROR_EXIT; @@ -1389,10 +1454,18 @@ 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"); + 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++) { @@ -1471,13 +1544,52 @@ void decode_packets(OpenStreamSource *ctx, } } -void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::string selected_phone_id) { +void openstream_worker(OpenStreamSource *ctx, + std::string base_srt_url, + std::string selected_phone_id, + bool usb_mode, + std::string configured_usb_host, + int phone_port) { avformat_network_init(); + if (usb_mode) { + in_addr configured_address = {}; + if (!configured_usb_host.empty() && + inet_pton(AF_INET, configured_usb_host.c_str(), &configured_address) != 1) { + set_slot_status(ctx, "Invalid USB phone IP"); + blog(LOG_ERROR, "[OpenStream] USB phone IP must be an IPv4 address"); + ctx->listener_running = false; + return; + } + std::optional host = configured_usb_host.empty() + ? usb_tether_gateway() + : std::optional(configured_usb_host); + while (!host && !ctx->stop_requested.load()) { + set_slot_status(ctx, "Waiting for USB tether"); + std::this_thread::sleep_for(std::chrono::seconds(1)); + host = usb_tether_gateway(); + } + if (!host) { + ctx->listener_running = false; + return; + } + PhoneDevice phone; + phone.instance_id = "usb:" + *host; + phone.name = "USB phone"; + phone.host = *host; + phone.port = phone_port; + phone.control_port = 9001; + phone.latency_ms = kUsbLatencyMs; + set_active_phone(ctx, phone); + base_srt_url = "srt://" + *host + ":" + std::to_string(phone_port) + + "?mode=caller&latency=" + std::to_string(kUsbLatencyMs); + blog(LOG_INFO, "[OpenStream] Using USB tether at %s", base_srt_url.c_str()); + } + while (!ctx->stop_requested.load()) { std::string srt_url = base_srt_url; std::optional reserved_phone; - if (srt_url == "openstream:auto") { + if (!usb_mode && srt_url == "openstream:auto") { std::optional phone; set_slot_status(ctx, "Waiting"); while (!ctx->stop_requested.load()) { @@ -1504,7 +1616,7 @@ void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::str "[OpenStream] Connecting source to phone %s at %s", phone->name.c_str(), srt_url.c_str()); - } else { + } else if (!usb_mode) { set_active_phone(ctx, std::nullopt); } @@ -1540,7 +1652,7 @@ void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::str if (reserved_phone.has_value()) { set_slot_status(ctx, "Reconnecting"); set_active_phone(ctx, reserved_phone); - } else { + } else if (!usb_mode) { set_active_phone(ctx, std::nullopt); } std::this_thread::sleep_for(std::chrono::milliseconds(500)); @@ -1564,7 +1676,7 @@ void openstream_worker(OpenStreamSource *ctx, std::string base_srt_url, std::str if (reserved_phone.has_value()) { set_slot_status(ctx, "Reconnecting"); set_active_phone(ctx, reserved_phone); - } else { + } else if (!usb_mode) { set_active_phone(ctx, std::nullopt); } if (!ctx->stop_requested.load()) { @@ -1615,6 +1727,8 @@ void openstream_start_worker(OpenStreamSource *ctx) { int listener_port = kDefaultListenerPort; int latency_ms = 120; int bitrate_mbps = 50; + bool usb_mode = false; + std::string usb_host; std::string source_name; std::string instance_id; std::string slot_id; @@ -1627,6 +1741,8 @@ void openstream_start_worker(OpenStreamSource *ctx) { listener_port = ctx->listener_port; latency_ms = ctx->latency_ms; bitrate_mbps = ctx->bitrate_mbps; + usb_mode = ctx->usb_mode; + usb_host = ctx->usb_host; source_name = ctx->device_name.empty() ? kOpenStreamSourceName : ctx->device_name; instance_id = ctx->instance_id; slot_id = ctx->slot_id; @@ -1642,6 +1758,8 @@ void openstream_start_worker(OpenStreamSource *ctx) { ctx->active_listener_port == listener_port && ctx->active_latency_ms == latency_ms && ctx->active_bitrate_mbps == bitrate_mbps && + ctx->active_usb_mode == usb_mode && + ctx->active_usb_host == usb_host && ctx->active_device_name == source_name && ctx->active_slot_id == slot_id && ctx->active_slot_label == slot_label && @@ -1657,6 +1775,8 @@ void openstream_start_worker(OpenStreamSource *ctx) { ctx->active_listener_port = listener_port; ctx->active_latency_ms = latency_ms; ctx->active_bitrate_mbps = bitrate_mbps; + ctx->active_usb_mode = usb_mode; + ctx->active_usb_host = usb_host; ctx->active_device_name = source_name; ctx->active_slot_id = slot_id; ctx->active_slot_label = slot_label; @@ -1665,16 +1785,26 @@ void openstream_start_worker(OpenStreamSource *ctx) { ctx->stop_requested = false; ctx->listener_running = true; ctx->phone_connected = false; - ctx->discovery.start(listener_port, - latency_ms, - bitrate_mbps, - source_name, - instance_id, - slot_id, - slot_label, - pairing_url, - &ctx->slot_busy); - ctx->worker = std::thread(openstream_worker, ctx, srt_url, selected_phone_id); + if (usb_mode) { + ctx->discovery.stop(); + } else { + ctx->discovery.start(listener_port, + latency_ms, + bitrate_mbps, + source_name, + instance_id, + slot_id, + slot_label, + pairing_url, + &ctx->slot_busy); + } + ctx->worker = std::thread(openstream_worker, + ctx, + srt_url, + selected_phone_id, + usb_mode, + usb_host, + listener_port); } void openstream_update(void *data, obs_data_t *settings) { @@ -1692,6 +1822,8 @@ void openstream_update(void *data, obs_data_t *settings) { 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->usb_mode = obs_data_get_bool(settings, "usb_mode"); + ctx->usb_host = obs_data_get_string(settings, "usb_host"); const char *slot_id = obs_data_get_string(settings, "slot_id"); if (slot_id && slot_id[0] != '\0') { ctx->slot_id = slot_id; @@ -1725,8 +1857,10 @@ void openstream_update(void *data, obs_data_t *settings) { ctx->slot_id, ctx->slot_label, ctx->instance_id); - ctx->pairing_hint = "Open OpenStream on your phone, choose " + ctx->slot_label + - ", and keep both devices on the same Wi-Fi. Pairing URL is in Advanced."; + ctx->pairing_hint = ctx->usb_mode + ? "Enable USB tethering on the phone, select USB in OpenStream, then start this slot." + : "Open OpenStream on your phone, choose " + ctx->slot_label + + ", and keep both devices on the same Wi-Fi. Pairing URL is in Advanced."; const std::vector phones = ctx->phone_discovery.devices(); ctx->phone_target_hint = "Waiting for a phone to choose " + ctx->slot_label; if (ctx->selected_phone_id == PhoneDiscoveryReceiver::kAutoPhoneId) { @@ -1816,6 +1950,8 @@ void openstream_defaults(obs_data_t *settings) { 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_bool(settings, "usb_mode", false); + obs_data_set_default_string(settings, "usb_host", ""); obs_data_set_default_double(settings, "cam_zoom", 1.0); } @@ -1950,6 +2086,16 @@ obs_properties_t *openstream_properties(void *data) { // ── Camera remote controls ── obs_properties_t *advanced_group = obs_properties_create(); obs_properties_add_bool(advanced_group, "listener_enabled", "Listen for this camera slot"); + obs_property_t *usb_mode = + obs_properties_add_bool(advanced_group, "usb_mode", "USB tethered connection"); + obs_property_set_long_description( + usb_mode, + "Uses Android USB tethering for a wired 30 ms SRT profile. No ADB tunnel or second encoder is used."); + obs_property_t *usb_host = + obs_properties_add_text(advanced_group, "usb_host", "USB phone IP (optional)", OBS_TEXT_DEFAULT); + obs_property_set_long_description( + usb_host, + "Leave blank to detect the USB/RNDIS gateway automatically; set the phone tether IP only if detection fails."); obs_properties_add_text(advanced_group, "source_instance_id", "Source instance ID", OBS_TEXT_INFO); obs_properties_add_text(advanced_group, "slot_id", "Slot ID", OBS_TEXT_INFO); obs_property_t *listener_port = @@ -2040,6 +2186,77 @@ obs_properties_t *openstream_properties(void *data) { return props; } +const char *openstream_usb_get_name(void *) { + return "OpenStream USB"; +} + +void openstream_usb_defaults(obs_data_t *settings) { + openstream_defaults(settings); + obs_data_set_default_string(settings, "device_name", "OpenStream USB"); + obs_data_set_default_string( + settings, + "usb_help", + "Enable USB tethering on Android, connect the cable, then click Connect USB."); + obs_data_set_default_bool(settings, "usb_mode", true); + obs_data_set_default_bool(settings, "listener_enabled", true); +} + +void openstream_usb_update(void *data, obs_data_t *settings) { + // This source has one job: connect over USB. Keep the transport forced even + // when an old scene collection carries Wi-Fi-era settings. + obs_data_set_bool(settings, "usb_mode", true); + obs_data_set_string(settings, "srt_url", "openstream:auto"); + openstream_update(data, settings); +} + +obs_properties_t *openstream_usb_properties(void *data) { + obs_properties_t *props = obs_properties_create(); + obs_property_t *status = obs_properties_add_text( + props, "slot_status", "Connection", OBS_TEXT_INFO); + obs_property_text_set_info_word_wrap(status, true); + obs_property_set_long_description( + status, + "OpenStream USB connects to the phone over Android USB tethering/RNDIS using the existing A/V stream."); + + obs_property_t *host = obs_properties_add_text( + props, "usb_host", "USB phone IP (optional)", OBS_TEXT_DEFAULT); + obs_property_set_long_description( + host, + "Leave blank to detect the Windows USB/RNDIS gateway automatically."); + obs_properties_add_int( + props, "listener_port", "Phone SRT port", 1024, 65535, 1); + + obs_properties_add_text( + props, + "usb_help", + "Setup", + OBS_TEXT_INFO); + obs_property_t *connect = obs_properties_add_button( + props, + "usb_connect", + "Connect USB", + [](obs_properties_t *, obs_property_t *, void *data) { + auto *ctx = static_cast(data); + if (!ctx) return false; + openstream_start_worker(ctx); + return true; + }); + obs_property_set_long_description( + connect, + "Connects this source over USB tethering. Enable USB tethering on Android first."); + obs_properties_add_button( + props, + "usb_disconnect", + "Disconnect", + [](obs_properties_t *, obs_property_t *, void *data) { + auto *ctx = static_cast(data); + if (!ctx) return false; + openstream_stop_worker(ctx); + return true; + }); + return props; +} + obs_source_info openstream_source_info = { .id = "openstream_phone_v8_source", .type = OBS_SOURCE_TYPE_INPUT, @@ -2063,6 +2280,18 @@ obs_source_info openstream_legacy_source_info = { .get_properties = openstream_properties, .update = openstream_update, }; + +obs_source_info openstream_usb_source_info = { + .id = "openstream_usb_source", + .type = OBS_SOURCE_TYPE_INPUT, + .output_flags = OBS_SOURCE_ASYNC_VIDEO | OBS_SOURCE_AUDIO, + .get_name = openstream_usb_get_name, + .create = openstream_create, + .destroy = openstream_destroy, + .get_defaults = openstream_usb_defaults, + .get_properties = openstream_usb_properties, + .update = openstream_usb_update, +}; } // namespace bool openstream_is_camera_source(obs_source_t *source) { @@ -2138,6 +2367,7 @@ bool obs_module_load(void) { #endif obs_register_source(&openstream_source_info); obs_register_source(&openstream_legacy_source_info); + obs_register_source(&openstream_usb_source_info); openstream_register_dock(); blog(LOG_INFO, "[OpenStream] OBS plugin loaded: V8 — video + audio + remote controls (Made by @yashas.vm)"); return true; diff --git a/tests/media_clock_test.cpp b/tests/media_clock_test.cpp new file mode 100644 index 0000000..a7d0d8a --- /dev/null +++ b/tests/media_clock_test.cpp @@ -0,0 +1,11 @@ +#include "media-clock.hpp" + +#include + +int main() { + MediaClock clock; + assert(!clock.map(-1, 1'000'000'000).has_value()); + assert(clock.map(10'000'000, 1'000'000'000) == 1'000'000'000); + assert(clock.map(30'000'000, 9'000'000'000) == 1'020'000'000); + assert(clock.map(15'000'000, 20'000'000'000) == 1'005'000'000); +} diff --git a/tests/test_repo_contract.py b/tests/test_repo_contract.py index 65cab08..c6be2e4 100644 --- a/tests/test_repo_contract.py +++ b/tests/test_repo_contract.py @@ -91,6 +91,44 @@ def test_audio_path_uses_adts_aac_and_obs_planar_formats() -> None: assert "obs_source_output_audio" in source +def test_obs_preserves_shared_source_timestamps_and_bounds_control_queue() -> None: + source = read("obs-plugin/src/openstream-source.cpp") + clock = read("obs-plugin/src/media-clock.hpp") + control = read("obs-plugin/src/async-control-client.cpp") + control_header = read("obs-plugin/src/async-control-client.hpp") + + assert "best_effort_timestamp" in source + assert "media_clock.map" in source + assert "timestamp = os_gettime_ns()" not in source + assert "source_origin_ns_" in clock + assert "kQueueCapacity = 16" in control_header + assert "commands_.size() >= kQueueCapacity" in control + + +def test_usb_mode_uses_rndis_gateway_and_not_adb_forwarding() -> None: + source = read("obs-plugin/src/openstream-source.cpp") + native = read("android/app/src/main/cpp/openstream_srt.cpp") + transport = read("android/app/src/main/java/dev/openstream/app/stream/TransportMode.kt") + + assert "GAA_FLAG_INCLUDE_GATEWAYS" in source + assert 'label.find(L"ndis")' in source + assert 'set_slot_status(ctx, "Waiting for USB tether")' in source + assert "kUsbLatencyMs = 30" in source + assert "kMinSrtLatencyMs = 30" in native + assert "USB_LATENCY_MS = 30" in transport + assert "adb forward" not in source.lower() + + +def test_openstream_usb_source_is_connect_only() -> None: + source = read("obs-plugin/src/openstream-source.cpp") + + assert '"OpenStream USB"' in source + assert '"openstream_usb_source"' in source + assert 'obs_data_set_bool(settings, "usb_mode", true)' in source + assert '"usb_connect"' in source + assert '"usb_disconnect"' in source + + def test_obs_plugin_routes_multiple_phones_by_selected_slot() -> None: source = read("obs-plugin/src/openstream-source.cpp") assert "selected_phone_id" in source