pub mod active;
pub mod ladder;
mod catalog;
mod config;
mod error;
mod feed;
mod pipeline;
mod rung;
pub use config::Config;
pub use ladder::{Ladder, Rung};
pub use error::Error;
pub async fn run(
source: moq_net::broadcast::Consumer,
output: moq_net::broadcast::Producer,
config: Config,
) -> Result<(), Error> {
Transcoder::new(source, output, config)?.run().await
}
pub struct Transcoder {
source: moq_net::broadcast::Consumer,
output: moq_net::broadcast::Producer,
config: Config,
derived: moq_mux::catalog::Producer,
dynamic: moq_net::broadcast::Dynamic,
active: active::Producer,
}
impl Transcoder {
pub fn new(
source: moq_net::broadcast::Consumer,
mut output: moq_net::broadcast::Producer,
config: Config,
) -> Result<Self, Error> {
let derived = moq_mux::catalog::Producer::new(&mut output, moq_mux::catalog::Config::default())?;
let dynamic = output.dynamic();
Ok(Self {
source,
output,
config,
derived,
dynamic,
active: active::Producer::default(),
})
}
pub fn active(&self) -> active::Consumer {
self.active.consume()
}
pub async fn run(self) -> Result<(), Error> {
let Self {
source,
output,
config,
mut derived,
mut dynamic,
active,
} = self;
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 mut decoders = catalog::Decoders::new(config.feed_decoder());
let (source_name, source_config, snapshot) = loop {
let Some(snapshot) = catalogs.next().await? else {
return Err(Error::NoSource);
};
match catalog::choose_source(&snapshot.video, &mut decoders).await {
Ok((name, config)) => break (name, config, snapshot),
Err(Error::NoSource) => tracing::debug!("no transcodable rendition yet; waiting for a catalog update"),
Err(err) => return Err(err),
}
};
let mut ladder = pipeline::Pipeline::new(
source.clone(),
config.clone(),
active,
decoders,
source_name,
source_config,
)
.await?;
{
let mut guard = derived.modify()?;
catalog::populate(&mut guard, &snapshot, ladder.rungs(), config.source.as_ref())?;
guard.commit()?;
}
let mut tasks = tokio::task::JoinSet::new();
loop {
tokio::select! {
request = dynamic.requested_track() => {
let Ok(request) = request else { break };
match ladder.rung(request.name())? {
Some(rung) => { tasks.spawn(rung::serve(rung, request)); }
None => request.reject(moq_net::Error::NotFound),
}
},
update = catalogs.next() => match update {
Ok(Some(snapshot)) => {
ladder.follow(&snapshot.video).await?;
let mut guard = derived.modify()?;
catalog::populate(&mut guard, &snapshot, ladder.rungs(), config.source.as_ref())?;
guard.commit()?;
}
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,
size: (u32, u32),
}
impl Source {
fn resize(&mut self, width: u32, height: u32) {
self.publish(width, height, None);
}
fn describe(&mut self, description: Option<bytes::Bytes>) {
let (width, height) = self.size;
self.publish(width, height, description);
}
fn publish(&mut self, width: u32, height: u32, description: Option<bytes::Bytes>) {
let mut video = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
video.coded_width = Some(width);
video.coded_height = Some(height);
video.bitrate = Some(1_000_000);
video.framerate = Some(30.0);
video.description = description;
self.size = (width, height);
let mut guard = self.catalog.modify().unwrap();
guard.video = hang::catalog::Video::default();
guard.video.insert("video", video).unwrap();
}
}
fn source_catalog(width: u32, height: u32) -> Source {
let mut broadcast = moq_net::broadcast::Info::default().produce();
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
let track = broadcast
.create_track("video", hang::container::track_info(hang::catalog::PRIORITY.video))
.unwrap();
let mut source = Source {
broadcast,
catalog,
_track: track,
size: (width, height),
};
source.resize(width, height);
source
}
async fn await_catalog(
catalogs: &mut moq_mux::catalog::hang::Consumer<()>,
ready: impl Fn(&moq_mux::catalog::hang::Catalog) -> bool,
) -> moq_mux::catalog::hang::Catalog {
loop {
let snapshot = catalogs.next().await.unwrap().unwrap();
if ready(&snapshot) {
return snapshot;
}
}
}
async fn subscribe(consumer: &moq_net::broadcast::Consumer, name: &str) -> moq_net::track::Subscriber {
let track = loop {
match consumer.track(name) {
Ok(track) => break track,
Err(moq_net::Error::NotFound) => tokio::task::yield_now().await,
Err(err) => panic!("track {name}: {err}"),
}
};
track.subscribe(None).await.unwrap()
}
fn nal_types(annexb: &[u8]) -> Vec<u8> {
let mut types = Vec::new();
let mut i = 0;
while i + 3 < annexb.len() {
if annexb[i..i + 3] == [0, 0, 1] {
types.push(annexb[i + 3] & 0x1f);
i += 3;
} else {
i += 1;
}
}
types
}
fn write_keyframe(group: &mut moq_net::group::Producer) {
let mut encoder = moq_video::encode::Encoder::new(&{
let mut config = moq_video::encode::Config::new(320, 240, moq_video::Rate::new(30, 1).unwrap());
config.kind = moq_video::encode::Kind::Software;
config
})
.unwrap();
encoder.cut().unwrap();
let gray = vec![0x80u8; 320 * 240 * 4];
for encoded in encoder.encode(&gray_frame(&gray, 0)).unwrap() {
hang::container::Frame {
timestamp: encoded.timestamp,
payload: encoded.payload,
}
.write_to(group)
.unwrap();
}
}
fn gray_frame(rgba: &[u8], timestamp: u64) -> moq_video::Frame {
let surface = moq_video::Surface::rgba(rgba, moq_video::Size::new(320, 240)).unwrap();
moq_video::Frame::new(surface, moq_net::Timestamp::from_micros(timestamp).unwrap())
}
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, moq_mux::catalog::Config::default()).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.modify().unwrap().video.insert("video", video).unwrap();
let info = hang::container::track_info(hang::catalog::PRIORITY.video);
let 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, moq_video::Rate::new(30, 1).unwrap());
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;
if index == 0 {
encoder.cut().unwrap();
}
for encoded in encoder.encode(&gray_frame(&gray, timestamp)).unwrap() {
let frame = hang::container::Frame {
timestamp: encoded.timestamp,
payload: encoded.payload,
};
frame.write_to(&mut group).unwrap();
}
}
group.finish().unwrap();
}
Source {
broadcast,
catalog,
_track: track,
size: (320, 240),
}
}
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, moq_mux::catalog::Config::default()).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.modify().unwrap().video.insert("video", video).unwrap();
let info = hang::container::track_info(hang::catalog::PRIORITY.video);
let track = broadcast.create_track("video", info).unwrap();
let source = Source {
broadcast,
catalog,
_track: track.clone(),
size: (320, 240),
};
let task = tokio::spawn(async move {
let mut encoder = moq_video::encode::Sink::open(&{
let mut config = moq_video::encode::Config::new(320, 240, moq_video::Rate::new(30, 1).unwrap());
config.kind = moq_video::encode::Kind::Software;
config
})
.await
.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;
if index == 0 {
encoder.cut().await.unwrap();
}
for encoded in encoder.encode(gray_frame(&gray, timestamp)).await.unwrap() {
let frame = hang::container::Frame {
timestamp: encoded.timestamp,
payload: encoded.payload,
};
frame.write_to(&mut group).unwrap();
}
}
group.finish().unwrap();
}
std::future::pending::<()>().await;
});
(source, task)
}
#[tokio::test]
async fn live_multi_rung() {
let (source, producer_task) = source_broadcast_live(3, 5);
let config = Config {
ladder: Ladder::new([
Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000)),
Rung::new(60, moq_net::bandwidth::Rate::from_bps(50_000)),
])
.unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
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().ordered()));
}
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"
);
while group.read_frame().await.unwrap().is_some() {}
assert_eq!(group.frame_count(), 5, "{name} dropped frames");
}
producer_task.abort();
transcoder.abort();
}
#[cfg_attr(
target_os = "windows",
ignore = "explicit live-DXVA GPU probe; VideoProcessorBlt can hang on affected drivers"
)]
#[tokio::test]
async fn live_multi_rung_hardware() {
if !hardware_available() {
eprintln!("skipping: no hardware decoder + encoder available");
return;
}
let (source, producer_task) = source_broadcast_live(3, 5);
let config = Config {
ladder: Ladder::new([
Rung::new(180, moq_net::bandwidth::Rate::from_bps(200_000)),
Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000)),
])
.unwrap(),
encoder: moq_video::encode::Kind::Hardware,
decoder: moq_video::decode::Kind::Hardware,
source: None,
..Default::default()
};
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().ordered()));
}
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"
);
while group.read_frame().await.unwrap().is_some() {}
assert_eq!(group.frame_count(), 5, "{name} dropped frames");
}
producer_task.abort();
transcoder.abort();
}
fn hardware_available() -> bool {
let mut encode = moq_video::encode::Config::new(160, 120, moq_video::Rate::new(30, 1).unwrap());
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()
}
#[cfg(feature = "vaapi")]
fn vaapi_decoder_available() -> bool {
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::Named("vaapi".to_string());
moq_video::decode::Decoder::new(&video, &decode).is_ok()
}
#[cfg(feature = "vaapi")]
#[tokio::test]
async fn vaapi_fetch_keeps_the_buffered_tail() {
let source = source_broadcast(1, 5);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Named("vaapi".to_string()),
source: None,
..Default::default()
};
if !vaapi_decoder_available() {
eprintln!("skipping: no VA-API H.264 decoder");
return;
}
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();
while fetched.read_frame().await.unwrap().is_some() {}
assert_eq!(
fetched.frame_count(),
5,
"VAAPI dropped the fetched group's buffered tail"
);
transcoder.abort();
}
#[cfg(feature = "vaapi")]
#[tokio::test]
async fn vaapi_live_keeps_the_buffered_tail_in_its_group() {
if !vaapi_decoder_available() {
eprintln!("skipping: no VA-API H.264 decoder");
return;
}
let (source, producer_task) = source_broadcast_live(1, 5);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Named("vaapi".to_string()),
source: None,
..Default::default()
};
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 subscriber = track.subscribe(None).await.unwrap();
let mut group = subscriber.recv_group().await.unwrap().unwrap();
while group.read_frame().await.unwrap().is_some() {}
assert_eq!(
group.frame_count(),
5,
"VAAPI moved the live group's buffered tail past its end"
);
producer_task.abort();
transcoder.abort();
}
#[cfg_attr(
target_os = "windows",
ignore = "explicit live-DXVA GPU probe; VideoProcessorBlt can hang on affected drivers"
)]
#[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 {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Hardware,
decoder: moq_video::decode::Kind::Hardware,
source: None,
..Default::default()
};
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"
);
while fetched.read_frame().await.unwrap().is_some() {}
assert_eq!(fetched.frame_count(), 5, "hardware transcode dropped frames");
transcoder.abort();
}
#[tokio::test]
async fn end_to_end() {
let source = source_broadcast(2, 5);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: Some(moq_net::path::RelativeOwned::from(".".to_string())),
..Default::default()
};
let origin = moq_tokio::origin::spawn();
let output = origin.create_broadcast("room/transcode").unwrap();
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()
.ordered();
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 mut timestamps = Vec::new();
let mut first_payload = None;
while let Some(payload) = fetched.read_frame().await.unwrap() {
let frame = hang::container::Frame::decode(payload.payload).unwrap();
assert!(!frame.payload.is_empty());
timestamps.push(frame.timestamp.as_micros());
first_payload = first_payload.or(Some(frame.payload));
}
let types = nal_types(&first_payload.expect("the group had no frames"));
assert!(types.contains(&7), "group does not open with an SPS: {types:?}");
assert!(types.contains(&8), "group does not open with a PPS: {types:?}");
assert!(types.contains(&5), "group does not open with an IDR: {types:?}");
assert_eq!(timestamps, (0..5).map(|i| i * 33_333).collect::<Vec<u128>>());
let total = fetched.finished().await.unwrap();
assert_eq!(total, 5);
transcoder.abort();
}
#[tokio::test]
async fn reports_active_rungs() {
let source = source_broadcast(2, 5);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = Transcoder::new(source.broadcast.consume(), output, config).unwrap();
let mut active = transcoder.active();
let driver = tokio::spawn(transcoder.run());
let update = active.next().await.unwrap();
let rendition = update.rendition;
assert_eq!(rendition.name(), "video/120p");
assert_eq!(rendition.size().height, 120);
assert_eq!(rendition.bitrate(), moq_net::bandwidth::Rate::from_bps(100_000));
assert!(!update.encoding, "encoding before anyone asked");
assert_eq!(rendition.frames(), 0);
let mut subscriber = consumer
.track("video/120p")
.unwrap()
.subscribe(None)
.await
.unwrap()
.ordered();
let update = active.next().await.unwrap();
assert_eq!(update.rendition.name(), "video/120p");
assert!(update.encoding);
let mut group = subscriber.next_group().await.unwrap().unwrap();
group.read_frame().await.unwrap().unwrap();
assert!(rendition.frames() > 0);
assert!(rendition.bytes() > 0);
drop(group);
drop(subscriber);
let update = active.next().await.unwrap();
assert_eq!(update.rendition.name(), "video/120p");
assert!(!update.encoding);
assert!(rendition.frames() > 0);
assert!(rendition.bytes() > 0);
driver.abort();
}
#[tokio::test]
async fn a_resized_rung_takes_a_fresh_name() {
let mut source = source_catalog(640, 360);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
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());
let derived = await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/120p")
})
.await;
assert_eq!(
derived.video.renditions.get("video/120p").and_then(|v| v.coded_width),
Some(212),
"640x360 should give 120p a 212 wide picture"
);
let mut retired = subscribe(&consumer, "video/120p").await;
source.resize(480, 360);
let derived = await_catalog(&mut catalogs, |snapshot| {
!snapshot.video.renditions.contains_key("video/120p")
})
.await;
let replacement = derived
.video
.renditions
.get("video/120p.2")
.expect("the resized rung was not republished under a fresh name");
assert_eq!(replacement.coded_width, Some(160));
assert_eq!(replacement.coded_height, Some(120));
let ended = tokio::time::timeout(std::time::Duration::from_secs(5), retired.recv_group())
.await
.expect("the retired rung never ended its track")
.expect("the retired rung aborted instead of finishing");
assert!(ended.is_none(), "expected a clean end, got a group");
subscribe(&consumer, "video/120p.2").await;
transcoder.abort();
}
#[tokio::test]
async fn ladder_follows_a_source_resize() {
let mut source = source_catalog(640, 360);
let config = Config {
ladder: Ladder::new([
Rung::new(360, moq_net::bandwidth::Rate::from_bps(900_000)),
Rung::new(240, moq_net::bandwidth::Rate::from_bps(300_000)),
Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000)),
])
.unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: Some(moq_net::path::RelativeOwned::from(".".to_string())),
..Default::default()
};
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());
let derived = await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/360p")
})
.await;
assert!(derived.video.renditions.contains_key("video/240p"));
assert!(derived.video.renditions.contains_key("video/120p"));
assert_eq!(
derived.video.renditions.get("video").and_then(|v| v.coded_width),
Some(640),
"the passthrough entry should describe the source"
);
let mut retired = subscribe(&consumer, "video/240p").await;
let mut kept = subscribe(&consumer, "video/120p").await;
source.resize(320, 180);
let derived = await_catalog(&mut catalogs, |snapshot| {
!snapshot.video.renditions.contains_key("video/360p")
})
.await;
assert!(
!derived.video.renditions.contains_key("video/240p"),
"240p outlived the resize"
);
let rung = derived
.video
.renditions
.get("video/120p")
.expect("120p was retired too");
assert_eq!(rung.coded_width, Some(212));
assert_eq!(rung.coded_height, Some(120));
assert_eq!(
derived.video.renditions.get("video").and_then(|v| v.coded_width),
Some(320),
"the passthrough entry should follow the source"
);
let ended = tokio::time::timeout(std::time::Duration::from_secs(5), retired.recv_group())
.await
.expect("the retired rung never ended its track")
.expect("the retired rung aborted instead of finishing");
assert!(ended.is_none(), "expected a clean end, got a group");
assert!(
tokio::time::timeout(std::time::Duration::from_millis(100), kept.recv_group())
.await
.is_err(),
"a rung that still fits was retired anyway"
);
transcoder.abort();
}
#[tokio::test]
async fn retirement_rides_out_an_open_live_group() {
let mut source = source_catalog(320, 240);
let mut group = source._track.create_group(0u64.into()).unwrap();
write_keyframe(&mut group);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let catalog = 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(catalog.subscribe(None).await.unwrap());
await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/120p")
})
.await;
let rung = consumer.track("video/120p").unwrap();
rung.query().await.unwrap();
while rung.latest() != Some(0) {
tokio::task::yield_now().await;
}
let mut fetched = rung.fetch_group(0, None).await.unwrap();
source.resize(160, 90);
tokio::time::timeout(
std::time::Duration::from_secs(5),
await_catalog(&mut catalogs, |snapshot| {
!snapshot.video.renditions.contains_key("video/120p")
}),
)
.await
.expect("the ladder never retired the rung");
group.finish().unwrap();
let finished = tokio::time::timeout(std::time::Duration::from_secs(5), async {
while fetched.read_frame().await?.is_some() {}
fetched.finished().await
})
.await
.expect("the accepted fetch never finished");
assert!(finished.is_ok(), "retirement aborted the accepted group: {finished:?}");
transcoder.abort();
}
#[tokio::test]
async fn retirement_waits_for_an_unclaimed_fetch() {
let mut source = source_catalog(320, 240);
let source_fetches = source._track.dynamic();
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(source.broadcast.consume(), output, config));
let catalog = 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(catalog.subscribe(None).await.unwrap());
await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/120p")
})
.await;
let rung = consumer.track("video/120p").unwrap();
rung.query().await.unwrap();
assert!(
source._track.subscription_changed().await.unwrap().is_some(),
"the rung never subscribed to the live source",
);
let fetching = rung.fetch_group(7, None);
let request = source_fetches.requested_group().await.expect("the source track closed");
assert_eq!(request.sequence(), 7);
source.resize(160, 90);
await_catalog(&mut catalogs, |snapshot| {
!snapshot.video.renditions.contains_key("video/120p")
})
.await;
assert!(
source._track.subscription_changed().await.unwrap().is_none(),
"the rung kept its live source subscription after retirement",
);
let mut group = request.accept(None).unwrap();
write_keyframe(&mut group);
group.finish().unwrap();
let mut fetched = fetching
.await
.expect("retirement finished the track before the fetch claimed its group");
let frames = async {
while fetched.read_frame().await?.is_some() {}
fetched.finished().await
}
.await
.expect("retirement aborted the accepted group");
assert!(frames > 0, "the fetch claimed its group but produced no frames");
transcoder.abort();
}
#[tokio::test]
async fn a_rebuilt_decode_renames_every_rung() {
let mut source = source_catalog(320, 240);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
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());
await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/120p")
})
.await;
let mut retired = subscribe(&consumer, "video/120p").await;
source.describe(Some(bytes::Bytes::from_static(&[0x01, 0x42, 0x00, 0x1e])));
let derived = tokio::time::timeout(
std::time::Duration::from_secs(5),
await_catalog(&mut catalogs, |snapshot| {
snapshot.video.renditions.contains_key("video/120p.2")
}),
)
.await
.expect("the rebuilt decode kept the retired rung name");
assert!(
!derived.video.renditions.contains_key("video/120p"),
"the retired name is still advertised"
);
assert_eq!(
derived.video.renditions.get("video/120p.2").and_then(|v| v.coded_width),
Some(160),
"the replacement should serve the same picture under a new name"
);
let ended = tokio::time::timeout(std::time::Duration::from_secs(5), retired.recv_group())
.await
.expect("the retired rung never ended its track")
.expect("the retired rung aborted instead of finishing");
assert!(ended.is_none(), "expected a clean end, got a group");
subscribe(&consumer, "video/120p.2").await;
transcoder.abort();
}
fn mixed_codec_source() -> (
moq_net::broadcast::Producer,
moq_mux::catalog::Producer,
[moq_net::track::Producer; 2],
) {
let mut broadcast = moq_net::broadcast::Info::default().produce();
let mut catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
let mut hevc = hang::catalog::VideoConfig::new(hang::catalog::H265 {
in_band: true,
profile_space: 0,
profile_idc: 1,
profile_compatibility_flags: [0x60, 0, 0, 0],
tier_flag: false,
level_idc: 120,
constraint_flags: [0x90, 0, 0, 0, 0, 0],
});
hevc.coded_width = Some(1920);
hevc.coded_height = Some(1080);
hevc.bitrate = Some(6_000_000);
let mut avc = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
avc.coded_width = Some(640);
avc.coded_height = Some(360);
avc.bitrate = Some(1_000_000);
let mut guard = catalog.modify().unwrap();
guard.video.insert("hevc", hevc).unwrap();
guard.video.insert("avc", avc).unwrap();
guard.commit().unwrap();
let info = hang::container::track_info(hang::catalog::PRIORITY.video);
let tracks = [
broadcast.create_track("hevc", info.clone()).unwrap(),
broadcast.create_track("avc", info).unwrap(),
];
(broadcast, catalog, tracks)
}
#[tokio::test]
async fn an_undecodable_larger_rendition_is_not_the_source() {
let (broadcast, _catalog, _tracks) = mixed_codec_source();
let config = Config {
ladder: Ladder::new([
Rung::new(720, moq_net::bandwidth::Rate::from_bps(2_500_000)),
Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000)),
])
.unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
let output = moq_net::broadcast::Info::default().produce();
let consumer = output.consume();
let transcoder = tokio::spawn(run(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());
let derived = await_catalog(&mut catalogs, |snapshot| !snapshot.video.renditions.is_empty()).await;
let names: Vec<_> = derived.video.renditions.keys().map(String::as_str).collect();
assert_eq!(
names,
["video/120p"],
"the ladder was sized against the H.265 rendition"
);
transcoder.abort();
}
#[tokio::test]
async fn a_missing_decoder_is_refused() {
let (broadcast, _catalog, _tracks) = mixed_codec_source();
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Named("missing".to_string()),
source: None,
..Default::default()
};
let output = moq_net::broadcast::Info::default().produce();
let result = tokio::time::timeout(
std::time::Duration::from_secs(5),
run(broadcast.consume(), output, config),
)
.await
.expect("run kept waiting for a source it can never decode");
match result {
Err(Error::Video(moq_video::Error::UnknownDecoder { name, codec, .. })) => {
assert_eq!(name, "missing");
assert_eq!(codec, moq_video::decode::Codec::H265);
}
other => panic!("expected the decoder's refusal, got {other:?}"),
}
}
#[tokio::test]
async fn shuts_down_on_source_end() {
let source = source_broadcast(1, 3);
let config = Config {
ladder: Ladder::new([Rung::new(120, moq_net::bandwidth::Rate::from_bps(100_000))]).unwrap(),
encoder: moq_video::encode::Kind::Software,
decoder: moq_video::decode::Kind::Software,
source: None,
..Default::default()
};
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();
}
}