From bdd694a1af04ccbe3cb3024820ca26a5864db7a9 Mon Sep 17 00:00:00 2001 From: Joaquin Bartaburu Date: Sat, 6 Jun 2026 02:45:51 -0300 Subject: [PATCH 1/3] feat(moqsrc): add configurable max-latency-ms GObject property Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/src/source/imp.rs | 29 ++++++++++++++++++++++++++--- 1 file changed, 26 insertions(+), 3 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 013ec73a96..9d80d4cc62 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -20,11 +20,23 @@ static RUNTIME: LazyLock = LazyLock::new(|| { .expect("spawn tokio runtime") }); -#[derive(Debug, Clone, Default)] +#[derive(Debug, Clone)] struct Settings { url: Option, broadcast: Option, tls_disable_verify: bool, + max_latency_ms: u64, +} + +impl Default for Settings { + fn default() -> Self { + Self { + url: None, + broadcast: None, + tls_disable_verify: false, + max_latency_ms: 1000, + } + } } #[derive(Debug, Clone)] @@ -32,6 +44,7 @@ struct ResolvedSettings { url: url::Url, broadcast: String, tls_disable_verify: bool, + max_latency: Duration, } impl TryFrom for ResolvedSettings { @@ -46,6 +59,7 @@ impl TryFrom for ResolvedSettings { .context("broadcast property is required")? .clone(), tls_disable_verify: value.tls_disable_verify, + max_latency: Duration::from_millis(value.max_latency_ms), }) } } @@ -183,6 +197,13 @@ impl ObjectImpl for MoqSrc { .blurb("Disable TLS certificate verification") .default_value(false) .build(), + glib::ParamSpecUInt64::builder("max-latency-ms") + .nick("Max Latency (ms)") + .blurb("Drop groups older than this when the buffer span exceeds it") + .default_value(1000) + .minimum(0) + .maximum(u32::MAX as u64) + .build(), ] }); PROPS.as_ref() @@ -194,6 +215,7 @@ impl ObjectImpl for MoqSrc { "url" => settings.url = 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(), _ => unreachable!(), } } @@ -204,6 +226,7 @@ impl ObjectImpl for MoqSrc { "url" => settings.url.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(), _ => unreachable!(), } } @@ -448,7 +471,7 @@ async fn run_session( let track_ref = moq_net::Track::new(&track_name); let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) - .with_latency(Duration::from_secs(1)); + .with_latency(settings.max_latency); tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); } @@ -462,7 +485,7 @@ async fn run_session( let track_ref = moq_net::Track::new(&track_name); let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) - .with_latency(Duration::from_secs(1)); + .with_latency(settings.max_latency); tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); } From 273b1be6516829e13e0e893a70e372e55bafec09 Mon Sep 17 00:00:00 2001 From: Joaquin Bartaburu Date: Mon, 8 Jun 2026 11:04:37 -0300 Subject: [PATCH 2/3] feat(moqsrc): propagate max-latency-ms changes to active track consumers Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/src/source/imp.rs | 25 ++++++++++++++++++++----- 1 file changed, 20 insertions(+), 5 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 9d80d4cc62..1b1b7e9e4c 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -130,15 +130,17 @@ struct PadHandle { struct SessionController { shutdown: watch::Sender, + latency: watch::Sender, join: tokio::task::JoinHandle<()>, } impl SessionController { fn start(settings: ResolvedSettings, control_tx: mpsc::UnboundedSender) -> Self { let (shutdown_tx, mut shutdown_rx) = watch::channel(false); + let (latency_tx, latency_rx) = watch::channel(settings.max_latency); let control_for_error = control_tx.clone(); let join = RUNTIME.spawn(async move { - let result = run_session(settings, control_tx, &mut shutdown_rx).await; + let result = run_session(settings, control_tx, &mut shutdown_rx, latency_rx).await; if let Err(err) = result { let _ = control_for_error.send(ControlMessage::ReportError(err)); } @@ -146,6 +148,7 @@ impl SessionController { Self { shutdown: shutdown_tx, + latency: latency_tx, join, } } @@ -215,7 +218,13 @@ impl ObjectImpl for MoqSrc { "url" => settings.url = 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(), + "max-latency-ms" => { + let ms: u64 = value.get().unwrap(); + settings.max_latency_ms = ms; + if let Some(session) = self.session.lock().unwrap().as_ref() { + let _ = session.latency.send(Duration::from_millis(ms)); + } + } _ => unreachable!(), } } @@ -436,6 +445,7 @@ async fn run_session( settings: ResolvedSettings, control_tx: mpsc::UnboundedSender, shutdown: &mut watch::Receiver, + latency: watch::Receiver, ) -> Result<()> { let mut config = moq_native::ClientConfig::default(); config.tls.disable_verify = Some(settings.tls_disable_verify); @@ -472,7 +482,7 @@ async fn run_session( let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) .with_latency(settings.max_latency); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone(), latency.clone())); } for (track_name, config) in catalog.audio.renditions { @@ -486,7 +496,7 @@ async fn run_session( let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) .with_latency(settings.max_latency); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone(), latency.clone())); } let _ = control_tx.send(ControlMessage::NoMorePads); @@ -521,8 +531,9 @@ fn spawn_track_pump( descriptor: TrackDescriptor, pad_endpoint: PadEndpoint, shutdown: watch::Receiver, + latency: watch::Receiver, ) -> tokio::task::JoinHandle<()> { - RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown)) + RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown, latency)) } async fn run_track_pump( @@ -530,6 +541,7 @@ async fn run_track_pump( descriptor: TrackDescriptor, pad_endpoint: PadEndpoint, mut shutdown: watch::Receiver, + mut latency: watch::Receiver, ) { let mut reference_ts = None; loop { @@ -538,6 +550,9 @@ async fn run_track_pump( pad_endpoint.send(PadMessage::Drop); break; } + _ = latency.changed() => { + track.set_latency(*latency.borrow()); + } frame = track.read() => { match frame { Ok(Some(frame)) => { From 867e8bbee584bb8d8a7967604eb883f97ae5e8a1 Mon Sep 17 00:00:00 2001 From: Joaquin Bartaburu Date: Mon, 8 Jun 2026 12:35:50 -0300 Subject: [PATCH 3/3] style(moq-gst): fix rustfmt formatting in spawn_track_pump calls --- rs/moq-gst/src/source/imp.rs | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 1b1b7e9e4c..ee8cad2d11 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -482,7 +482,13 @@ async fn run_session( let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) .with_latency(settings.max_latency); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone(), latency.clone())); + tasks.push(spawn_track_pump( + track, + descriptor, + endpoint, + shutdown.clone(), + latency.clone(), + )); } for (track_name, config) in catalog.audio.renditions { @@ -496,7 +502,13 @@ async fn run_session( let track_consumer = broadcast.subscribe_track(&track_ref)?; let track = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy) .with_latency(settings.max_latency); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone(), latency.clone())); + tasks.push(spawn_track_pump( + track, + descriptor, + endpoint, + shutdown.clone(), + latency.clone(), + )); } let _ = control_tx.send(ControlMessage::NoMorePads);