From 5eb39503371c187e3d223d7b4fdd0a438dc9518a Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Wed, 10 Jun 2026 15:22:15 -0300 Subject: [PATCH 1/5] feat(moq-gst): add group-order feature and ascending/max-latency-ms properties - Declare group-order Cargo feature - Add max_latency_ms and ascending to Settings/ResolvedSettings - Expose max-latency-ms and ascending as GStreamer properties on moqsrc - Set track_ref.ordered = ascending behind #[cfg(feature = "group-order")] - Extract subscribe_track() helper to eliminate duplicated cfg block - Consumer latency uses Duration::MAX when ascending, else max_latency_ms Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/Cargo.toml | 71 +- rs/moq-gst/src/source/imp.rs | 1373 +++++++++++++++++----------------- 2 files changed, 744 insertions(+), 700 deletions(-) diff --git a/rs/moq-gst/Cargo.toml b/rs/moq-gst/Cargo.toml index 2b3dbd431f..765d24a5a5 100644 --- a/rs/moq-gst/Cargo.toml +++ b/rs/moq-gst/Cargo.toml @@ -1,34 +1,37 @@ -[package] -name = "moq-gst" -description = "Media over QUIC - GStreamer plugin" -authors = ["Luke Curley"] -repository = "https://github.com/kixelated/moq" -license = "MIT OR Apache-2.0" - -version = "0.2.2" -edition = "2021" -rust-version.workspace = true -publish = false - -[lib] -name = "gstmoq" -crate-type = ["cdylib", "rlib"] -test = false -doctest = false - -[dependencies] - -anyhow = { version = "1", features = ["backtrace"] } -bytes = "1" -gst = { package = "gstreamer", version = "0.23" } -hang = { workspace = true } -moq-mux = { workspace = true } -moq-native = { workspace = true, default-features = true } -moq-net = { workspace = true } -tokio = { workspace = true, features = ["full"] } -tracing = "0.1" -tracing-subscriber = "0.3" -url = "2" - -[build-dependencies] -gst-plugin-version-helper = "0.8" +[package] +name = "moq-gst" +description = "Media over QUIC - GStreamer plugin" +authors = ["Luke Curley"] +repository = "https://github.com/kixelated/moq" +license = "MIT OR Apache-2.0" + +version = "0.2.2" +edition = "2021" +rust-version.workspace = true +publish = false + +[lib] +name = "gstmoq" +crate-type = ["cdylib", "rlib"] +test = false +doctest = false + +[features] +group-order = [] + +[dependencies] + +anyhow = { version = "1", features = ["backtrace"] } +bytes = "1" +gst = { package = "gstreamer", version = "0.23" } +hang = { workspace = true } +moq-mux = { workspace = true } +moq-native = { workspace = true, default-features = true } +moq-net = { workspace = true } +tokio = { workspace = true, features = ["full"] } +tracing = "0.1" +tracing-subscriber = "0.3" +url = "2" + +[build-dependencies] +gst-plugin-version-helper = "0.8" diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 013ec73a96..a4abf5591e 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -1,666 +1,707 @@ -use std::collections::HashMap; -use std::sync::{LazyLock, Mutex}; -use std::time::Duration; - -use anyhow::{bail, Context, Result}; -use gst::glib; -use gst::prelude::*; -use gst::subclass::prelude::*; -use tokio::sync::{mpsc, oneshot, watch}; - -use hang::moq_net; - -static CAT: LazyLock = - LazyLock::new(|| gst::DebugCategory::new("moq-src", gst::DebugColorFlags::empty(), Some("MoQ Source Element"))); - -static RUNTIME: LazyLock = LazyLock::new(|| { - tokio::runtime::Builder::new_multi_thread() - .enable_all() - .build() - .expect("spawn tokio runtime") -}); - -#[derive(Debug, Clone, Default)] -struct Settings { - url: Option, - broadcast: Option, - tls_disable_verify: bool, -} - -#[derive(Debug, Clone)] -struct ResolvedSettings { - url: url::Url, - broadcast: String, - tls_disable_verify: bool, -} - -impl TryFrom for ResolvedSettings { - type Error = anyhow::Error; - - fn try_from(value: Settings) -> Result { - 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(), - tls_disable_verify: value.tls_disable_verify, - }) - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -enum TrackKind { - Video, - Audio, -} - -impl TrackKind { - fn template_name(&self) -> &'static str { - match self { - TrackKind::Video => "video_%u", - TrackKind::Audio => "audio_%u", - } - } -} - -#[derive(Debug, Clone)] -struct TrackDescriptor { - kind: TrackKind, - name: String, -} - -impl TrackDescriptor { - fn pad_name(&self) -> String { - match self.kind { - TrackKind::Video => format!("video_{}", self.name), - TrackKind::Audio => format!("audio_{}", self.name), - } - } -} - -#[derive(Debug)] -enum ControlMessage { - CreatePad { - descriptor: TrackDescriptor, - caps: gst::Caps, - reply: oneshot::Sender, - }, - NoMorePads, - ReportError(anyhow::Error), -} - -#[derive(Debug)] -enum PadMessage { - Buffer(gst::Buffer), - Eos, - Drop, -} - -#[derive(Debug, Clone)] -struct PadEndpoint { - sender: mpsc::UnboundedSender, -} - -impl PadEndpoint { - fn send(&self, msg: PadMessage) -> bool { - self.sender.send(msg).is_ok() - } -} - -struct PadHandle { - sender: mpsc::UnboundedSender, - task: glib::JoinHandle<()>, -} - -struct SessionController { - shutdown: 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 control_for_error = control_tx.clone(); - let join = RUNTIME.spawn(async move { - let result = run_session(settings, control_tx, &mut shutdown_rx).await; - if let Err(err) = result { - let _ = control_for_error.send(ControlMessage::ReportError(err)); - } - }); - - Self { - shutdown: shutdown_tx, - join, - } - } - - fn stop(self) { - let _ = self.shutdown.send(true); - RUNTIME.spawn(async move { - if let Err(err) = self.join.await { - gst::warning!(CAT, "session task ended with error: {err:?}"); - } - }); - } -} - -#[derive(Default)] -pub struct MoqSrc { - settings: Mutex, - pads: Mutex>, - control_task: Mutex>>, - control_sender: Mutex>>, - session: Mutex>, -} - -#[glib::object_subclass] -impl ObjectSubclass for MoqSrc { - const NAME: &'static str = "MoqSrc"; - type Type = super::MoqSrc; - type ParentType = gst::Element; - - fn new() -> Self { - Self::default() - } -} - -impl ObjectImpl for MoqSrc { - fn properties() -> &'static [glib::ParamSpec] { - static PROPS: LazyLock> = LazyLock::new(|| { - vec![ - glib::ParamSpecString::builder("url") - .nick("Source URL") - .blurb("Connect to the given URL") - .build(), - glib::ParamSpecString::builder("broadcast") - .nick("Broadcast") - .blurb("The broadcast name to subscribe to") - .build(), - glib::ParamSpecBoolean::builder("tls-disable-verify") - .nick("TLS Disable Verify") - .blurb("Disable TLS certificate verification") - .default_value(false) - .build(), - ] - }); - PROPS.as_ref() - } - - fn set_property(&self, _id: usize, value: &glib::Value, pspec: &glib::ParamSpec) { - let mut settings = self.settings.lock().unwrap(); - match pspec.name() { - "url" => settings.url = value.get().unwrap(), - "broadcast" => settings.broadcast = value.get().unwrap(), - "tls-disable-verify" => settings.tls_disable_verify = value.get().unwrap(), - _ => unreachable!(), - } - } - - fn property(&self, _id: usize, pspec: &glib::ParamSpec) -> glib::Value { - let settings = self.settings.lock().unwrap(); - match pspec.name() { - "url" => settings.url.to_value(), - "broadcast" => settings.broadcast.to_value(), - "tls-disable-verify" => settings.tls_disable_verify.to_value(), - _ => unreachable!(), - } - } -} - -impl GstObjectImpl for MoqSrc {} -impl ElementImpl for MoqSrc { - fn metadata() -> Option<&'static gst::subclass::ElementMetadata> { - static META: LazyLock = LazyLock::new(|| { - gst::subclass::ElementMetadata::new( - "MoQ Src", - "Source/Network/MoQ", - "Receives media over the network via MoQ", - "Luke Curley , Steve McFarlin ", - ) - }); - Some(&*META) - } - - fn pad_templates() -> &'static [gst::PadTemplate] { - static PAD_TEMPLATES: LazyLock> = LazyLock::new(|| { - vec![ - gst::PadTemplate::new( - "video_%u", - gst::PadDirection::Src, - gst::PadPresence::Sometimes, - &gst::Caps::new_any(), - ) - .unwrap(), - gst::PadTemplate::new( - "audio_%u", - gst::PadDirection::Src, - gst::PadPresence::Sometimes, - &gst::Caps::new_any(), - ) - .unwrap(), - ] - }); - PAD_TEMPLATES.as_ref() - } - - fn change_state(&self, transition: gst::StateChange) -> Result { - match transition { - gst::StateChange::ReadyToPaused => { - if let Err(err) = self.start_session() { - gst::error!(CAT, obj = self.obj(), "failed to start session: {err:?}"); - return Err(gst::StateChangeError); - } - let success = self.parent_change_state(transition)?; - let result = match success { - gst::StateChangeSuccess::Async => gst::StateChangeSuccess::Async, - _ => gst::StateChangeSuccess::NoPreroll, - }; - Ok(result) - } - gst::StateChange::PausedToReady => { - self.stop_session(); - self.parent_change_state(transition) - } - _ => self.parent_change_state(transition), - } - } -} - -impl MoqSrc { - fn start_session(&self) -> Result<()> { - let settings = { - let settings = self.settings.lock().unwrap().clone(); - ResolvedSettings::try_from(settings)? - }; - - let (control_tx, control_rx) = mpsc::unbounded_channel(); - let obj = self.obj(); - let weak = obj.downgrade(); - let context = glib::MainContext::default(); - let control_task = spawn_main_context_forwarder(&context, control_rx, move |msg| { - if let Some(obj) = weak.upgrade() { - obj.imp().handle_control_message(msg); - true - } else { - false - } - }); - - *self.control_task.lock().unwrap() = Some(control_task); - *self.control_sender.lock().unwrap() = Some(control_tx.clone()); - - let session = SessionController::start(settings, control_tx); - *self.session.lock().unwrap() = Some(session); - Ok(()) - } - - fn stop_session(&self) { - if let Some(session) = self.session.lock().unwrap().take() { - session.stop(); - } - - if let Some(control_task) = self.control_task.lock().unwrap().take() { - control_task.abort(); - } - - let handles = self.pads.lock().unwrap().drain().collect::>(); - for (name, handle) in handles { - gst::debug!(CAT, "dropping pad {name}"); - let _ = handle.sender.send(PadMessage::Drop); - handle.task.abort(); - } - - *self.control_sender.lock().unwrap() = None; - } - - fn handle_control_message(&self, msg: ControlMessage) { - match msg { - ControlMessage::CreatePad { - descriptor, - caps, - reply, - } => { - if let Err(err) = self.create_pad(descriptor, caps, reply) { - 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:?}"]); - } - } - } - - fn create_pad( - &self, - descriptor: TrackDescriptor, - caps: gst::Caps, - reply: oneshot::Sender, - ) -> Result<()> { - let obj = self.obj(); - let templ = obj - .element_class() - .pad_template(descriptor.kind.template_name()) - .context("missing pad template")?; - - let pad = gst::Pad::builder_from_template(&templ) - .name(descriptor.pad_name()) - .build(); - - pad.set_active(true)?; - - let stream_start = gst::event::StreamStart::builder(&descriptor.name) - .group_id(gst::GroupId::next()) - .build(); - pad.push_event(stream_start); - pad.push_event(gst::event::Caps::new(&caps)); - pad.push_event(gst::event::Segment::new(&gst::FormattedSegment::::new())); - - obj.add_pad(&pad)?; - - let (pad_tx, pad_rx) = mpsc::unbounded_channel(); - let pad_clone = pad.clone(); - let weak = obj.downgrade(); - let context = glib::MainContext::default(); - let task = spawn_main_context_forwarder(&context, pad_rx, move |msg| { - if let Some(obj) = weak.upgrade() { - let imp = obj.imp(); - imp.dispatch_pad_message(&pad_clone, msg) - } else { - false - } - }); - - self.pads.lock().unwrap().insert( - descriptor.pad_name(), - PadHandle { - sender: pad_tx.clone(), - task, - }, - ); - - let _ = reply.send(PadEndpoint { sender: pad_tx }); - Ok(()) - } - - fn dispatch_pad_message(&self, pad: &gst::Pad, msg: PadMessage) -> bool { - match msg { - PadMessage::Buffer(buffer) => { - if let Err(err) = pad.push(buffer) { - gst::warning!(CAT, "failed to push buffer: {err:?}"); - return false; - } - true - } - PadMessage::Eos => { - pad.push_event(gst::event::Eos::builder().build()); - true - } - PadMessage::Drop => { - let _ = pad.set_active(false); - let _ = self.obj().remove_pad(pad); - false - } - } - } -} - -async fn run_session( - settings: ResolvedSettings, - 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?; - - // Wait for the broadcast to be announced. Synchronous lookup would race the gossip of - // announcements that happens after the session is established. - tracing::info!(broadcast = %settings.broadcast, "waiting for broadcast to be announced"); - let broadcast = tokio::select! { - broadcast = origin_consumer.announced_broadcast(&settings.broadcast) => broadcast - .context("broadcast not allowed or origin closed")?, - _ = shutdown.changed() => return Ok(()), - }; - - 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 tasks = Vec::new(); - - for (track_name, config) in catalog.video.renditions { - let descriptor = TrackDescriptor { - kind: TrackKind::Video, - name: track_name.clone(), - }; - let caps = video_caps(&config)?; - let endpoint = request_pad(&control_tx, descriptor.clone(), caps).await?; - 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())); - } - - for (track_name, config) in catalog.audio.renditions { - let descriptor = TrackDescriptor { - kind: TrackKind::Audio, - name: track_name.clone(), - }; - let caps = audio_caps(&config)?; - let endpoint = request_pad(&control_tx, descriptor.clone(), caps).await?; - 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())); - } - - let _ = control_tx.send(ControlMessage::NoMorePads); - - for task in tasks { - let _ = task.await; - } - - Ok(()) -} - -async fn request_pad( - control_tx: &mpsc::UnboundedSender, - descriptor: TrackDescriptor, - caps: gst::Caps, -) -> Result { - let (reply_tx, reply_rx) = oneshot::channel(); - control_tx - .send(ControlMessage::CreatePad { - descriptor, - caps, - reply: reply_tx, - }) - .map_err(|_| anyhow::anyhow!("control plane shut down"))?; - - let endpoint = reply_rx.await.context("pad creation cancelled")?; - Ok(endpoint) -} - -fn spawn_track_pump( - track: moq_mux::container::Consumer, - descriptor: TrackDescriptor, - pad_endpoint: PadEndpoint, - shutdown: watch::Receiver, -) -> tokio::task::JoinHandle<()> { - RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown)) -} - -async fn run_track_pump( - mut track: moq_mux::container::Consumer, - descriptor: TrackDescriptor, - pad_endpoint: PadEndpoint, - mut shutdown: watch::Receiver, -) { - let mut reference_ts = None; - loop { - tokio::select! { - _ = shutdown.changed() => { - pad_endpoint.send(PadMessage::Drop); - break; - } - frame = track.read() => { - match frame { - Ok(Some(frame)) => { - let timestamp = frame.timestamp; - let is_keyframe = frame.keyframe; - let payload = frame.payload; - let mut buffer = gst::Buffer::from_slice(payload.to_vec()); - let buffer_mut = buffer.get_mut().unwrap(); - - let pts = match reference_ts { - Some(reference) => { - let delta: Duration = (timestamp - reference).into(); - gst::ClockTime::from_nseconds(delta.as_nanos() as u64) - } - None => { - reference_ts = Some(timestamp); - gst::ClockTime::ZERO - } - }; - buffer_mut.set_pts(Some(pts)); - - let mut flags = buffer_mut.flags(); - match descriptor.kind { - TrackKind::Video => { - if is_keyframe { - flags.remove(gst::BufferFlags::DELTA_UNIT); - } else { - flags.insert(gst::BufferFlags::DELTA_UNIT); - } - } - TrackKind::Audio => { - flags.remove(gst::BufferFlags::DELTA_UNIT); - } - } - buffer_mut.set_flags(flags); - - if !pad_endpoint.send(PadMessage::Buffer(buffer)) { - break; - } - } - Ok(None) => { - pad_endpoint.send(PadMessage::Eos); - pad_endpoint.send(PadMessage::Drop); - break; - } - Err(err) => { - gst::warning!(CAT, "track {} failed: {err:?}", descriptor.name); - pad_endpoint.send(PadMessage::Drop); - break; - } - } - } - } - } -} - -fn video_caps(config: &hang::catalog::VideoConfig) -> Result { - use hang::catalog::VideoCodec; - - let caps = match &config.codec { - VideoCodec::H264(_) => { - let mut builder = gst::Caps::builder("video/x-h264").field("alignment", "au"); - if let Some(description) = &config.description { - builder = builder - .field("stream-format", "avc") - .field("codec_data", gst::Buffer::from_slice(description.clone())); - } else { - builder = builder.field("stream-format", "annexb"); - } - builder.build() - } - VideoCodec::H265(h265) => { - let mut builder = gst::Caps::builder("video/x-h265").field("alignment", "au"); - match &config.description { - Some(description) => { - let format = if h265.in_band { "hev1" } else { "hvc1" }; - builder = builder - .field("stream-format", format) - .field("codec_data", gst::Buffer::from_slice(description.clone())); - } - None => { - let format = if h265.in_band { "hev1" } else { "byte-stream" }; - builder = builder.field("stream-format", format); - } - } - builder.build() - } - VideoCodec::AV1(_) => { - let mut builder = gst::Caps::builder("video/x-av1"); - if let Some(description) = &config.description { - builder = builder.field("codec_data", gst::Buffer::from_slice(description.clone())); - } - builder.build() - } - other => bail!("unsupported video codec: {other:?}"), - }; - Ok(caps) -} - -fn audio_caps(config: &hang::catalog::AudioConfig) -> Result { - let caps = match &config.codec { - hang::catalog::AudioCodec::AAC(_) => { - let mut builder = gst::Caps::builder("audio/mpeg") - .field("mpegversion", 4) - .field("rate", config.sample_rate) - .field("channels", config.channel_count); - if let Some(description) = &config.description { - builder = builder - .field("codec_data", gst::Buffer::from_slice(description.clone())) - .field("stream-format", "aac"); - } else { - builder = builder.field("stream-format", "adts"); - } - builder.build() - } - hang::catalog::AudioCodec::Opus => { - let mut builder = gst::Caps::builder("audio/x-opus") - .field("rate", config.sample_rate) - .field("channels", config.channel_count); - if let Some(description) = &config.description { - builder = builder - .field("codec_data", gst::Buffer::from_slice(description.clone())) - .field("stream-format", "ogg"); - } - builder.build() - } - other => bail!("unsupported audio codec: {other:?}"), - }; - Ok(caps) -} - -fn spawn_main_context_forwarder( - context: &glib::MainContext, - mut rx: mpsc::UnboundedReceiver, - mut handler: F, -) -> glib::JoinHandle<()> -where - T: Send + 'static, - F: FnMut(T) -> bool + 'static, -{ - let ctx = context.clone(); - ctx.spawn_local(async move { - while let Some(msg) = rx.recv().await { - if !handler(msg) { - break; - } - } - }) -} +use std::collections::HashMap; +use std::sync::{LazyLock, Mutex}; +use std::time::Duration; + +use anyhow::{bail, Context, Result}; +use gst::glib; +use gst::prelude::*; +use gst::subclass::prelude::*; +use tokio::sync::{mpsc, oneshot, watch}; + +use hang::moq_net; + +static CAT: LazyLock = + LazyLock::new(|| gst::DebugCategory::new("moq-src", gst::DebugColorFlags::empty(), Some("MoQ Source Element"))); + +static RUNTIME: LazyLock = LazyLock::new(|| { + tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build() + .expect("spawn tokio runtime") +}); + +#[derive(Debug, Clone)] +struct Settings { + url: Option, + broadcast: Option, + tls_disable_verify: bool, + max_latency_ms: u64, + ascending: bool, +} + +impl Default for Settings { + fn default() -> Self { + Self { + url: None, + broadcast: None, + tls_disable_verify: false, + max_latency_ms: 1000, + ascending: false, + } + } +} + +#[derive(Debug, Clone)] +struct ResolvedSettings { + url: url::Url, + broadcast: String, + tls_disable_verify: bool, + max_latency_ms: u64, + ascending: bool, +} + +impl TryFrom for ResolvedSettings { + type Error = anyhow::Error; + + fn try_from(value: Settings) -> Result { + 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(), + tls_disable_verify: value.tls_disable_verify, + max_latency_ms: value.max_latency_ms, + ascending: value.ascending, + }) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +enum TrackKind { + Video, + Audio, +} + +impl TrackKind { + fn template_name(&self) -> &'static str { + match self { + TrackKind::Video => "video_%u", + TrackKind::Audio => "audio_%u", + } + } +} + +#[derive(Debug, Clone)] +struct TrackDescriptor { + kind: TrackKind, + name: String, +} + +impl TrackDescriptor { + fn pad_name(&self) -> String { + match self.kind { + TrackKind::Video => format!("video_{}", self.name), + TrackKind::Audio => format!("audio_{}", self.name), + } + } +} + +#[derive(Debug)] +enum ControlMessage { + CreatePad { + descriptor: TrackDescriptor, + caps: gst::Caps, + reply: oneshot::Sender, + }, + NoMorePads, + ReportError(anyhow::Error), +} + +#[derive(Debug)] +enum PadMessage { + Buffer(gst::Buffer), + Eos, + Drop, +} + +#[derive(Debug, Clone)] +struct PadEndpoint { + sender: mpsc::UnboundedSender, +} + +impl PadEndpoint { + fn send(&self, msg: PadMessage) -> bool { + self.sender.send(msg).is_ok() + } +} + +struct PadHandle { + sender: mpsc::UnboundedSender, + task: glib::JoinHandle<()>, +} + +struct SessionController { + shutdown: 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 control_for_error = control_tx.clone(); + let join = RUNTIME.spawn(async move { + let result = run_session(settings, control_tx, &mut shutdown_rx).await; + if let Err(err) = result { + let _ = control_for_error.send(ControlMessage::ReportError(err)); + } + }); + + Self { + shutdown: shutdown_tx, + join, + } + } + + fn stop(self) { + let _ = self.shutdown.send(true); + RUNTIME.spawn(async move { + if let Err(err) = self.join.await { + gst::warning!(CAT, "session task ended with error: {err:?}"); + } + }); + } +} + +#[derive(Default)] +pub struct MoqSrc { + settings: Mutex, + pads: Mutex>, + control_task: Mutex>>, + control_sender: Mutex>>, + session: Mutex>, +} + +#[glib::object_subclass] +impl ObjectSubclass for MoqSrc { + const NAME: &'static str = "MoqSrc"; + type Type = super::MoqSrc; + type ParentType = gst::Element; + + fn new() -> Self { + Self::default() + } +} + +impl ObjectImpl for MoqSrc { + fn properties() -> &'static [glib::ParamSpec] { + static PROPS: LazyLock> = LazyLock::new(|| { + vec![ + glib::ParamSpecString::builder("url") + .nick("Source URL") + .blurb("Connect to the given URL") + .build(), + glib::ParamSpecString::builder("broadcast") + .nick("Broadcast") + .blurb("The broadcast name to subscribe to") + .build(), + glib::ParamSpecBoolean::builder("tls-disable-verify") + .nick("TLS Disable Verify") + .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 to stay at the live edge. Ignored when ascending=true.") + .default_value(1000) + .build(), + glib::ParamSpecBoolean::builder("ascending") + .nick("Ascending") + .blurb("Deliver groups oldest-first; default false delivers newest-first (live edge)") + .default_value(false) + .build(), + ] + }); + PROPS.as_ref() + } + + fn set_property(&self, _id: usize, value: &glib::Value, pspec: &glib::ParamSpec) { + let mut settings = self.settings.lock().unwrap(); + match pspec.name() { + "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(), + "ascending" => settings.ascending = value.get().unwrap(), + _ => unreachable!(), + } + } + + fn property(&self, _id: usize, pspec: &glib::ParamSpec) -> glib::Value { + let settings = self.settings.lock().unwrap(); + match pspec.name() { + "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(), + "ascending" => settings.ascending.to_value(), + _ => unreachable!(), + } + } +} + +impl GstObjectImpl for MoqSrc {} +impl ElementImpl for MoqSrc { + fn metadata() -> Option<&'static gst::subclass::ElementMetadata> { + static META: LazyLock = LazyLock::new(|| { + gst::subclass::ElementMetadata::new( + "MoQ Src", + "Source/Network/MoQ", + "Receives media over the network via MoQ", + "Luke Curley , Steve McFarlin ", + ) + }); + Some(&*META) + } + + fn pad_templates() -> &'static [gst::PadTemplate] { + static PAD_TEMPLATES: LazyLock> = LazyLock::new(|| { + vec![ + gst::PadTemplate::new( + "video_%u", + gst::PadDirection::Src, + gst::PadPresence::Sometimes, + &gst::Caps::new_any(), + ) + .unwrap(), + gst::PadTemplate::new( + "audio_%u", + gst::PadDirection::Src, + gst::PadPresence::Sometimes, + &gst::Caps::new_any(), + ) + .unwrap(), + ] + }); + PAD_TEMPLATES.as_ref() + } + + fn change_state(&self, transition: gst::StateChange) -> Result { + match transition { + gst::StateChange::ReadyToPaused => { + if let Err(err) = self.start_session() { + gst::error!(CAT, obj = self.obj(), "failed to start session: {err:?}"); + return Err(gst::StateChangeError); + } + let success = self.parent_change_state(transition)?; + let result = match success { + gst::StateChangeSuccess::Async => gst::StateChangeSuccess::Async, + _ => gst::StateChangeSuccess::NoPreroll, + }; + Ok(result) + } + gst::StateChange::PausedToReady => { + self.stop_session(); + self.parent_change_state(transition) + } + _ => self.parent_change_state(transition), + } + } +} + +impl MoqSrc { + fn start_session(&self) -> Result<()> { + let settings = { + let settings = self.settings.lock().unwrap().clone(); + ResolvedSettings::try_from(settings)? + }; + + let (control_tx, control_rx) = mpsc::unbounded_channel(); + let obj = self.obj(); + let weak = obj.downgrade(); + let context = glib::MainContext::default(); + let control_task = spawn_main_context_forwarder(&context, control_rx, move |msg| { + if let Some(obj) = weak.upgrade() { + obj.imp().handle_control_message(msg); + true + } else { + false + } + }); + + *self.control_task.lock().unwrap() = Some(control_task); + *self.control_sender.lock().unwrap() = Some(control_tx.clone()); + + let session = SessionController::start(settings, control_tx); + *self.session.lock().unwrap() = Some(session); + Ok(()) + } + + fn stop_session(&self) { + if let Some(session) = self.session.lock().unwrap().take() { + session.stop(); + } + + if let Some(control_task) = self.control_task.lock().unwrap().take() { + control_task.abort(); + } + + let handles = self.pads.lock().unwrap().drain().collect::>(); + for (name, handle) in handles { + gst::debug!(CAT, "dropping pad {name}"); + let _ = handle.sender.send(PadMessage::Drop); + handle.task.abort(); + } + + *self.control_sender.lock().unwrap() = None; + } + + fn handle_control_message(&self, msg: ControlMessage) { + match msg { + ControlMessage::CreatePad { + descriptor, + caps, + reply, + } => { + if let Err(err) = self.create_pad(descriptor, caps, reply) { + 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:?}"]); + } + } + } + + fn create_pad( + &self, + descriptor: TrackDescriptor, + caps: gst::Caps, + reply: oneshot::Sender, + ) -> Result<()> { + let obj = self.obj(); + let templ = obj + .element_class() + .pad_template(descriptor.kind.template_name()) + .context("missing pad template")?; + + let pad = gst::Pad::builder_from_template(&templ) + .name(descriptor.pad_name()) + .build(); + + pad.set_active(true)?; + + let stream_start = gst::event::StreamStart::builder(&descriptor.name) + .group_id(gst::GroupId::next()) + .build(); + pad.push_event(stream_start); + pad.push_event(gst::event::Caps::new(&caps)); + pad.push_event(gst::event::Segment::new(&gst::FormattedSegment::::new())); + + obj.add_pad(&pad)?; + + let (pad_tx, pad_rx) = mpsc::unbounded_channel(); + let pad_clone = pad.clone(); + let weak = obj.downgrade(); + let context = glib::MainContext::default(); + let task = spawn_main_context_forwarder(&context, pad_rx, move |msg| { + if let Some(obj) = weak.upgrade() { + let imp = obj.imp(); + imp.dispatch_pad_message(&pad_clone, msg) + } else { + false + } + }); + + self.pads.lock().unwrap().insert( + descriptor.pad_name(), + PadHandle { + sender: pad_tx.clone(), + task, + }, + ); + + let _ = reply.send(PadEndpoint { sender: pad_tx }); + Ok(()) + } + + fn dispatch_pad_message(&self, pad: &gst::Pad, msg: PadMessage) -> bool { + match msg { + PadMessage::Buffer(buffer) => { + if let Err(err) = pad.push(buffer) { + gst::warning!(CAT, "failed to push buffer: {err:?}"); + return false; + } + true + } + PadMessage::Eos => { + pad.push_event(gst::event::Eos::builder().build()); + true + } + PadMessage::Drop => { + let _ = pad.set_active(false); + let _ = self.obj().remove_pad(pad); + false + } + } + } +} + +async fn run_session( + settings: ResolvedSettings, + 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?; + + // Wait for the broadcast to be announced. Synchronous lookup would race the gossip of + // announcements that happens after the session is established. + tracing::info!(broadcast = %settings.broadcast, "waiting for broadcast to be announced"); + let broadcast = tokio::select! { + broadcast = origin_consumer.announced_broadcast(&settings.broadcast) => broadcast + .context("broadcast not allowed or origin closed")?, + _ = shutdown.changed() => return Ok(()), + }; + + 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 consumer_latency = if settings.ascending { + Duration::MAX + } else { + Duration::from_millis(settings.max_latency_ms) + }; + + 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())); + } + + 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())); + } + + let _ = control_tx.send(ControlMessage::NoMorePads); + + for task in tasks { + let _ = task.await; + } + + Ok(()) +} + +fn subscribe_track( + broadcast: &moq_net::BroadcastConsumer, + name: &str, + ascending: bool, + latency: Duration, +) -> Result> { + #[allow(unused_mut)] + let mut track_ref = moq_net::Track::new(name); + #[cfg(feature = "group-order")] + { track_ref.ordered = ascending; } + #[cfg(not(feature = "group-order"))] + { let _ = ascending; } + let consumer = broadcast.subscribe_track(&track_ref)?; + Ok(moq_mux::container::Consumer::new(consumer, moq_mux::catalog::hang::Container::Legacy) + .with_latency(latency)) +} + +async fn request_pad( + control_tx: &mpsc::UnboundedSender, + descriptor: TrackDescriptor, + caps: gst::Caps, +) -> Result { + let (reply_tx, reply_rx) = oneshot::channel(); + control_tx + .send(ControlMessage::CreatePad { + descriptor, + caps, + reply: reply_tx, + }) + .map_err(|_| anyhow::anyhow!("control plane shut down"))?; + + let endpoint = reply_rx.await.context("pad creation cancelled")?; + Ok(endpoint) +} + +fn spawn_track_pump( + track: moq_mux::container::Consumer, + descriptor: TrackDescriptor, + pad_endpoint: PadEndpoint, + shutdown: watch::Receiver, +) -> tokio::task::JoinHandle<()> { + RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown)) +} + +async fn run_track_pump( + mut track: moq_mux::container::Consumer, + descriptor: TrackDescriptor, + pad_endpoint: PadEndpoint, + mut shutdown: watch::Receiver, +) { + let mut reference_ts = None; + loop { + tokio::select! { + _ = shutdown.changed() => { + pad_endpoint.send(PadMessage::Drop); + break; + } + frame = track.read() => { + match frame { + Ok(Some(frame)) => { + let timestamp = frame.timestamp; + let is_keyframe = frame.keyframe; + let payload = frame.payload; + let mut buffer = gst::Buffer::from_slice(payload.to_vec()); + let buffer_mut = buffer.get_mut().unwrap(); + + let pts = match reference_ts { + Some(reference) => { + let delta: Duration = (timestamp - reference).into(); + gst::ClockTime::from_nseconds(delta.as_nanos() as u64) + } + None => { + reference_ts = Some(timestamp); + gst::ClockTime::ZERO + } + }; + buffer_mut.set_pts(Some(pts)); + + let mut flags = buffer_mut.flags(); + match descriptor.kind { + TrackKind::Video => { + if is_keyframe { + flags.remove(gst::BufferFlags::DELTA_UNIT); + } else { + flags.insert(gst::BufferFlags::DELTA_UNIT); + } + } + TrackKind::Audio => { + flags.remove(gst::BufferFlags::DELTA_UNIT); + } + } + buffer_mut.set_flags(flags); + + if !pad_endpoint.send(PadMessage::Buffer(buffer)) { + break; + } + } + Ok(None) => { + pad_endpoint.send(PadMessage::Eos); + pad_endpoint.send(PadMessage::Drop); + break; + } + Err(err) => { + gst::warning!(CAT, "track {} failed: {err:?}", descriptor.name); + pad_endpoint.send(PadMessage::Drop); + break; + } + } + } + } + } +} + +fn video_caps(config: &hang::catalog::VideoConfig) -> Result { + use hang::catalog::VideoCodec; + + let caps = match &config.codec { + VideoCodec::H264(_) => { + let mut builder = gst::Caps::builder("video/x-h264").field("alignment", "au"); + if let Some(description) = &config.description { + builder = builder + .field("stream-format", "avc") + .field("codec_data", gst::Buffer::from_slice(description.clone())); + } else { + builder = builder.field("stream-format", "annexb"); + } + builder.build() + } + VideoCodec::H265(h265) => { + let mut builder = gst::Caps::builder("video/x-h265").field("alignment", "au"); + match &config.description { + Some(description) => { + let format = if h265.in_band { "hev1" } else { "hvc1" }; + builder = builder + .field("stream-format", format) + .field("codec_data", gst::Buffer::from_slice(description.clone())); + } + None => { + let format = if h265.in_band { "hev1" } else { "byte-stream" }; + builder = builder.field("stream-format", format); + } + } + builder.build() + } + VideoCodec::AV1(_) => { + let mut builder = gst::Caps::builder("video/x-av1"); + if let Some(description) = &config.description { + builder = builder.field("codec_data", gst::Buffer::from_slice(description.clone())); + } + builder.build() + } + other => bail!("unsupported video codec: {other:?}"), + }; + Ok(caps) +} + +fn audio_caps(config: &hang::catalog::AudioConfig) -> Result { + let caps = match &config.codec { + hang::catalog::AudioCodec::AAC(_) => { + let mut builder = gst::Caps::builder("audio/mpeg") + .field("mpegversion", 4) + .field("rate", config.sample_rate) + .field("channels", config.channel_count); + if let Some(description) = &config.description { + builder = builder + .field("codec_data", gst::Buffer::from_slice(description.clone())) + .field("stream-format", "aac"); + } else { + builder = builder.field("stream-format", "adts"); + } + builder.build() + } + hang::catalog::AudioCodec::Opus => { + let mut builder = gst::Caps::builder("audio/x-opus") + .field("rate", config.sample_rate) + .field("channels", config.channel_count); + if let Some(description) = &config.description { + builder = builder + .field("codec_data", gst::Buffer::from_slice(description.clone())) + .field("stream-format", "ogg"); + } + builder.build() + } + other => bail!("unsupported audio codec: {other:?}"), + }; + Ok(caps) +} + +fn spawn_main_context_forwarder( + context: &glib::MainContext, + mut rx: mpsc::UnboundedReceiver, + mut handler: F, +) -> glib::JoinHandle<()> +where + T: Send + 'static, + F: FnMut(T) -> bool + 'static, +{ + let ctx = context.clone(); + ctx.spawn_local(async move { + while let Some(msg) = rx.recv().await { + if !handler(msg) { + break; + } + } + }) +} From b28cf6ca8f6092fab16d469bb041773b5af5baeb Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Thu, 11 Jun 2026 11:54:59 -0300 Subject: [PATCH 2/5] feat(moq-gst): decouple ascending ordering from latency budget ascending and max_latency_ms are now orthogonal: the consumer always uses Duration::from_millis(max_latency_ms) regardless of ascending mode. The subscriber sets G_MAXUINT64 as the default when ascending=true and no explicit MOQ_MAX_LATENCY_MS is provided, preserving the unlimited-budget behaviour of bare moq-ascending while allowing moq-125-asc and similar variants to apply a real budget with ascending group ordering. Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/src/source/imp.rs | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index a4abf5591e..093a306c78 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -468,11 +468,7 @@ async fn run_session( let mut catalog = moq_mux::catalog::hang::Consumer::new(catalog_track); let catalog = catalog.next().await?.context("catalog missing")?.clone(); - let consumer_latency = if settings.ascending { - Duration::MAX - } else { - Duration::from_millis(settings.max_latency_ms) - }; + let consumer_latency = Duration::from_millis(settings.max_latency_ms); let mut tasks = Vec::new(); From bed380f19f95b1a0cb930f667ce47e05b1eb8c56 Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Fri, 12 Jun 2026 12:50:53 -0300 Subject: [PATCH 3/5] fix(moq-gst): add mutable_ready to properties and warn when ascending is no-op Add .mutable_ready() to all five moqsrc property builders so settings can be changed in READY state. Without this, writes after ReadyToPaused are silently accepted but never applied to the running session. Emit a gst::warning! when ascending=true is set but the group-order feature is not compiled in, so callers get an explicit signal instead of silent wrong behavior. Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/src/source/imp.rs | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index 093a306c78..f78a8eb918 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -191,25 +191,30 @@ impl ObjectImpl for MoqSrc { glib::ParamSpecString::builder("url") .nick("Source URL") .blurb("Connect to the given URL") + .mutable_ready() .build(), glib::ParamSpecString::builder("broadcast") .nick("Broadcast") .blurb("The broadcast name to subscribe to") + .mutable_ready() .build(), glib::ParamSpecBoolean::builder("tls-disable-verify") .nick("TLS Disable Verify") .blurb("Disable TLS certificate verification") .default_value(false) + .mutable_ready() .build(), glib::ParamSpecUInt64::builder("max-latency-ms") .nick("Max latency (ms)") .blurb("Drop groups older than this to stay at the live edge. Ignored when ascending=true.") .default_value(1000) + .mutable_ready() .build(), glib::ParamSpecBoolean::builder("ascending") .nick("Ascending") .blurb("Deliver groups oldest-first; default false delivers newest-first (live edge)") .default_value(false) + .mutable_ready() .build(), ] }); @@ -223,7 +228,13 @@ impl ObjectImpl for MoqSrc { "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(), - "ascending" => settings.ascending = value.get().unwrap(), + "ascending" => { + settings.ascending = value.get().unwrap(); + #[cfg(not(feature = "group-order"))] + if settings.ascending { + gst::warning!(CAT, "ascending=true has no effect: build with --features group-order"); + } + } _ => unreachable!(), } } From 581ec5fc3cc156e38c2e7ea670b27423421074d8 Mon Sep 17 00:00:00 2001 From: santi-ferreiro Date: Fri, 12 Jun 2026 16:39:09 -0300 Subject: [PATCH 4/5] Update rs/moq-gst/src/source/imp.rs Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- rs/moq-gst/src/source/imp.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index f78a8eb918..ec908f3e85 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -206,7 +206,7 @@ impl ObjectImpl for MoqSrc { .build(), glib::ParamSpecUInt64::builder("max-latency-ms") .nick("Max latency (ms)") - .blurb("Drop groups older than this to stay at the live edge. Ignored when ascending=true.") + .blurb("Drop groups older than this to stay at the live edge.") .default_value(1000) .mutable_ready() .build(), From 0da76156c7fae19689585edd6f80e68c011d3c3b Mon Sep 17 00:00:00 2001 From: Santiago Ferreiro Date: Fri, 12 Jun 2026 17:04:32 -0300 Subject: [PATCH 5/5] chore(moq-gst): normalize line endings to LF Co-Authored-By: Claude Sonnet 4.6 --- rs/moq-gst/Cargo.toml | 74 +- rs/moq-gst/src/source/imp.rs | 1428 +++++++++++++++++----------------- 2 files changed, 751 insertions(+), 751 deletions(-) diff --git a/rs/moq-gst/Cargo.toml b/rs/moq-gst/Cargo.toml index 765d24a5a5..398ff30dc0 100644 --- a/rs/moq-gst/Cargo.toml +++ b/rs/moq-gst/Cargo.toml @@ -1,37 +1,37 @@ -[package] -name = "moq-gst" -description = "Media over QUIC - GStreamer plugin" -authors = ["Luke Curley"] -repository = "https://github.com/kixelated/moq" -license = "MIT OR Apache-2.0" - -version = "0.2.2" -edition = "2021" -rust-version.workspace = true -publish = false - -[lib] -name = "gstmoq" -crate-type = ["cdylib", "rlib"] -test = false -doctest = false - -[features] -group-order = [] - -[dependencies] - -anyhow = { version = "1", features = ["backtrace"] } -bytes = "1" -gst = { package = "gstreamer", version = "0.23" } -hang = { workspace = true } -moq-mux = { workspace = true } -moq-native = { workspace = true, default-features = true } -moq-net = { workspace = true } -tokio = { workspace = true, features = ["full"] } -tracing = "0.1" -tracing-subscriber = "0.3" -url = "2" - -[build-dependencies] -gst-plugin-version-helper = "0.8" +[package] +name = "moq-gst" +description = "Media over QUIC - GStreamer plugin" +authors = ["Luke Curley"] +repository = "https://github.com/kixelated/moq" +license = "MIT OR Apache-2.0" + +version = "0.2.2" +edition = "2021" +rust-version.workspace = true +publish = false + +[lib] +name = "gstmoq" +crate-type = ["cdylib", "rlib"] +test = false +doctest = false + +[features] +group-order = [] + +[dependencies] + +anyhow = { version = "1", features = ["backtrace"] } +bytes = "1" +gst = { package = "gstreamer", version = "0.23" } +hang = { workspace = true } +moq-mux = { workspace = true } +moq-native = { workspace = true, default-features = true } +moq-net = { workspace = true } +tokio = { workspace = true, features = ["full"] } +tracing = "0.1" +tracing-subscriber = "0.3" +url = "2" + +[build-dependencies] +gst-plugin-version-helper = "0.8" diff --git a/rs/moq-gst/src/source/imp.rs b/rs/moq-gst/src/source/imp.rs index ec908f3e85..f97e4f5053 100644 --- a/rs/moq-gst/src/source/imp.rs +++ b/rs/moq-gst/src/source/imp.rs @@ -1,714 +1,714 @@ -use std::collections::HashMap; -use std::sync::{LazyLock, Mutex}; -use std::time::Duration; - -use anyhow::{bail, Context, Result}; -use gst::glib; -use gst::prelude::*; -use gst::subclass::prelude::*; -use tokio::sync::{mpsc, oneshot, watch}; - -use hang::moq_net; - -static CAT: LazyLock = - LazyLock::new(|| gst::DebugCategory::new("moq-src", gst::DebugColorFlags::empty(), Some("MoQ Source Element"))); - -static RUNTIME: LazyLock = LazyLock::new(|| { - tokio::runtime::Builder::new_multi_thread() - .enable_all() - .build() - .expect("spawn tokio runtime") -}); - -#[derive(Debug, Clone)] -struct Settings { - url: Option, - broadcast: Option, - tls_disable_verify: bool, - max_latency_ms: u64, - ascending: bool, -} - -impl Default for Settings { - fn default() -> Self { - Self { - url: None, - broadcast: None, - tls_disable_verify: false, - max_latency_ms: 1000, - ascending: false, - } - } -} - -#[derive(Debug, Clone)] -struct ResolvedSettings { - url: url::Url, - broadcast: String, - tls_disable_verify: bool, - max_latency_ms: u64, - ascending: bool, -} - -impl TryFrom for ResolvedSettings { - type Error = anyhow::Error; - - fn try_from(value: Settings) -> Result { - 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(), - tls_disable_verify: value.tls_disable_verify, - max_latency_ms: value.max_latency_ms, - ascending: value.ascending, - }) - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -enum TrackKind { - Video, - Audio, -} - -impl TrackKind { - fn template_name(&self) -> &'static str { - match self { - TrackKind::Video => "video_%u", - TrackKind::Audio => "audio_%u", - } - } -} - -#[derive(Debug, Clone)] -struct TrackDescriptor { - kind: TrackKind, - name: String, -} - -impl TrackDescriptor { - fn pad_name(&self) -> String { - match self.kind { - TrackKind::Video => format!("video_{}", self.name), - TrackKind::Audio => format!("audio_{}", self.name), - } - } -} - -#[derive(Debug)] -enum ControlMessage { - CreatePad { - descriptor: TrackDescriptor, - caps: gst::Caps, - reply: oneshot::Sender, - }, - NoMorePads, - ReportError(anyhow::Error), -} - -#[derive(Debug)] -enum PadMessage { - Buffer(gst::Buffer), - Eos, - Drop, -} - -#[derive(Debug, Clone)] -struct PadEndpoint { - sender: mpsc::UnboundedSender, -} - -impl PadEndpoint { - fn send(&self, msg: PadMessage) -> bool { - self.sender.send(msg).is_ok() - } -} - -struct PadHandle { - sender: mpsc::UnboundedSender, - task: glib::JoinHandle<()>, -} - -struct SessionController { - shutdown: 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 control_for_error = control_tx.clone(); - let join = RUNTIME.spawn(async move { - let result = run_session(settings, control_tx, &mut shutdown_rx).await; - if let Err(err) = result { - let _ = control_for_error.send(ControlMessage::ReportError(err)); - } - }); - - Self { - shutdown: shutdown_tx, - join, - } - } - - fn stop(self) { - let _ = self.shutdown.send(true); - RUNTIME.spawn(async move { - if let Err(err) = self.join.await { - gst::warning!(CAT, "session task ended with error: {err:?}"); - } - }); - } -} - -#[derive(Default)] -pub struct MoqSrc { - settings: Mutex, - pads: Mutex>, - control_task: Mutex>>, - control_sender: Mutex>>, - session: Mutex>, -} - -#[glib::object_subclass] -impl ObjectSubclass for MoqSrc { - const NAME: &'static str = "MoqSrc"; - type Type = super::MoqSrc; - type ParentType = gst::Element; - - fn new() -> Self { - Self::default() - } -} - -impl ObjectImpl for MoqSrc { - fn properties() -> &'static [glib::ParamSpec] { - static PROPS: LazyLock> = LazyLock::new(|| { - vec![ - glib::ParamSpecString::builder("url") - .nick("Source URL") - .blurb("Connect to the given URL") - .mutable_ready() - .build(), - glib::ParamSpecString::builder("broadcast") - .nick("Broadcast") - .blurb("The broadcast name to subscribe to") - .mutable_ready() - .build(), - glib::ParamSpecBoolean::builder("tls-disable-verify") - .nick("TLS Disable Verify") - .blurb("Disable TLS certificate verification") - .default_value(false) - .mutable_ready() - .build(), - glib::ParamSpecUInt64::builder("max-latency-ms") - .nick("Max latency (ms)") - .blurb("Drop groups older than this to stay at the live edge.") - .default_value(1000) - .mutable_ready() - .build(), - glib::ParamSpecBoolean::builder("ascending") - .nick("Ascending") - .blurb("Deliver groups oldest-first; default false delivers newest-first (live edge)") - .default_value(false) - .mutable_ready() - .build(), - ] - }); - PROPS.as_ref() - } - - fn set_property(&self, _id: usize, value: &glib::Value, pspec: &glib::ParamSpec) { - let mut settings = self.settings.lock().unwrap(); - match pspec.name() { - "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(), - "ascending" => { - settings.ascending = value.get().unwrap(); - #[cfg(not(feature = "group-order"))] - if settings.ascending { - gst::warning!(CAT, "ascending=true has no effect: build with --features group-order"); - } - } - _ => unreachable!(), - } - } - - fn property(&self, _id: usize, pspec: &glib::ParamSpec) -> glib::Value { - let settings = self.settings.lock().unwrap(); - match pspec.name() { - "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(), - "ascending" => settings.ascending.to_value(), - _ => unreachable!(), - } - } -} - -impl GstObjectImpl for MoqSrc {} -impl ElementImpl for MoqSrc { - fn metadata() -> Option<&'static gst::subclass::ElementMetadata> { - static META: LazyLock = LazyLock::new(|| { - gst::subclass::ElementMetadata::new( - "MoQ Src", - "Source/Network/MoQ", - "Receives media over the network via MoQ", - "Luke Curley , Steve McFarlin ", - ) - }); - Some(&*META) - } - - fn pad_templates() -> &'static [gst::PadTemplate] { - static PAD_TEMPLATES: LazyLock> = LazyLock::new(|| { - vec![ - gst::PadTemplate::new( - "video_%u", - gst::PadDirection::Src, - gst::PadPresence::Sometimes, - &gst::Caps::new_any(), - ) - .unwrap(), - gst::PadTemplate::new( - "audio_%u", - gst::PadDirection::Src, - gst::PadPresence::Sometimes, - &gst::Caps::new_any(), - ) - .unwrap(), - ] - }); - PAD_TEMPLATES.as_ref() - } - - fn change_state(&self, transition: gst::StateChange) -> Result { - match transition { - gst::StateChange::ReadyToPaused => { - if let Err(err) = self.start_session() { - gst::error!(CAT, obj = self.obj(), "failed to start session: {err:?}"); - return Err(gst::StateChangeError); - } - let success = self.parent_change_state(transition)?; - let result = match success { - gst::StateChangeSuccess::Async => gst::StateChangeSuccess::Async, - _ => gst::StateChangeSuccess::NoPreroll, - }; - Ok(result) - } - gst::StateChange::PausedToReady => { - self.stop_session(); - self.parent_change_state(transition) - } - _ => self.parent_change_state(transition), - } - } -} - -impl MoqSrc { - fn start_session(&self) -> Result<()> { - let settings = { - let settings = self.settings.lock().unwrap().clone(); - ResolvedSettings::try_from(settings)? - }; - - let (control_tx, control_rx) = mpsc::unbounded_channel(); - let obj = self.obj(); - let weak = obj.downgrade(); - let context = glib::MainContext::default(); - let control_task = spawn_main_context_forwarder(&context, control_rx, move |msg| { - if let Some(obj) = weak.upgrade() { - obj.imp().handle_control_message(msg); - true - } else { - false - } - }); - - *self.control_task.lock().unwrap() = Some(control_task); - *self.control_sender.lock().unwrap() = Some(control_tx.clone()); - - let session = SessionController::start(settings, control_tx); - *self.session.lock().unwrap() = Some(session); - Ok(()) - } - - fn stop_session(&self) { - if let Some(session) = self.session.lock().unwrap().take() { - session.stop(); - } - - if let Some(control_task) = self.control_task.lock().unwrap().take() { - control_task.abort(); - } - - let handles = self.pads.lock().unwrap().drain().collect::>(); - for (name, handle) in handles { - gst::debug!(CAT, "dropping pad {name}"); - let _ = handle.sender.send(PadMessage::Drop); - handle.task.abort(); - } - - *self.control_sender.lock().unwrap() = None; - } - - fn handle_control_message(&self, msg: ControlMessage) { - match msg { - ControlMessage::CreatePad { - descriptor, - caps, - reply, - } => { - if let Err(err) = self.create_pad(descriptor, caps, reply) { - 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:?}"]); - } - } - } - - fn create_pad( - &self, - descriptor: TrackDescriptor, - caps: gst::Caps, - reply: oneshot::Sender, - ) -> Result<()> { - let obj = self.obj(); - let templ = obj - .element_class() - .pad_template(descriptor.kind.template_name()) - .context("missing pad template")?; - - let pad = gst::Pad::builder_from_template(&templ) - .name(descriptor.pad_name()) - .build(); - - pad.set_active(true)?; - - let stream_start = gst::event::StreamStart::builder(&descriptor.name) - .group_id(gst::GroupId::next()) - .build(); - pad.push_event(stream_start); - pad.push_event(gst::event::Caps::new(&caps)); - pad.push_event(gst::event::Segment::new(&gst::FormattedSegment::::new())); - - obj.add_pad(&pad)?; - - let (pad_tx, pad_rx) = mpsc::unbounded_channel(); - let pad_clone = pad.clone(); - let weak = obj.downgrade(); - let context = glib::MainContext::default(); - let task = spawn_main_context_forwarder(&context, pad_rx, move |msg| { - if let Some(obj) = weak.upgrade() { - let imp = obj.imp(); - imp.dispatch_pad_message(&pad_clone, msg) - } else { - false - } - }); - - self.pads.lock().unwrap().insert( - descriptor.pad_name(), - PadHandle { - sender: pad_tx.clone(), - task, - }, - ); - - let _ = reply.send(PadEndpoint { sender: pad_tx }); - Ok(()) - } - - fn dispatch_pad_message(&self, pad: &gst::Pad, msg: PadMessage) -> bool { - match msg { - PadMessage::Buffer(buffer) => { - if let Err(err) = pad.push(buffer) { - gst::warning!(CAT, "failed to push buffer: {err:?}"); - return false; - } - true - } - PadMessage::Eos => { - pad.push_event(gst::event::Eos::builder().build()); - true - } - PadMessage::Drop => { - let _ = pad.set_active(false); - let _ = self.obj().remove_pad(pad); - false - } - } - } -} - -async fn run_session( - settings: ResolvedSettings, - 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?; - - // Wait for the broadcast to be announced. Synchronous lookup would race the gossip of - // announcements that happens after the session is established. - tracing::info!(broadcast = %settings.broadcast, "waiting for broadcast to be announced"); - let broadcast = tokio::select! { - broadcast = origin_consumer.announced_broadcast(&settings.broadcast) => broadcast - .context("broadcast not allowed or origin closed")?, - _ = shutdown.changed() => return Ok(()), - }; - - 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 consumer_latency = Duration::from_millis(settings.max_latency_ms); - - 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())); - } - - 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())); - } - - let _ = control_tx.send(ControlMessage::NoMorePads); - - for task in tasks { - let _ = task.await; - } - - Ok(()) -} - -fn subscribe_track( - broadcast: &moq_net::BroadcastConsumer, - name: &str, - ascending: bool, - latency: Duration, -) -> Result> { - #[allow(unused_mut)] - let mut track_ref = moq_net::Track::new(name); - #[cfg(feature = "group-order")] - { track_ref.ordered = ascending; } - #[cfg(not(feature = "group-order"))] - { let _ = ascending; } - let consumer = broadcast.subscribe_track(&track_ref)?; - Ok(moq_mux::container::Consumer::new(consumer, moq_mux::catalog::hang::Container::Legacy) - .with_latency(latency)) -} - -async fn request_pad( - control_tx: &mpsc::UnboundedSender, - descriptor: TrackDescriptor, - caps: gst::Caps, -) -> Result { - let (reply_tx, reply_rx) = oneshot::channel(); - control_tx - .send(ControlMessage::CreatePad { - descriptor, - caps, - reply: reply_tx, - }) - .map_err(|_| anyhow::anyhow!("control plane shut down"))?; - - let endpoint = reply_rx.await.context("pad creation cancelled")?; - Ok(endpoint) -} - -fn spawn_track_pump( - track: moq_mux::container::Consumer, - descriptor: TrackDescriptor, - pad_endpoint: PadEndpoint, - shutdown: watch::Receiver, -) -> tokio::task::JoinHandle<()> { - RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown)) -} - -async fn run_track_pump( - mut track: moq_mux::container::Consumer, - descriptor: TrackDescriptor, - pad_endpoint: PadEndpoint, - mut shutdown: watch::Receiver, -) { - let mut reference_ts = None; - loop { - tokio::select! { - _ = shutdown.changed() => { - pad_endpoint.send(PadMessage::Drop); - break; - } - frame = track.read() => { - match frame { - Ok(Some(frame)) => { - let timestamp = frame.timestamp; - let is_keyframe = frame.keyframe; - let payload = frame.payload; - let mut buffer = gst::Buffer::from_slice(payload.to_vec()); - let buffer_mut = buffer.get_mut().unwrap(); - - let pts = match reference_ts { - Some(reference) => { - let delta: Duration = (timestamp - reference).into(); - gst::ClockTime::from_nseconds(delta.as_nanos() as u64) - } - None => { - reference_ts = Some(timestamp); - gst::ClockTime::ZERO - } - }; - buffer_mut.set_pts(Some(pts)); - - let mut flags = buffer_mut.flags(); - match descriptor.kind { - TrackKind::Video => { - if is_keyframe { - flags.remove(gst::BufferFlags::DELTA_UNIT); - } else { - flags.insert(gst::BufferFlags::DELTA_UNIT); - } - } - TrackKind::Audio => { - flags.remove(gst::BufferFlags::DELTA_UNIT); - } - } - buffer_mut.set_flags(flags); - - if !pad_endpoint.send(PadMessage::Buffer(buffer)) { - break; - } - } - Ok(None) => { - pad_endpoint.send(PadMessage::Eos); - pad_endpoint.send(PadMessage::Drop); - break; - } - Err(err) => { - gst::warning!(CAT, "track {} failed: {err:?}", descriptor.name); - pad_endpoint.send(PadMessage::Drop); - break; - } - } - } - } - } -} - -fn video_caps(config: &hang::catalog::VideoConfig) -> Result { - use hang::catalog::VideoCodec; - - let caps = match &config.codec { - VideoCodec::H264(_) => { - let mut builder = gst::Caps::builder("video/x-h264").field("alignment", "au"); - if let Some(description) = &config.description { - builder = builder - .field("stream-format", "avc") - .field("codec_data", gst::Buffer::from_slice(description.clone())); - } else { - builder = builder.field("stream-format", "annexb"); - } - builder.build() - } - VideoCodec::H265(h265) => { - let mut builder = gst::Caps::builder("video/x-h265").field("alignment", "au"); - match &config.description { - Some(description) => { - let format = if h265.in_band { "hev1" } else { "hvc1" }; - builder = builder - .field("stream-format", format) - .field("codec_data", gst::Buffer::from_slice(description.clone())); - } - None => { - let format = if h265.in_band { "hev1" } else { "byte-stream" }; - builder = builder.field("stream-format", format); - } - } - builder.build() - } - VideoCodec::AV1(_) => { - let mut builder = gst::Caps::builder("video/x-av1"); - if let Some(description) = &config.description { - builder = builder.field("codec_data", gst::Buffer::from_slice(description.clone())); - } - builder.build() - } - other => bail!("unsupported video codec: {other:?}"), - }; - Ok(caps) -} - -fn audio_caps(config: &hang::catalog::AudioConfig) -> Result { - let caps = match &config.codec { - hang::catalog::AudioCodec::AAC(_) => { - let mut builder = gst::Caps::builder("audio/mpeg") - .field("mpegversion", 4) - .field("rate", config.sample_rate) - .field("channels", config.channel_count); - if let Some(description) = &config.description { - builder = builder - .field("codec_data", gst::Buffer::from_slice(description.clone())) - .field("stream-format", "aac"); - } else { - builder = builder.field("stream-format", "adts"); - } - builder.build() - } - hang::catalog::AudioCodec::Opus => { - let mut builder = gst::Caps::builder("audio/x-opus") - .field("rate", config.sample_rate) - .field("channels", config.channel_count); - if let Some(description) = &config.description { - builder = builder - .field("codec_data", gst::Buffer::from_slice(description.clone())) - .field("stream-format", "ogg"); - } - builder.build() - } - other => bail!("unsupported audio codec: {other:?}"), - }; - Ok(caps) -} - -fn spawn_main_context_forwarder( - context: &glib::MainContext, - mut rx: mpsc::UnboundedReceiver, - mut handler: F, -) -> glib::JoinHandle<()> -where - T: Send + 'static, - F: FnMut(T) -> bool + 'static, -{ - let ctx = context.clone(); - ctx.spawn_local(async move { - while let Some(msg) = rx.recv().await { - if !handler(msg) { - break; - } - } - }) -} +use std::collections::HashMap; +use std::sync::{LazyLock, Mutex}; +use std::time::Duration; + +use anyhow::{bail, Context, Result}; +use gst::glib; +use gst::prelude::*; +use gst::subclass::prelude::*; +use tokio::sync::{mpsc, oneshot, watch}; + +use hang::moq_net; + +static CAT: LazyLock = + LazyLock::new(|| gst::DebugCategory::new("moq-src", gst::DebugColorFlags::empty(), Some("MoQ Source Element"))); + +static RUNTIME: LazyLock = LazyLock::new(|| { + tokio::runtime::Builder::new_multi_thread() + .enable_all() + .build() + .expect("spawn tokio runtime") +}); + +#[derive(Debug, Clone)] +struct Settings { + url: Option, + broadcast: Option, + tls_disable_verify: bool, + max_latency_ms: u64, + ascending: bool, +} + +impl Default for Settings { + fn default() -> Self { + Self { + url: None, + broadcast: None, + tls_disable_verify: false, + max_latency_ms: 1000, + ascending: false, + } + } +} + +#[derive(Debug, Clone)] +struct ResolvedSettings { + url: url::Url, + broadcast: String, + tls_disable_verify: bool, + max_latency_ms: u64, + ascending: bool, +} + +impl TryFrom for ResolvedSettings { + type Error = anyhow::Error; + + fn try_from(value: Settings) -> Result { + 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(), + tls_disable_verify: value.tls_disable_verify, + max_latency_ms: value.max_latency_ms, + ascending: value.ascending, + }) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +enum TrackKind { + Video, + Audio, +} + +impl TrackKind { + fn template_name(&self) -> &'static str { + match self { + TrackKind::Video => "video_%u", + TrackKind::Audio => "audio_%u", + } + } +} + +#[derive(Debug, Clone)] +struct TrackDescriptor { + kind: TrackKind, + name: String, +} + +impl TrackDescriptor { + fn pad_name(&self) -> String { + match self.kind { + TrackKind::Video => format!("video_{}", self.name), + TrackKind::Audio => format!("audio_{}", self.name), + } + } +} + +#[derive(Debug)] +enum ControlMessage { + CreatePad { + descriptor: TrackDescriptor, + caps: gst::Caps, + reply: oneshot::Sender, + }, + NoMorePads, + ReportError(anyhow::Error), +} + +#[derive(Debug)] +enum PadMessage { + Buffer(gst::Buffer), + Eos, + Drop, +} + +#[derive(Debug, Clone)] +struct PadEndpoint { + sender: mpsc::UnboundedSender, +} + +impl PadEndpoint { + fn send(&self, msg: PadMessage) -> bool { + self.sender.send(msg).is_ok() + } +} + +struct PadHandle { + sender: mpsc::UnboundedSender, + task: glib::JoinHandle<()>, +} + +struct SessionController { + shutdown: 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 control_for_error = control_tx.clone(); + let join = RUNTIME.spawn(async move { + let result = run_session(settings, control_tx, &mut shutdown_rx).await; + if let Err(err) = result { + let _ = control_for_error.send(ControlMessage::ReportError(err)); + } + }); + + Self { + shutdown: shutdown_tx, + join, + } + } + + fn stop(self) { + let _ = self.shutdown.send(true); + RUNTIME.spawn(async move { + if let Err(err) = self.join.await { + gst::warning!(CAT, "session task ended with error: {err:?}"); + } + }); + } +} + +#[derive(Default)] +pub struct MoqSrc { + settings: Mutex, + pads: Mutex>, + control_task: Mutex>>, + control_sender: Mutex>>, + session: Mutex>, +} + +#[glib::object_subclass] +impl ObjectSubclass for MoqSrc { + const NAME: &'static str = "MoqSrc"; + type Type = super::MoqSrc; + type ParentType = gst::Element; + + fn new() -> Self { + Self::default() + } +} + +impl ObjectImpl for MoqSrc { + fn properties() -> &'static [glib::ParamSpec] { + static PROPS: LazyLock> = LazyLock::new(|| { + vec![ + glib::ParamSpecString::builder("url") + .nick("Source URL") + .blurb("Connect to the given URL") + .mutable_ready() + .build(), + glib::ParamSpecString::builder("broadcast") + .nick("Broadcast") + .blurb("The broadcast name to subscribe to") + .mutable_ready() + .build(), + glib::ParamSpecBoolean::builder("tls-disable-verify") + .nick("TLS Disable Verify") + .blurb("Disable TLS certificate verification") + .default_value(false) + .mutable_ready() + .build(), + glib::ParamSpecUInt64::builder("max-latency-ms") + .nick("Max latency (ms)") + .blurb("Drop groups older than this to stay at the live edge.") + .default_value(1000) + .mutable_ready() + .build(), + glib::ParamSpecBoolean::builder("ascending") + .nick("Ascending") + .blurb("Deliver groups oldest-first; default false delivers newest-first (live edge)") + .default_value(false) + .mutable_ready() + .build(), + ] + }); + PROPS.as_ref() + } + + fn set_property(&self, _id: usize, value: &glib::Value, pspec: &glib::ParamSpec) { + let mut settings = self.settings.lock().unwrap(); + match pspec.name() { + "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(), + "ascending" => { + settings.ascending = value.get().unwrap(); + #[cfg(not(feature = "group-order"))] + if settings.ascending { + gst::warning!(CAT, "ascending=true has no effect: build with --features group-order"); + } + } + _ => unreachable!(), + } + } + + fn property(&self, _id: usize, pspec: &glib::ParamSpec) -> glib::Value { + let settings = self.settings.lock().unwrap(); + match pspec.name() { + "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(), + "ascending" => settings.ascending.to_value(), + _ => unreachable!(), + } + } +} + +impl GstObjectImpl for MoqSrc {} +impl ElementImpl for MoqSrc { + fn metadata() -> Option<&'static gst::subclass::ElementMetadata> { + static META: LazyLock = LazyLock::new(|| { + gst::subclass::ElementMetadata::new( + "MoQ Src", + "Source/Network/MoQ", + "Receives media over the network via MoQ", + "Luke Curley , Steve McFarlin ", + ) + }); + Some(&*META) + } + + fn pad_templates() -> &'static [gst::PadTemplate] { + static PAD_TEMPLATES: LazyLock> = LazyLock::new(|| { + vec![ + gst::PadTemplate::new( + "video_%u", + gst::PadDirection::Src, + gst::PadPresence::Sometimes, + &gst::Caps::new_any(), + ) + .unwrap(), + gst::PadTemplate::new( + "audio_%u", + gst::PadDirection::Src, + gst::PadPresence::Sometimes, + &gst::Caps::new_any(), + ) + .unwrap(), + ] + }); + PAD_TEMPLATES.as_ref() + } + + fn change_state(&self, transition: gst::StateChange) -> Result { + match transition { + gst::StateChange::ReadyToPaused => { + if let Err(err) = self.start_session() { + gst::error!(CAT, obj = self.obj(), "failed to start session: {err:?}"); + return Err(gst::StateChangeError); + } + let success = self.parent_change_state(transition)?; + let result = match success { + gst::StateChangeSuccess::Async => gst::StateChangeSuccess::Async, + _ => gst::StateChangeSuccess::NoPreroll, + }; + Ok(result) + } + gst::StateChange::PausedToReady => { + self.stop_session(); + self.parent_change_state(transition) + } + _ => self.parent_change_state(transition), + } + } +} + +impl MoqSrc { + fn start_session(&self) -> Result<()> { + let settings = { + let settings = self.settings.lock().unwrap().clone(); + ResolvedSettings::try_from(settings)? + }; + + let (control_tx, control_rx) = mpsc::unbounded_channel(); + let obj = self.obj(); + let weak = obj.downgrade(); + let context = glib::MainContext::default(); + let control_task = spawn_main_context_forwarder(&context, control_rx, move |msg| { + if let Some(obj) = weak.upgrade() { + obj.imp().handle_control_message(msg); + true + } else { + false + } + }); + + *self.control_task.lock().unwrap() = Some(control_task); + *self.control_sender.lock().unwrap() = Some(control_tx.clone()); + + let session = SessionController::start(settings, control_tx); + *self.session.lock().unwrap() = Some(session); + Ok(()) + } + + fn stop_session(&self) { + if let Some(session) = self.session.lock().unwrap().take() { + session.stop(); + } + + if let Some(control_task) = self.control_task.lock().unwrap().take() { + control_task.abort(); + } + + let handles = self.pads.lock().unwrap().drain().collect::>(); + for (name, handle) in handles { + gst::debug!(CAT, "dropping pad {name}"); + let _ = handle.sender.send(PadMessage::Drop); + handle.task.abort(); + } + + *self.control_sender.lock().unwrap() = None; + } + + fn handle_control_message(&self, msg: ControlMessage) { + match msg { + ControlMessage::CreatePad { + descriptor, + caps, + reply, + } => { + if let Err(err) = self.create_pad(descriptor, caps, reply) { + 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:?}"]); + } + } + } + + fn create_pad( + &self, + descriptor: TrackDescriptor, + caps: gst::Caps, + reply: oneshot::Sender, + ) -> Result<()> { + let obj = self.obj(); + let templ = obj + .element_class() + .pad_template(descriptor.kind.template_name()) + .context("missing pad template")?; + + let pad = gst::Pad::builder_from_template(&templ) + .name(descriptor.pad_name()) + .build(); + + pad.set_active(true)?; + + let stream_start = gst::event::StreamStart::builder(&descriptor.name) + .group_id(gst::GroupId::next()) + .build(); + pad.push_event(stream_start); + pad.push_event(gst::event::Caps::new(&caps)); + pad.push_event(gst::event::Segment::new(&gst::FormattedSegment::::new())); + + obj.add_pad(&pad)?; + + let (pad_tx, pad_rx) = mpsc::unbounded_channel(); + let pad_clone = pad.clone(); + let weak = obj.downgrade(); + let context = glib::MainContext::default(); + let task = spawn_main_context_forwarder(&context, pad_rx, move |msg| { + if let Some(obj) = weak.upgrade() { + let imp = obj.imp(); + imp.dispatch_pad_message(&pad_clone, msg) + } else { + false + } + }); + + self.pads.lock().unwrap().insert( + descriptor.pad_name(), + PadHandle { + sender: pad_tx.clone(), + task, + }, + ); + + let _ = reply.send(PadEndpoint { sender: pad_tx }); + Ok(()) + } + + fn dispatch_pad_message(&self, pad: &gst::Pad, msg: PadMessage) -> bool { + match msg { + PadMessage::Buffer(buffer) => { + if let Err(err) = pad.push(buffer) { + gst::warning!(CAT, "failed to push buffer: {err:?}"); + return false; + } + true + } + PadMessage::Eos => { + pad.push_event(gst::event::Eos::builder().build()); + true + } + PadMessage::Drop => { + let _ = pad.set_active(false); + let _ = self.obj().remove_pad(pad); + false + } + } + } +} + +async fn run_session( + settings: ResolvedSettings, + 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?; + + // Wait for the broadcast to be announced. Synchronous lookup would race the gossip of + // announcements that happens after the session is established. + tracing::info!(broadcast = %settings.broadcast, "waiting for broadcast to be announced"); + let broadcast = tokio::select! { + broadcast = origin_consumer.announced_broadcast(&settings.broadcast) => broadcast + .context("broadcast not allowed or origin closed")?, + _ = shutdown.changed() => return Ok(()), + }; + + 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 consumer_latency = Duration::from_millis(settings.max_latency_ms); + + 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())); + } + + 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())); + } + + let _ = control_tx.send(ControlMessage::NoMorePads); + + for task in tasks { + let _ = task.await; + } + + Ok(()) +} + +fn subscribe_track( + broadcast: &moq_net::BroadcastConsumer, + name: &str, + ascending: bool, + latency: Duration, +) -> Result> { + #[allow(unused_mut)] + let mut track_ref = moq_net::Track::new(name); + #[cfg(feature = "group-order")] + { track_ref.ordered = ascending; } + #[cfg(not(feature = "group-order"))] + { let _ = ascending; } + let consumer = broadcast.subscribe_track(&track_ref)?; + Ok(moq_mux::container::Consumer::new(consumer, moq_mux::catalog::hang::Container::Legacy) + .with_latency(latency)) +} + +async fn request_pad( + control_tx: &mpsc::UnboundedSender, + descriptor: TrackDescriptor, + caps: gst::Caps, +) -> Result { + let (reply_tx, reply_rx) = oneshot::channel(); + control_tx + .send(ControlMessage::CreatePad { + descriptor, + caps, + reply: reply_tx, + }) + .map_err(|_| anyhow::anyhow!("control plane shut down"))?; + + let endpoint = reply_rx.await.context("pad creation cancelled")?; + Ok(endpoint) +} + +fn spawn_track_pump( + track: moq_mux::container::Consumer, + descriptor: TrackDescriptor, + pad_endpoint: PadEndpoint, + shutdown: watch::Receiver, +) -> tokio::task::JoinHandle<()> { + RUNTIME.spawn(run_track_pump(track, descriptor, pad_endpoint, shutdown)) +} + +async fn run_track_pump( + mut track: moq_mux::container::Consumer, + descriptor: TrackDescriptor, + pad_endpoint: PadEndpoint, + mut shutdown: watch::Receiver, +) { + let mut reference_ts = None; + loop { + tokio::select! { + _ = shutdown.changed() => { + pad_endpoint.send(PadMessage::Drop); + break; + } + frame = track.read() => { + match frame { + Ok(Some(frame)) => { + let timestamp = frame.timestamp; + let is_keyframe = frame.keyframe; + let payload = frame.payload; + let mut buffer = gst::Buffer::from_slice(payload.to_vec()); + let buffer_mut = buffer.get_mut().unwrap(); + + let pts = match reference_ts { + Some(reference) => { + let delta: Duration = (timestamp - reference).into(); + gst::ClockTime::from_nseconds(delta.as_nanos() as u64) + } + None => { + reference_ts = Some(timestamp); + gst::ClockTime::ZERO + } + }; + buffer_mut.set_pts(Some(pts)); + + let mut flags = buffer_mut.flags(); + match descriptor.kind { + TrackKind::Video => { + if is_keyframe { + flags.remove(gst::BufferFlags::DELTA_UNIT); + } else { + flags.insert(gst::BufferFlags::DELTA_UNIT); + } + } + TrackKind::Audio => { + flags.remove(gst::BufferFlags::DELTA_UNIT); + } + } + buffer_mut.set_flags(flags); + + if !pad_endpoint.send(PadMessage::Buffer(buffer)) { + break; + } + } + Ok(None) => { + pad_endpoint.send(PadMessage::Eos); + pad_endpoint.send(PadMessage::Drop); + break; + } + Err(err) => { + gst::warning!(CAT, "track {} failed: {err:?}", descriptor.name); + pad_endpoint.send(PadMessage::Drop); + break; + } + } + } + } + } +} + +fn video_caps(config: &hang::catalog::VideoConfig) -> Result { + use hang::catalog::VideoCodec; + + let caps = match &config.codec { + VideoCodec::H264(_) => { + let mut builder = gst::Caps::builder("video/x-h264").field("alignment", "au"); + if let Some(description) = &config.description { + builder = builder + .field("stream-format", "avc") + .field("codec_data", gst::Buffer::from_slice(description.clone())); + } else { + builder = builder.field("stream-format", "annexb"); + } + builder.build() + } + VideoCodec::H265(h265) => { + let mut builder = gst::Caps::builder("video/x-h265").field("alignment", "au"); + match &config.description { + Some(description) => { + let format = if h265.in_band { "hev1" } else { "hvc1" }; + builder = builder + .field("stream-format", format) + .field("codec_data", gst::Buffer::from_slice(description.clone())); + } + None => { + let format = if h265.in_band { "hev1" } else { "byte-stream" }; + builder = builder.field("stream-format", format); + } + } + builder.build() + } + VideoCodec::AV1(_) => { + let mut builder = gst::Caps::builder("video/x-av1"); + if let Some(description) = &config.description { + builder = builder.field("codec_data", gst::Buffer::from_slice(description.clone())); + } + builder.build() + } + other => bail!("unsupported video codec: {other:?}"), + }; + Ok(caps) +} + +fn audio_caps(config: &hang::catalog::AudioConfig) -> Result { + let caps = match &config.codec { + hang::catalog::AudioCodec::AAC(_) => { + let mut builder = gst::Caps::builder("audio/mpeg") + .field("mpegversion", 4) + .field("rate", config.sample_rate) + .field("channels", config.channel_count); + if let Some(description) = &config.description { + builder = builder + .field("codec_data", gst::Buffer::from_slice(description.clone())) + .field("stream-format", "aac"); + } else { + builder = builder.field("stream-format", "adts"); + } + builder.build() + } + hang::catalog::AudioCodec::Opus => { + let mut builder = gst::Caps::builder("audio/x-opus") + .field("rate", config.sample_rate) + .field("channels", config.channel_count); + if let Some(description) = &config.description { + builder = builder + .field("codec_data", gst::Buffer::from_slice(description.clone())) + .field("stream-format", "ogg"); + } + builder.build() + } + other => bail!("unsupported audio codec: {other:?}"), + }; + Ok(caps) +} + +fn spawn_main_context_forwarder( + context: &glib::MainContext, + mut rx: mpsc::UnboundedReceiver, + mut handler: F, +) -> glib::JoinHandle<()> +where + T: Send + 'static, + F: FnMut(T) -> bool + 'static, +{ + let ctx = context.clone(); + ctx.spawn_local(async move { + while let Some(msg) = rx.recv().await { + if !handler(msg) { + break; + } + } + }) +}