diff --git a/doc/bin/gstreamer.md b/doc/bin/gstreamer.md index bfb6818b07..26f37f7df4 100644 --- a/doc/bin/gstreamer.md +++ b/doc/bin/gstreamer.md @@ -30,6 +30,44 @@ Both elements support the following properties: For `http://` URLs, `moq-native` automatically fetches the server's certificate fingerprint from `/certificate.sha256` and verifies TLS against it. You don't need `tls-disable-verify` for local development. ::: +### moqsink listen mode + +By default `moqsink` dials a relay (`url`). It can instead run its own QUIC/WebTransport server and serve the broadcast to subscribers that dial it directly, with no relay in between. These properties are moqsink-only and `listen` is mutually exclusive with `url`. + +| Property | Type | Description | +| -------------- | ------ | ------------------------------------------------------------------------------------ | +| `listen` | string | Bind address `host:port` (e.g. `0.0.0.0:4443`). Runs a server in listen mode. Mutually exclusive with `url`. | +| `tls-generate` | string | Comma-separated hostnames for a self-signed certificate. Listen mode only. | + +The self-signed certificate is only trusted if the subscriber uses the fingerprint or sets `tls-disable-verify`. The server serves its published broadcast to any connecting peer regardless of the dialed URL path; the subscriber selects the broadcast by `broadcast` name. + +```bash +# Publisher: listen instead of dialing a relay. +gst-launch-1.0 videotestsrc ! x264enc ! \ + moqsink name=mux listen=0.0.0.0:4443 tls-generate=localhost broadcast=bbb + +# Subscriber: dial the publisher directly. +gst-launch-1.0 moqsrc url=https://localhost:4443 broadcast=bbb tls-disable-verify=true ! ... +``` + +### moqsrc listen mode + +By default `moqsrc` dials a relay or publisher (`url`). It can instead run its own QUIC/WebTransport server and consume the broadcast from a publisher that dials it directly, with no relay in between. This is the mirror of moqsink listen mode: the subscriber is the server and the publisher is the client, so the publisher stays the dialing peer (useful when the publisher must survive its own address changing, since QUIC connection migration is a client-only capability). These properties mirror moqsink and `listen` is mutually exclusive with `url`. + +| Property | Type | Description | +| -------------- | ------ | ------------------------------------------------------------------------------------ | +| `listen` | string | Bind address `host:port` (e.g. `0.0.0.0:4443`). Runs a server in listen mode. Mutually exclusive with `url`. | +| `tls-generate` | string | Comma-separated hostnames for a self-signed certificate. Listen mode only. | + +```bash +# Subscriber: listen instead of dialing. +gst-launch-1.0 moqsrc name=src listen=0.0.0.0:4443 tls-generate=localhost broadcast=bbb ! ... + +# Publisher: dial the subscriber directly. +gst-launch-1.0 videotestsrc ! x264enc ! \ + moqsink name=mux url=https://localhost:4443 broadcast=bbb tls-disable-verify=true +``` + ## Prerequisites The plugin requires GStreamer development libraries. It is **not** built by default since most users don't have them installed. diff --git a/rs/moq-gst/src/sink/imp.rs b/rs/moq-gst/src/sink/imp.rs index a4ae7169ef..5602c87f07 100644 --- a/rs/moq-gst/src/sink/imp.rs +++ b/rs/moq-gst/src/sink/imp.rs @@ -36,13 +36,21 @@ static CAT: LazyLock = LazyLock::new(|| { #[derive(Debug, Clone, Default)] struct Settings { url: Option, + listen: Option, + tls_generate: Option, broadcast: Option, tls_disable_verify: bool, } +#[derive(Debug, Clone)] +enum Transport { + Client { url: Url }, + Server { bind: String, tls_generate: Vec }, +} + #[derive(Debug, Clone)] struct ResolvedSettings { - url: Url, + transport: Transport, broadcast: String, tls_disable_verify: bool, } @@ -51,18 +59,49 @@ impl TryFrom for ResolvedSettings { type Error = anyhow::Error; fn try_from(value: Settings) -> Result { + let broadcast = value + .broadcast + .as_ref() + .context("broadcast property is required")? + .clone(); + + let url = value.url.as_deref().filter(|s| !s.is_empty()); + let listen = value.listen.as_deref().filter(|s| !s.is_empty()); + + let transport = match (url, listen) { + (Some(url), None) => Transport::Client { url: Url::parse(url)? }, + (None, Some(bind)) => { + let tls_generate = value.tls_generate.as_deref().map(parse_hostnames).unwrap_or_default(); + anyhow::ensure!( + !tls_generate.is_empty(), + "tls-generate is required in listen mode to generate a self-signed certificate" + ); + Transport::Server { + bind: bind.to_string(), + tls_generate, + } + } + (Some(_), Some(_)) => anyhow::bail!("set exactly one of `url` or `listen`, not both"), + (None, None) => anyhow::bail!("one of `url` or `listen` is required"), + }; + Ok(Self { - url: Url::parse(value.url.as_ref().context("url property is required")?)?, - broadcast: value - .broadcast - .as_ref() - .context("broadcast property is required")? - .clone(), + transport, + broadcast, tls_disable_verify: value.tls_disable_verify, }) } } +fn parse_hostnames(value: &str) -> Vec { + value + .split(',') + .map(str::trim) + .filter(|host| !host.is_empty()) + .map(str::to_string) + .collect() +} + #[derive(Debug)] struct SessionHandle { sender: mpsc::UnboundedSender, @@ -87,7 +126,7 @@ struct PadState { struct RuntimeState { #[allow(dead_code)] - session: moq_net::Session, + session: Option, broadcast: moq_net::BroadcastProducer, catalog: moq_mux::catalog::hang::Producer, pads: HashMap, @@ -132,7 +171,15 @@ impl ObjectImpl for MoqSink { vec![ glib::ParamSpecString::builder("url") .nick("Source URL") - .blurb("Connect to the given URL") + .blurb("Dial the given URL (client mode). Mutually exclusive with `listen`") + .build(), + glib::ParamSpecString::builder("listen") + .nick("Listen address") + .blurb("Run a server bound to host:port and serve the broadcast to subscribers that dial it. Mutually exclusive with `url`") + .build(), + glib::ParamSpecString::builder("tls-generate") + .nick("TLS generate hostnames") + .blurb("Comma-separated hostnames for a self-signed certificate (listen mode only)") .build(), glib::ParamSpecString::builder("broadcast") .nick("Broadcast") @@ -152,6 +199,8 @@ impl ObjectImpl for MoqSink { let mut settings = self.settings.lock().unwrap(); match pspec.name() { "url" => settings.url = value.get().unwrap(), + "listen" => settings.listen = value.get().unwrap(), + "tls-generate" => settings.tls_generate = value.get().unwrap(), "broadcast" => settings.broadcast = value.get().unwrap(), "tls-disable-verify" => settings.tls_disable_verify = value.get().unwrap(), _ => unreachable!(), @@ -162,6 +211,8 @@ impl ObjectImpl for MoqSink { let settings = self.settings.lock().unwrap(); match pspec.name() { "url" => settings.url.to_value(), + "listen" => settings.listen.to_value(), + "tls-generate" => settings.tls_generate.to_value(), "broadcast" => settings.broadcast.to_value(), "tls-disable-verify" => settings.tls_disable_verify.to_value(), _ => unreachable!(), @@ -382,11 +433,6 @@ async fn run_session( mut rx: mpsc::UnboundedReceiver, element_weak: gst::glib::WeakRef, ) -> Result<()> { - let mut client_config = moq_native::ClientConfig::default(); - client_config.tls.disable_verify = Some(settings.tls_disable_verify); - - let client = client_config.init()?; - let origin = moq_net::Origin::random().produce(); let mut broadcast = moq_net::Broadcast::new().produce(); let broadcast_consumer = broadcast.consume(); @@ -399,8 +445,57 @@ async fn run_session( settings.broadcast ); - let client = client.with_publish(origin.consume()); - let session = client.connect(settings.url.clone()).await?; + let mut shutdown_tx = None; + + let (session, accept_task) = match settings.transport.clone() { + Transport::Client { url } => { + let mut client_config = moq_native::ClientConfig::default(); + client_config.tls.disable_verify = Some(settings.tls_disable_verify); + let client = client_config.init()?.with_publish(origin.consume()); + let session = client.connect(url).await?; + (Some(session), None) + } + Transport::Server { bind, tls_generate } => { + let mut server_config = moq_native::ServerConfig::default(); + server_config.bind = Some(bind.clone()); + server_config.tls.generate = tls_generate; + let mut server = server_config.init()?; + gst::info!(CAT, "moqsink listening on {bind}"); + + // Signals accepted sessions to close when the element stops, so detached + // session tasks don't leak connections past run_session. + let (tx, _) = tokio::sync::broadcast::channel::<()>(1); + let accept_tx = tx.clone(); + + let task = RUNTIME.spawn(async move { + let mut conn_id: u64 = 0; + while let Some(request) = server.accept().await { + let origin_consumer = origin.consume(); + let id = conn_id; + conn_id += 1; + let mut shutdown_rx = accept_tx.subscribe(); + tokio::spawn(async move { + match request.with_publish(origin_consumer).ok().await { + Ok(mut session) => { + gst::info!(CAT, "moqsink accepted session {id}"); + tokio::select! { + _ = shutdown_rx.recv() => session.close(moq_net::Error::Cancel), + res = session.closed() => { + if let Err(err) = res { + gst::warning!(CAT, "session {id} closed with error: {err:?}"); + } + } + } + } + Err(err) => gst::warning!(CAT, "failed to accept session {id}: {err:?}"), + } + }); + } + }); + shutdown_tx = Some(tx); + (None, Some(task)) + } + }; let mut runtime = RuntimeState { session, @@ -440,6 +535,14 @@ async fn run_session( } } + if let Some(task) = accept_task { + task.abort(); + } + + if let Some(tx) = shutdown_tx { + let _ = tx.send(()); + } + Ok(()) } diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index f97e4f5053..75b8f03077 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -1,4 +1,4 @@ -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::sync::{LazyLock, Mutex}; use std::time::Duration; @@ -23,6 +23,8 @@ static RUNTIME: LazyLock = LazyLock::new(|| { #[derive(Debug, Clone)] struct Settings { url: Option, + listen: Option, + tls_generate: Option, broadcast: Option, tls_disable_verify: bool, max_latency_ms: u64, @@ -33,6 +35,8 @@ impl Default for Settings { fn default() -> Self { Self { url: None, + listen: None, + tls_generate: None, broadcast: None, tls_disable_verify: false, max_latency_ms: 1000, @@ -41,9 +45,15 @@ impl Default for Settings { } } +#[derive(Debug, Clone)] +enum Transport { + Client { url: url::Url }, + Server { bind: String, tls_generate: Vec }, +} + #[derive(Debug, Clone)] struct ResolvedSettings { - url: url::Url, + transport: Transport, broadcast: String, tls_disable_verify: bool, max_latency_ms: u64, @@ -54,13 +64,37 @@ impl TryFrom for ResolvedSettings { type Error = anyhow::Error; fn try_from(value: Settings) -> Result { + let broadcast = value + .broadcast + .as_ref() + .context("broadcast property is required")? + .clone(); + + let url = value.url.as_deref().filter(|s| !s.is_empty()); + let listen = value.listen.as_deref().filter(|s| !s.is_empty()); + + let transport = match (url, listen) { + (Some(url), None) => Transport::Client { + url: url::Url::parse(url)?, + }, + (None, Some(bind)) => { + let tls_generate = value.tls_generate.as_deref().map(parse_hostnames).unwrap_or_default(); + anyhow::ensure!( + !tls_generate.is_empty(), + "tls-generate is required in listen mode to generate a self-signed certificate" + ); + Transport::Server { + bind: bind.to_string(), + tls_generate, + } + } + (Some(_), Some(_)) => anyhow::bail!("set exactly one of `url` or `listen`, not both"), + (None, None) => anyhow::bail!("one of `url` or `listen` is required"), + }; + Ok(Self { - url: url::Url::parse(value.url.as_ref().context("url property is required")?)?, - broadcast: value - .broadcast - .as_ref() - .context("broadcast property is required")? - .clone(), + transport, + broadcast, tls_disable_verify: value.tls_disable_verify, max_latency_ms: value.max_latency_ms, ascending: value.ascending, @@ -68,6 +102,15 @@ impl TryFrom for ResolvedSettings { } } +fn parse_hostnames(value: &str) -> Vec { + value + .split(',') + .map(str::trim) + .filter(|host| !host.is_empty()) + .map(str::to_string) + .collect() +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] enum TrackKind { Video, @@ -105,7 +148,6 @@ enum ControlMessage { caps: gst::Caps, reply: oneshot::Sender, }, - NoMorePads, ReportError(anyhow::Error), } @@ -190,7 +232,17 @@ impl ObjectImpl for MoqSrc { vec![ glib::ParamSpecString::builder("url") .nick("Source URL") - .blurb("Connect to the given URL") + .blurb("Dial the given URL (client mode). Mutually exclusive with `listen`") + .mutable_ready() + .build(), + glib::ParamSpecString::builder("listen") + .nick("Listen address") + .blurb("Run a server bound to host:port and consume the broadcast from a publisher that dials it. Mutually exclusive with `url`") + .mutable_ready() + .build(), + glib::ParamSpecString::builder("tls-generate") + .nick("TLS generate hostnames") + .blurb("Comma-separated hostnames for a self-signed certificate (listen mode only)") .mutable_ready() .build(), glib::ParamSpecString::builder("broadcast") @@ -225,6 +277,8 @@ impl ObjectImpl for MoqSrc { let mut settings = self.settings.lock().unwrap(); match pspec.name() { "url" => settings.url = value.get().unwrap(), + "listen" => settings.listen = value.get().unwrap(), + "tls-generate" => settings.tls_generate = value.get().unwrap(), "broadcast" => settings.broadcast = value.get().unwrap(), "tls-disable-verify" => settings.tls_disable_verify = value.get().unwrap(), "max-latency-ms" => settings.max_latency_ms = value.get().unwrap(), @@ -243,6 +297,8 @@ impl ObjectImpl for MoqSrc { let settings = self.settings.lock().unwrap(); match pspec.name() { "url" => settings.url.to_value(), + "listen" => settings.listen.to_value(), + "tls-generate" => settings.tls_generate.to_value(), "broadcast" => settings.broadcast.to_value(), "tls-disable-verify" => settings.tls_disable_verify.to_value(), "max-latency-ms" => settings.max_latency_ms.to_value(), @@ -369,9 +425,6 @@ impl MoqSrc { gst::error!(CAT, obj = self.obj(), "failed to create pad: {err:?}"); } } - ControlMessage::NoMorePads => { - self.obj().no_more_pads(); - } ControlMessage::ReportError(err) => { gst::element_error!(self.obj(), gst::CoreError::Failed, ("session error"), ["{err:?}"]); } @@ -457,14 +510,29 @@ async fn run_session( control_tx: mpsc::UnboundedSender, shutdown: &mut watch::Receiver, ) -> Result<()> { - let mut config = moq_native::ClientConfig::default(); - config.tls.disable_verify = Some(settings.tls_disable_verify); - let origin = moq_net::Origin::random().produce(); let origin_consumer = origin.consume(); - let client = config.init()?.with_consume(origin); - let _session = client.connect(settings.url.clone()).await?; + let _session = match settings.transport.clone() { + Transport::Client { url } => { + let mut config = moq_native::ClientConfig::default(); + config.tls.disable_verify = Some(settings.tls_disable_verify); + let client = config.init()?.with_consume(origin); + client.connect(url).await? + } + Transport::Server { bind, tls_generate } => { + let mut server_config = moq_native::ServerConfig::default(); + server_config.bind = Some(bind.clone()); + server_config.tls.generate = tls_generate; + let mut server = server_config.init()?; + gst::info!(CAT, "moqsrc listening on {bind}"); + let request = tokio::select! { + request = server.accept() => request.context("listener closed before a session arrived")?, + _ = shutdown.changed() => return Ok(()), + }; + request.with_consume(origin).ok().await? + } + }; // Wait for the broadcast to be announced. Synchronous lookup would race the gossip of // announcements that happens after the session is established. @@ -476,28 +544,44 @@ async fn run_session( }; let catalog_track = broadcast.subscribe_track(&hang::catalog::Catalog::default_track())?; - let mut catalog = moq_mux::catalog::hang::Consumer::new(catalog_track); - let catalog = catalog.next().await?.context("catalog missing")?.clone(); + let mut catalog_updates = moq_mux::catalog::hang::Consumer::new(catalog_track); let consumer_latency = Duration::from_millis(settings.max_latency_ms); - + let mut known_tracks = HashSet::new(); let mut tasks = Vec::new(); - for (track_name, config) in catalog.video.renditions { - let descriptor = TrackDescriptor { kind: TrackKind::Video, name: track_name.clone() }; - let endpoint = request_pad(&control_tx, descriptor.clone(), video_caps(&config)?).await?; - let track = subscribe_track(&broadcast, &track_name, settings.ascending, consumer_latency)?; - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); - } + // Tracks register in the catalog at different times (audio at caps time, video + // at the first keyframe) and can be replaced under a new name mid-stream, so + // keep consuming catalog updates and subscribe to tracks as they appear. + loop { + let update = tokio::select! { + update = catalog_updates.next() => match update? { + Some(update) => update, + None => break, + }, + _ = shutdown.changed() => break, + }; - for (track_name, config) in catalog.audio.renditions { - let descriptor = TrackDescriptor { kind: TrackKind::Audio, name: track_name.clone() }; - let endpoint = request_pad(&control_tx, descriptor.clone(), audio_caps(&config)?).await?; - let track = subscribe_track(&broadcast, &track_name, settings.ascending, consumer_latency)?; - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); - } + for (track_name, config) in update.video.renditions { + if !known_tracks.insert(track_name.clone()) { + continue; + } + let descriptor = TrackDescriptor { kind: TrackKind::Video, name: track_name.clone() }; + let endpoint = request_pad(&control_tx, descriptor.clone(), video_caps(&config)?).await?; + let track = subscribe_track(&broadcast, &track_name, settings.ascending, consumer_latency)?; + tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + } - let _ = control_tx.send(ControlMessage::NoMorePads); + for (track_name, config) in update.audio.renditions { + if !known_tracks.insert(track_name.clone()) { + continue; + } + let descriptor = TrackDescriptor { kind: TrackKind::Audio, name: track_name.clone() }; + let endpoint = request_pad(&control_tx, descriptor.clone(), audio_caps(&config)?).await?; + let track = subscribe_track(&broadcast, &track_name, settings.ascending, consumer_latency)?; + tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + } + } for task in tasks { let _ = task.await; diff --git a/rs/moq-net/src/stats.rs b/rs/moq-net/src/stats.rs index bbb1b2b7ae..7189c1ac4b 100644 --- a/rs/moq-net/src/stats.rs +++ b/rs/moq-net/src/stats.rs @@ -1458,6 +1458,7 @@ mod tests { .subscribe_track(&Track { name: "publisher.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); let frame = read_frame(track).await; @@ -1486,6 +1487,7 @@ mod tests { .subscribe_track(&Track { name: "publisher.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); let frame = read_frame(track).await; @@ -1522,6 +1524,7 @@ mod tests { .subscribe_track(&Track { name: "publisher.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); let frame = read_frame(track).await; @@ -1650,6 +1653,7 @@ mod tests { .subscribe_track(&Track { name: "sessions.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); let frame = read_session_frame(track).await; @@ -1665,6 +1669,7 @@ mod tests { .subscribe_track(&Track { name: "internal/sessions.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); let snap = *read_session_frame(int_track).await.get("peer").expect("internal entry"); @@ -1716,6 +1721,7 @@ mod tests { .subscribe_track(&Track { name: "publisher.json".into(), priority: 0, + ordered: false, }) .expect("subscribe"); assert!( @@ -1730,6 +1736,7 @@ mod tests { .subscribe_track(&Track { name: name.into(), priority: 0, + ordered: false, }) .expect("subscribe"); let frame = read_frame(t).await;