From d9427e5d5287a47737964ee8a2846c60bbabb224 Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Tue, 23 Jun 2026 15:16:43 -0300 Subject: [PATCH 1/7] feat(moq-gst): add listen mode to moqsink for relay-less publishing Adds `listen` and `tls-generate` properties to moqsink so it can run its own QUIC/WebTransport server and serve a broadcast directly to subscribers that dial it, with no relay in between. `listen` is mutually exclusive with `url` (dial mode stays the default); the server branch in run_session mirrors the relay accept-and-serve path. Documents both properties in doc/bin/gstreamer.md. Co-Authored-By: Claude Opus 4.8 --- doc/bin/gstreamer.md | 20 +++++++ rs/moq-gst/src/sink/imp.rs | 114 +++++++++++++++++++++++++++++++------ 2 files changed, 118 insertions(+), 16 deletions(-) diff --git a/doc/bin/gstreamer.md b/doc/bin/gstreamer.md index bfb6818b07..d6fcd4aa52 100644 --- a/doc/bin/gstreamer.md +++ b/doc/bin/gstreamer.md @@ -30,6 +30,26 @@ 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`). When set, runs a server and ignores `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 ! ... +``` + ## 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..f9cadee480 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,46 @@ 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)) => Transport::Server { + bind: bind.to_string(), + tls_generate: value + .tls_generate + .as_deref() + .map(parse_hostnames) + .unwrap_or_default(), + }, + (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 +123,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 +168,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 +196,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 +208,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 +430,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 +442,43 @@ async fn run_session( settings.broadcast ); - let client = client.with_publish(origin.consume()); - let session = client.connect(settings.url.clone()).await?; + 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}"); + + 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; + tokio::spawn(async move { + match request.with_publish(origin_consumer).ok().await { + Ok(session) => { + gst::info!(CAT, "moqsink accepted session {id}"); + if let Err(err) = session.closed().await { + gst::warning!(CAT, "session {id} closed with error: {err:?}"); + } + } + Err(err) => gst::warning!(CAT, "failed to accept session {id}: {err:?}"), + } + }); + } + }); + (None, Some(task)) + } + }; let mut runtime = RuntimeState { session, @@ -440,6 +518,10 @@ async fn run_session( } } + if let Some(task) = accept_task { + task.abort(); + } + Ok(()) } From 94eed86017dfdd0050e9eaac02ac51edfa744ecb Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Wed, 24 Jun 2026 02:01:43 -0300 Subject: [PATCH 2/7] fix: group-order test build Add the `ordered` field to Track initializers in tests so the suite compiles after the group-order feature added that field to Track. Co-Authored-By: Claude Opus 4.8 --- rs/moq-net/src/stats.rs | 7 +++++++ 1 file changed, 7 insertions(+) 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; From 0d21a70d44bf77ddabcdb37b949b40aec7780ecd Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Wed, 24 Jun 2026 11:42:08 -0300 Subject: [PATCH 3/7] fix(moq-gst): address review on listen mode - require tls-generate in listen mode (server needs a cert; fail early with a clear error instead of a cryptic runtime init failure) - signal accepted sessions to close when the element stops, so detached session tasks don't leak connections past run_session - doc: state listen and url are mutually exclusive (the code already errors if both are set) Co-Authored-By: Claude Opus 4.8 --- doc/bin/gstreamer.md | 2 +- rs/moq-gst/src/sink/imp.rs | 43 ++++++++++++++++++++++++++++---------- 2 files changed, 33 insertions(+), 12 deletions(-) diff --git a/doc/bin/gstreamer.md b/doc/bin/gstreamer.md index d6fcd4aa52..a78e0a443a 100644 --- a/doc/bin/gstreamer.md +++ b/doc/bin/gstreamer.md @@ -36,7 +36,7 @@ By default `moqsink` dials a relay (`url`). It can instead run its own QUIC/WebT | Property | Type | Description | | -------------- | ------ | ------------------------------------------------------------------------------------ | -| `listen` | string | Bind address `host:port` (e.g. `0.0.0.0:4443`). When set, runs a server and ignores `url`. | +| `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. diff --git a/rs/moq-gst/src/sink/imp.rs b/rs/moq-gst/src/sink/imp.rs index f9cadee480..5602c87f07 100644 --- a/rs/moq-gst/src/sink/imp.rs +++ b/rs/moq-gst/src/sink/imp.rs @@ -70,14 +70,17 @@ impl TryFrom for ResolvedSettings { let transport = match (url, listen) { (Some(url), None) => Transport::Client { url: Url::parse(url)? }, - (None, Some(bind)) => Transport::Server { - bind: bind.to_string(), - tls_generate: value - .tls_generate - .as_deref() - .map(parse_hostnames) - .unwrap_or_default(), - }, + (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"), }; @@ -442,6 +445,8 @@ async fn run_session( settings.broadcast ); + let mut shutdown_tx = None; + let (session, accept_task) = match settings.transport.clone() { Transport::Client { url } => { let mut client_config = moq_native::ClientConfig::default(); @@ -457,18 +462,29 @@ async fn run_session( 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(session) => { + Ok(mut session) => { gst::info!(CAT, "moqsink accepted session {id}"); - if let Err(err) = session.closed().await { - gst::warning!(CAT, "session {id} closed with error: {err:?}"); + 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:?}"), @@ -476,6 +492,7 @@ async fn run_session( }); } }); + shutdown_tx = Some(tx); (None, Some(task)) } }; @@ -522,6 +539,10 @@ async fn run_session( task.abort(); } + if let Some(tx) = shutdown_tx { + let _ = tx.send(()); + } + Ok(()) } From d9babf7fb8120ba503b6a8a015c7251691cfd702 Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Fri, 3 Jul 2026 16:31:11 -0300 Subject: [PATCH 4/7] feat(moq-gst): add listen mode to moqsrc for relay-less subscribing Mirror of moqsink listen mode. moqsrc can run its own QUIC/WebTransport server (listen + tls-generate) and consume a broadcast from a publisher that dials it, with no relay. listen is mutually exclusive with url; the server branch in run_session accepts one session and consumes into the origin. Keeps the publisher as the dialing (QUIC client) peer, which is required for QUIC connection migration when the publisher's address changes. Co-Authored-By: Claude Opus 4.8 --- doc/bin/gstreamer.md | 18 +++++++ rs/moq-gst/src/source/imp.rs | 95 +++++++++++++++++++++++++++++++----- 2 files changed, 100 insertions(+), 13 deletions(-) diff --git a/doc/bin/gstreamer.md b/doc/bin/gstreamer.md index a78e0a443a..26f37f7df4 100644 --- a/doc/bin/gstreamer.md +++ b/doc/bin/gstreamer.md @@ -50,6 +50,24 @@ gst-launch-1.0 videotestsrc ! x264enc ! \ 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/source/imp.rs b/rs/moq-gst/src/source/imp.rs index f97e4f5053..76bfcc452a 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -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, @@ -190,7 +233,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 +278,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 +298,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(), @@ -457,14 +514,26 @@ 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 = server.accept().await.context("listener closed before a session arrived")?; + 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. From 05e69b86d42e9514ec984de1e2815d044df580a1 Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Mon, 6 Jul 2026 11:52:57 -0300 Subject: [PATCH 5/7] fix(moq-gst): wait for a non-empty catalog in moqsrc The publisher briefly publishes an empty catalog while replacing a track: the old importer drop removes its rendition before the successor re-adds it. moqsrc treated an empty catalog as a completed session and exited silently, permanently stalling the stream. Skip catalog updates with no tracks and wait for one that lists renditions. Co-Authored-By: Claude Opus 4.8 --- rs/moq-gst/src/source/imp.rs | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 76bfcc452a..16fe836c97 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -545,8 +545,17 @@ 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); + // The publisher briefly publishes an empty catalog while replacing a track: + // the old importer's drop removes its rendition before the successor re-adds it. + // An empty catalog is transient, so wait for an update that lists tracks. + let catalog = loop { + let update = catalog_updates.next().await?.context("catalog missing")?; + if !update.video.renditions.is_empty() || !update.audio.renditions.is_empty() { + break update; + } + tracing::debug!("ignoring catalog update with no tracks"); + }; let consumer_latency = Duration::from_millis(settings.max_latency_ms); From 7ec32485fa5c40032645c4278ba7f59887dc84dc Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Mon, 6 Jul 2026 12:16:42 -0300 Subject: [PATCH 6/7] fix(moq-gst): subscribe moqsrc tracks from catalog updates as they appear A single catalog snapshot races the publisher pipeline: audio registers at caps time, video only at the first keyframe, and a replaced track briefly publishes an empty catalog before reappearing under a new name. Reading one snapshot could yield no tracks (silent stall) or an audio-only or doomed track list. Keep consuming catalog updates and subscribe to each new track, deduplicated by name. Session errors now surface through the catalog stream, posting an element error instead of idling on a dead session. Drop the NoMorePads message since pads can now appear at any time. Co-Authored-By: Claude Opus 4.8 --- rs/moq-gst/src/source/imp.rs | 61 +++++++++++++++++++----------------- 1 file changed, 32 insertions(+), 29 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 16fe836c97..dfcdb88fa6 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; @@ -148,7 +148,6 @@ enum ControlMessage { caps: gst::Caps, reply: oneshot::Sender, }, - NoMorePads, ReportError(anyhow::Error), } @@ -426,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:?}"]); } @@ -546,36 +542,43 @@ async fn run_session( let catalog_track = broadcast.subscribe_track(&hang::catalog::Catalog::default_track())?; let mut catalog_updates = moq_mux::catalog::hang::Consumer::new(catalog_track); - // The publisher briefly publishes an empty catalog while replacing a track: - // the old importer's drop removes its rendition before the successor re-adds it. - // An empty catalog is transient, so wait for an update that lists tracks. - let catalog = loop { - let update = catalog_updates.next().await?.context("catalog missing")?; - if !update.video.renditions.is_empty() || !update.audio.renditions.is_empty() { - break update; - } - tracing::debug!("ignoring catalog update with no tracks"); - }; 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; From af7356091cddca945423da261a8a81ea4431835c Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Mon, 6 Jul 2026 18:04:37 -0300 Subject: [PATCH 7/7] fix(moq-gst): select on shutdown while moqsrc waits for an incoming session Stopping the element before a publisher dialed left the session task parked in server.accept() forever, keeping the UDP socket bound. Co-Authored-By: Claude Opus 4.8 --- rs/moq-gst/src/source/imp.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index dfcdb88fa6..75b8f03077 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -526,7 +526,10 @@ async fn run_session( server_config.tls.generate = tls_generate; let mut server = server_config.init()?; gst::info!(CAT, "moqsrc listening on {bind}"); - let request = server.accept().await.context("listener closed before a session arrived")?; + 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? } };