Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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**

---

Expand All @@ -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
Expand Down
4 changes: 3 additions & 1 deletion android/app/src/main/cpp/openstream_srt.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -467,7 +469,7 @@ std::optional<SrtUrl> 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) {
Expand Down
70 changes: 50 additions & 20 deletions android/app/src/main/java/dev/openstream/app/MainActivity.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 },
Expand Down Expand Up @@ -202,20 +206,20 @@ class MainActivity : Activity() {
override fun onStart() {
super.onStart()
activityStarted = true
phoneAdvertiser.start()
obsDiscoveryClient.start()
startDiscovery()
controlServer.start()
startPreviewIfAllowed()
startPhoneServerIfAllowed()
}

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
Expand All @@ -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()
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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()
}
}
}
Expand Down Expand Up @@ -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,
)
}

Expand Down Expand Up @@ -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) {
Expand Down
15 changes: 14 additions & 1 deletion android/app/src/main/java/dev/openstream/app/SettingsActivity.kt
Original file line number Diff line number Diff line change
Expand Up @@ -6,16 +6,20 @@ 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() {

private lateinit var inputObsHost: EditText
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
Expand All @@ -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)
Expand All @@ -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) {
Expand All @@ -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(
Expand All @@ -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()
Expand Down Expand Up @@ -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"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,13 +14,16 @@ 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,
private val onDevicesChanged: (List<DiscoveredObsDevice>) -> Unit,
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<String, DiscoveredObsDevice>()
private var socket: MulticastSocket? = null
Expand All @@ -29,17 +32,21 @@ 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()
}
}

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()
Expand All @@ -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)
Expand All @@ -92,14 +105,17 @@ class ObsDiscoveryClient(
publishDevices()
}
} catch (_: Exception) {
if (running.get()) {
if (running.get() && generation.get() == runGeneration) {
if (pruneExpired()) {
publishDevices()
}
}
}
}
udp.close()
synchronized(socketLock) {
if (socket === udp) socket = null
}
}

private fun joinDiscoveryMulticast(socket: MulticastSocket) {
Expand Down
Loading
Loading