diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 013ec73a96..ee8cad2d11 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), }) } } @@ -116,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)); } @@ -132,6 +148,7 @@ impl SessionController { Self { shutdown: shutdown_tx, + latency: latency_tx, join, } } @@ -183,6 +200,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 +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" => { + 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!(), } } @@ -204,6 +235,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!(), } } @@ -413,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); @@ -448,8 +481,14 @@ 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)); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + .with_latency(settings.max_latency); + tasks.push(spawn_track_pump( + track, + descriptor, + endpoint, + shutdown.clone(), + latency.clone(), + )); } for (track_name, config) in catalog.audio.renditions { @@ -462,8 +501,14 @@ 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)); - tasks.push(spawn_track_pump(track, descriptor, endpoint, shutdown.clone())); + .with_latency(settings.max_latency); + tasks.push(spawn_track_pump( + track, + descriptor, + endpoint, + shutdown.clone(), + latency.clone(), + )); } let _ = control_tx.send(ControlMessage::NoMorePads); @@ -498,8 +543,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( @@ -507,6 +553,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 { @@ -515,6 +562,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)) => {