use std::time::Duration;
use anyhow::Context;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
moq_native::Log::new(tracing::Level::DEBUG).init()?;
let origin = moq_net::Origin::random().produce();
let consumer = origin.consume();
tokio::select! {
res = run_session(origin) => res,
res = run_subscribe(consumer) => res,
}
}
async fn run_session(origin: moq_net::origin::Producer) -> anyhow::Result<()> {
let client = moq_native::ClientConfig::default().init()?;
let url = url::Url::parse("https://cdn.moq.dev/anon/video-example").unwrap();
let reconnect = client.with_subscriber(origin).reconnect(url);
Ok(reconnect.closed().await?)
}
async fn run_subscribe(consumer: moq_net::origin::Consumer) -> anyhow::Result<()> {
let moq_net::announce::Update { path, broadcast } = consumer.announced().next().await.context("origin closed")?;
let broadcast = broadcast.with_context(|| format!("broadcast unannounced: {path}"))?;
tracing::info!(%path, "broadcast announced");
let catalog_track = broadcast
.track(hang::Catalog::DEFAULT_NAME)?
.subscribe(hang::Catalog::default_subscription())
.await?;
let mut catalog = moq_mux::catalog::hang::Consumer::<()>::new(catalog_track);
let info = catalog.next().await?.ok_or_else(|| anyhow::anyhow!("no catalog"))?;
let (name, config) = info
.video
.renditions
.iter()
.next()
.ok_or_else(|| anyhow::anyhow!("no video renditions"))?;
tracing::info!(
%name,
codec = %config.codec,
width = ?config.coded_width,
height = ?config.coded_height,
"subscribing to video track"
);
let track_consumer = broadcast
.track(name)?
.subscribe(moq_net::track::Subscription::default().with_priority(1))
.await?;
let mut ordered = moq_mux::container::Consumer::new(track_consumer, moq_mux::catalog::hang::Container::Legacy)
.with_latency(Duration::from_millis(500));
while let Some(frame) = ordered.read().await? {
tracing::info!(
timestamp = ?frame.timestamp,
keyframe = frame.keyframe,
bytes = frame.payload.len(),
"received frame"
);
}
Ok(())
}