Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 57 additions & 7 deletions rs/moq-gst/src/source/imp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,18 +20,31 @@ static RUNTIME: LazyLock<tokio::runtime::Runtime> = LazyLock::new(|| {
.expect("spawn tokio runtime")
});

#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone)]
struct Settings {
url: Option<String>,
broadcast: Option<String>,
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)]
struct ResolvedSettings {
url: url::Url,
broadcast: String,
tls_disable_verify: bool,
max_latency: Duration,
}

impl TryFrom<Settings> for ResolvedSettings {
Expand All @@ -46,6 +59,7 @@ impl TryFrom<Settings> for ResolvedSettings {
.context("broadcast property is required")?
.clone(),
tls_disable_verify: value.tls_disable_verify,
max_latency: Duration::from_millis(value.max_latency_ms),
})
}
}
Expand Down Expand Up @@ -116,22 +130,25 @@ struct PadHandle {

struct SessionController {
shutdown: watch::Sender<bool>,
latency: watch::Sender<Duration>,
join: tokio::task::JoinHandle<()>,
}

impl SessionController {
fn start(settings: ResolvedSettings, control_tx: mpsc::UnboundedSender<ControlMessage>) -> 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));
}
});

Self {
shutdown: shutdown_tx,
latency: latency_tx,
join,
}
}
Expand Down Expand Up @@ -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()
Expand All @@ -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!(),
}
}
Expand All @@ -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!(),
}
}
Expand Down Expand Up @@ -413,6 +445,7 @@ async fn run_session(
settings: ResolvedSettings,
control_tx: mpsc::UnboundedSender<ControlMessage>,
shutdown: &mut watch::Receiver<bool>,
latency: watch::Receiver<Duration>,
) -> Result<()> {
let mut config = moq_native::ClientConfig::default();
config.tls.disable_verify = Some(settings.tls_disable_verify);
Expand Down Expand Up @@ -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 {
Expand All @@ -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);
Expand Down Expand Up @@ -498,15 +543,17 @@ fn spawn_track_pump(
descriptor: TrackDescriptor,
pad_endpoint: PadEndpoint,
shutdown: watch::Receiver<bool>,
latency: watch::Receiver<Duration>,
) -> 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(
mut track: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
descriptor: TrackDescriptor,
pad_endpoint: PadEndpoint,
mut shutdown: watch::Receiver<bool>,
mut latency: watch::Receiver<Duration>,
) {
let mut reference_ts = None;
loop {
Expand All @@ -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)) => {
Expand Down
Loading