mod catalog;
mod config;
mod error;
mod feed;
mod rung;
pub use config::{Config, Rung};
pub use error::Error;
pub async fn run(
source: moq_net::broadcast::Consumer,
mut output: moq_net::broadcast::Producer,
config: Config,
) -> Result<(), Error> {
let mut derived = moq_mux::catalog::Producer::new(&mut output)?;
let mut dynamic = output.dynamic();
let track = source
.track(hang::Catalog::DEFAULT_NAME)?
.subscribe(hang::Catalog::default_subscription())
.await?;
let mut catalogs = moq_mux::catalog::hang::Consumer::<()>::new(track);
let (source_name, source_config, snapshot) = loop {
let Some(snapshot) = catalogs.next().await? else {
return Err(Error::NoSource);
};
match catalog::choose_source(&snapshot.video) {
Ok((name, config)) => break (name, config, snapshot),
Err(_) => tracing::debug!("no transcodable rendition yet; waiting for a catalog update"),
}
};
let rungs = catalog::resolve_rungs(&config.rungs, &source_name, &source_config)?;
tracing::info!(source = %source_name, rungs = rungs.len(), "transcoding");
let feed = feed::Feed::new(
source.track(&source_name)?,
source_config.clone(),
config.decoder.clone(),
);
let entries: Vec<_> = rungs
.iter()
.map(|rung| (rung.name.clone(), catalog::rung_entry(rung, &source_config)))
.collect();
{
let mut guard = derived.lock();
catalog::populate(&mut guard, &snapshot, &entries, config.source.as_ref())?;
}
let mut tasks = tokio::task::JoinSet::new();
loop {
tokio::select! {
request = dynamic.requested_track() => {
let Ok(request) = request else { break };
match rungs.iter().find(|rung| rung.name == request.name()) {
Some(info) => {
let rung = rung::Rung {
source: source.track(&source_name)?,
feed: feed.clone(),
broadcast: source.clone(),
config: source_config.clone(),
encoder: config.encoder.clone(),
decoder: config.decoder.clone(),
info: info.clone(),
};
tasks.spawn(rung::serve(rung, request));
}
None => request.reject(moq_net::Error::NotFound),
}
},
update = catalogs.next() => match update {
Ok(Some(snapshot)) => {
let mut guard = derived.lock();
catalog::populate(&mut guard, &snapshot, &entries, config.source.as_ref())?;
}
Ok(None) => break,
Err(err) => {
tracing::debug!(%err, "source catalog ended");
break;
}
},
Some(result) = tasks.join_next() => match result {
Ok(Ok(())) => {}
Ok(Err(err)) => tracing::warn!(%err, "rung failed"),
Err(err) => tracing::warn!(%err, "rung panicked"),
}
}
}
tasks.shutdown().await;
derived.finish()?;
output.finish();
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
struct Source {
broadcast: moq_net::broadcast::Producer,
_catalog: moq_mux::catalog::Producer,
_track: moq_net::track::Producer,
}
fn source_broadcast(groups: u64, frames: u64) -> Source {
let mut broadcast = moq_net::broadcast::Info::default().produce();
let mut catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
let mut video = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
video.coded_width = Some(320);
video.coded_height = Some(240);
video.bitrate = Some(1_000_000);
video.framerate = Some(30.0);
catalog.lock().video.insert("video", video).unwrap();
let info = hang::container::track_info();
let mut track = broadcast.create_track("video", info).unwrap();
let mut encoder = moq_video::encode::Encoder::new(&{
let mut config = moq_video::encode::Config::new(320, 240, 30);
config.kind = moq_video::encode::Kind::Software;
config
})
.unwrap();
let gray = vec![0x80u8; 320 * 240 * 4];
for sequence in 0..groups {
let mut group = track.create_group(sequence.into()).unwrap();
for index in 0..frames {
let timestamp = (sequence * frames + index) * 33_333;
for payload in encoder
.encode_rgba(&gray, moq_video::Size::new(320, 240), index == 0)
.unwrap()
{
let frame = hang::container::Frame {
timestamp: moq_net::Timestamp::from_micros(timestamp).unwrap(),
payload,
};
frame.write_to(&mut group).unwrap();
}
}
group.finish().unwrap();
}
Source {
broadcast,
_catalog: catalog,
_track: track,
}
}
fn source_broadcast_live(groups: u64, frames: u64) -> (Source, tokio::task::JoinHandle<()>) {
let mut broadcast = moq_net::broadcast::Info::default().produce();
let mut catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
let mut video = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
video.coded_width = Some(320);
video.coded_height = Some(240);
video.bitrate = Some(1_000_000);
video.framerate = Some(30.0);
catalog.lock().video.insert("video", video).unwrap();
let info = hang::container::track_info();
let mut track = broadcast.create_track("video", info).unwrap();
let source = Source {
broadcast,
_catalog: catalog,
_track: track.clone(),
};
let task = tokio::spawn(async move {
let mut encoder = moq_video::encode::Encoder::new(&{
let mut config = moq_video::encode::Config::new(320, 240, 30);
config.kind = moq_video::encode::Kind::Software;
config
})
.unwrap();
let gray = vec![0x80u8; 320 * 240 * 4];
for sequence in 0..groups {
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
let mut group = track.create_group(sequence.into()).unwrap();
for index in 0..frames {
let timestamp = (sequence * frames + index) * 33_333;
for payload in encoder
.encode_rgba(&gray, moq_video::Size::new(320, 240), index == 0)
.unwrap()
{
let frame = hang::container::Frame {
timestamp: moq_net::Timestamp::from_micros(timestamp).unwrap(),
payload,
};
frame.write_to(&mut group).unwrap();
}
}
group.finish().unwrap();
}
std::future::pending::<()>().await;
});
(source, task)
}
#[tokio::test]
async fn live_multi_rung() {
tokio::time::pause();
let (source, producer_task) = source_broadcast_live(3, 5);
let config = Config {
rungs: vec![Rung::new(120, 100_000), Rung::new(60, 50_000)],
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let mut subscribers = Vec::new();
for name in ["video/120p", "video/60p"] {
let track = loop {
match consumer.track(name) {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("rung track {name}: {err}"),
}
};
subscribers.push((name, track.subscribe(None).await.unwrap()));
}
for (name, subscriber) in &mut subscribers {
let mut group = subscriber.next_group().await.unwrap().unwrap();
let payload = group.read_frame().await.unwrap().unwrap();
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(
frame.payload.starts_with(&[0, 0, 0, 1]) || frame.payload.starts_with(&[0, 0, 1]),
"{name} output is not Annex-B"
);
let total = group.finished().await.unwrap();
assert_eq!(total, 5, "{name} dropped frames");
}
producer_task.abort();
transcoder.abort();
}
#[tokio::test]
async fn live_multi_rung_hardware() {
if !hardware_available() {
eprintln!("skipping: no hardware decoder + encoder available");
return;
}
tokio::time::pause();
let (source, producer_task) = source_broadcast_live(3, 5);
let config = Config {
rungs: vec![Rung::new(180, 200_000), Rung::new(120, 100_000)],
encoder: moq_video::encode::Kind::Hardware,
decoder: moq_video::decode::Kind::Hardware,
source: None,
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let mut subscribers = Vec::new();
for name in ["video/180p", "video/120p"] {
let track = loop {
match consumer.track(name) {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("rung track {name}: {err}"),
}
};
subscribers.push((name, track.subscribe(None).await.unwrap()));
}
for (name, subscriber) in &mut subscribers {
let mut group = subscriber.next_group().await.unwrap().unwrap();
let payload = group.read_frame().await.unwrap().unwrap();
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(
frame.payload.starts_with(&[0, 0, 0, 1]) || frame.payload.starts_with(&[0, 0, 1]),
"{name} output is not Annex-B"
);
let total = group.finished().await.unwrap();
assert_eq!(total, 5, "{name} dropped frames");
}
producer_task.abort();
transcoder.abort();
}
fn hardware_available() -> bool {
let mut encode = moq_video::encode::Config::new(160, 120, 30);
encode.kind = moq_video::encode::Kind::Hardware;
if moq_video::encode::Encoder::new(&encode).is_err() {
return false;
}
let video = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
let mut decode = moq_video::decode::Config::new();
decode.kind = moq_video::decode::Kind::Hardware;
moq_video::decode::Decoder::new(&video, &decode).is_ok()
}
#[tokio::test]
async fn end_to_end_hardware() {
if !hardware_available() {
eprintln!("skipping: no hardware decoder + encoder available");
return;
}
let source = source_broadcast(2, 5);
let config = Config {
rungs: vec![Rung::new(120, 100_000)],
encoder: moq_video::encode::Kind::Hardware,
decoder: moq_video::decode::Kind::Hardware,
source: None,
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let track = loop {
match consumer.track("video/120p") {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("rung track: {err}"),
}
};
let mut fetched = track.fetch_group(0, None).await.unwrap();
let payload = fetched.read_frame().await.unwrap().unwrap();
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(
frame.payload.starts_with(&[0, 0, 0, 1]) || frame.payload.starts_with(&[0, 0, 1]),
"hardware rung output is not Annex-B"
);
let total = fetched.finished().await.unwrap();
assert_eq!(total, 5, "hardware transcode dropped frames");
transcoder.abort();
}
#[tokio::test]
async fn end_to_end() {
let source = source_broadcast(2, 5);
let config = Config {
rungs: vec![Rung::new(120, 100_000)],
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: Some(moq_net::PathRelativeOwned::from("..".to_string())),
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let track = loop {
match consumer.track(hang::Catalog::DEFAULT_NAME) {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("catalog track: {err}"),
}
};
let track = track.subscribe(None).await.unwrap();
let mut catalogs = moq_mux::catalog::hang::Consumer::<()>::new(track);
let derived = loop {
let snapshot = catalogs.next().await.unwrap().unwrap();
if snapshot.video.renditions.contains_key("video/120p") {
break snapshot;
}
};
let rung = derived.video.renditions.get("video/120p").expect("rung missing");
assert_eq!(rung.coded_width, Some(160));
assert_eq!(rung.coded_height, Some(120));
assert_eq!(rung.bitrate, Some(100_000));
assert!(rung.codec.to_string().starts_with("avc3."));
let passthrough = derived.video.renditions.get("video").expect("passthrough missing");
assert_eq!(passthrough.broadcast.as_ref().map(|b| b.as_ref()), Some(".."));
let mut subscriber = consumer.track("video/120p").unwrap().subscribe(None).await.unwrap();
let mut group = subscriber.next_group().await.unwrap().unwrap();
assert!(group.sequence <= 1, "unexpected sequence {}", group.sequence);
let payload = group.read_frame().await.unwrap().unwrap();
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(
frame.payload.starts_with(&[0, 0, 0, 1]) || frame.payload.starts_with(&[0, 0, 1]),
"rung output is not Annex-B"
);
let mut fetched = consumer
.track("video/120p")
.unwrap()
.fetch_group(0, None)
.await
.unwrap();
let payload = fetched.read_frame().await.unwrap().unwrap();
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(!frame.payload.is_empty());
let total = fetched.finished().await.unwrap();
assert_eq!(total, 5);
transcoder.abort();
}
#[tokio::test]
async fn shuts_down_on_source_end() {
let source = source_broadcast(1, 3);
let config = Config {
rungs: vec![Rung::new(120, 100_000)],
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let track = loop {
match consumer.track(hang::Catalog::DEFAULT_NAME) {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("catalog track: {err}"),
}
};
let mut catalogs = moq_mux::catalog::hang::Consumer::<()>::new(track.subscribe(None).await.unwrap());
catalogs.next().await.unwrap().unwrap();
drop(source);
let result = tokio::time::timeout(std::time::Duration::from_secs(5), transcoder).await;
result.expect("run did not shut down within 5s").unwrap().unwrap();
}
}