use hang::catalog::{AV1, Video, VideoCodec, VideoConfig};
use moq_net::path::RelativeOwned;
use crate::{Error, Ladder};
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct Resolved {
pub name: String,
pub height: u32,
pub size: moq_video::Size,
pub bitrate: moq_net::bandwidth::Rate,
pub framerate: Option<moq_video::Rate>,
}
impl Resolved {
pub fn same_shape(&self, other: &Self) -> bool {
self.height == other.height
&& self.size == other.size
&& self.bitrate == other.bitrate
&& self.framerate == other.framerate
}
}
#[derive(Default)]
pub(crate) struct Names {
revisions: std::collections::HashMap<u32, u32>,
}
impl Names {
pub fn mint(&mut self, height: u32) -> String {
let revision = self.revisions.entry(height).and_modify(|r| *r += 1).or_insert(1);
match *revision {
1 => format!("video/{height}p"),
revision => format!("video/{height}p.{revision}"),
}
}
}
#[derive(Clone, Debug)]
pub(crate) struct Published {
pub rung: Resolved,
pub entry: VideoConfig,
}
pub(crate) fn choose_source(video: &Video) -> Result<(String, VideoConfig), Error> {
video
.renditions
.iter()
.filter(|(_, config)| config.broadcast.is_none())
.filter(|(_, config)| can_decode(config))
.filter(|(_, config)| dimensions(config).is_some())
.max_by_key(|(_, config)| (config.coded_height, config.coded_width, config.bitrate))
.map(|(name, config)| (name.clone(), config.clone()))
.ok_or(Error::NoSource)
}
pub(crate) fn follow_source(video: &Video, current: &str) -> Result<(String, VideoConfig), Error> {
match video.renditions.get(current) {
Some(config) if config.broadcast.is_none() && can_decode(config) && dimensions(config).is_some() => {
Ok((current.to_string(), config.clone()))
}
_ => choose_source(video),
}
}
pub(crate) fn same_stream(old: &VideoConfig, new: &VideoConfig) -> bool {
old.codec == new.codec && old.container == new.container && old.description == new.description
}
fn dimensions(config: &VideoConfig) -> Option<(u64, u64)> {
match (config.coded_width, config.coded_height) {
(Some(w), Some(h)) if w > 0 && h > 0 => Some((w as u64, h as u64)),
_ => None,
}
}
fn can_decode(config: &VideoConfig) -> bool {
match &config.codec {
VideoCodec::H264(_) | VideoCodec::H265(_) => true,
VideoCodec::AV1(av1) => is_supported_av1(av1),
_ => false,
}
}
fn is_supported_av1(av1: &AV1) -> bool {
av1.bitdepth == 8 && !av1.mono_chrome && av1.chroma_subsampling_x && av1.chroma_subsampling_y
}
pub(crate) fn resolve_rungs(ladder: &Ladder, source_name: &str, source: &VideoConfig) -> Result<Vec<Resolved>, Error> {
let Some((source_width, source_height)) = dimensions(source) else {
return Err(Error::SourceDimensions(source_name.to_string()));
};
let framerate = source.framerate.map(moq_video::Rate::from_f64).transpose()?;
let mut resolved: Vec<Resolved> = Vec::new();
for rung in ladder.rungs() {
let height = rung.height as u64;
if height > source_height {
continue;
}
if height == source_height && source.bitrate.is_none() {
continue;
}
if source.bitrate.is_some_and(|bitrate| rung.bitrate.as_bps() >= bitrate) {
continue;
}
let width = ((source_width * height + source_height / 2) / source_height) & !1;
if width == 0 {
continue;
}
let height = height as u32;
resolved.push(Resolved {
name: format!("video/{height}p"),
height,
size: moq_video::Size::new(width as u32, height),
bitrate: rung.bitrate,
framerate,
});
}
Ok(resolved)
}
pub(crate) async fn rung_entry(
rung: &Resolved,
source: &VideoConfig,
encoder: &moq_video::encode::Kind,
) -> Result<VideoConfig, Error> {
let encode_rate = rung.framerate.unwrap_or(moq_video::Rate::new(30, 1).unwrap());
let mut config = moq_video::encode::Config::new(rung.size.width, rung.size.height, encode_rate);
config.bitrate = Some(rung.bitrate);
config.kind = encoder.clone();
let mut entry = config.probe().await?;
entry.framerate = rung.framerate.map(moq_video::Rate::as_f64);
entry.optimize_for_latency = source.optimize_for_latency;
Ok(entry)
}
pub(crate) fn inherit_stalled(rungs: &mut [Published], source: &VideoConfig) {
for published in rungs {
published.entry.stalled = source.stalled;
}
}
pub(crate) fn populate(
out: &mut moq_mux::catalog::hang::Catalog,
source: &moq_mux::catalog::hang::Catalog,
rungs: &[Published],
source_rel: Option<&RelativeOwned>,
) -> Result<(), Error> {
out.video = Video::default();
out.audio = hang::catalog::Audio::default();
out.archive = source.archive.clone();
out.clock = source.clock;
out.video.display = source.video.display.clone();
out.video.rotation = source.video.rotation;
out.video.flip = source.video.flip;
for published in rungs {
out.video.insert(&published.rung.name, published.entry.clone())?;
}
let Some(rel) = source_rel else {
return Ok(());
};
for (name, config) in &source.video.renditions {
if config.broadcast.is_some() {
continue;
}
let mut config = config.clone();
config.broadcast = Some(rel.clone());
if out.video.insert(name, config).is_err() {
tracing::warn!(rendition = %name, "source video rendition collides with a rung name; skipping");
}
}
for (name, config) in &source.audio.renditions {
if config.broadcast.is_some() {
continue;
}
let mut config = config.clone();
config.broadcast = Some(rel.clone());
if out.audio.insert(name, config).is_err() {
tracing::warn!(rendition = %name, "duplicate source audio rendition; skipping");
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use hang::catalog::H264;
use super::*;
use crate::Rung;
fn source(width: u32, height: u32, bitrate: Option<u64>) -> VideoConfig {
let mut config = VideoConfig::new(H264 {
inline: true,
profile: 0x64,
constraints: 0,
level: 40,
});
config.coded_width = Some(width);
config.coded_height = Some(height);
config.bitrate = bitrate;
config.framerate = Some(30.0);
config
}
#[test]
fn rungs_inherit_a_stalled_source() {
let mut source_catalog = moq_mux::catalog::hang::Catalog::default();
let mut src = source(1280, 720, Some(2_500_000));
src.stalled = Some(true);
source_catalog.video.insert("video", src).unwrap();
let mut published = [Published {
rung: Resolved {
name: "video/360p".into(),
height: 360,
size: moq_video::Size::new(640, 360),
bitrate: moq_net::bandwidth::Rate::from_bps(600_000),
framerate: Some(moq_video::Rate::new(30, 1).unwrap()),
},
entry: source(640, 360, Some(600_000)),
}];
inherit_stalled(&mut published, &source_catalog.video.renditions["video"]);
let mut out = moq_mux::catalog::hang::Catalog::default();
populate(&mut out, &source_catalog, &published, None).unwrap();
assert_eq!(
out.video.renditions.get("video/360p").and_then(|c| c.stalled),
Some(true)
);
let healthy = source(1920, 1080, Some(5_000_000));
source_catalog.video.insert("healthy", healthy.clone()).unwrap();
inherit_stalled(&mut published, &healthy);
populate(&mut out, &source_catalog, &published, None).unwrap();
assert_eq!(out.video.renditions["video/360p"].stalled, None);
}
#[test]
fn rungs_never_upscale() {
let rungs = crate::Config::default().ladder;
let resolved = resolve_rungs(&rungs, "video", &source(854, 480, Some(2_000_000))).unwrap();
let names: Vec<_> = resolved.iter().map(|r| r.name.as_str()).collect();
assert_eq!(names, ["video/240p", "video/360p", "video/480p"]);
}
#[test]
fn fractional_and_unknown_source_rates_stay_explicit() {
let ladder = crate::Config::default().ladder;
let mut fractional = source(1280, 720, Some(2_500_000));
fractional.framerate = Some(30_000.0 / 1_001.0);
let resolved = resolve_rungs(&ladder, "video", &fractional).unwrap();
assert_eq!(
resolved[0].framerate,
Some(moq_video::Rate::new(30_000, 1_001).unwrap())
);
fractional.framerate = None;
let resolved = resolve_rungs(&ladder, "video", &fractional).unwrap();
assert_eq!(resolved[0].framerate, None);
}
#[test]
fn filtering_preserves_order_across_source_changes() {
let ladder = Ladder::new([
Rung::new(720, moq_net::bandwidth::Rate::from_bps(2_500_000)),
Rung::new(241, moq_net::bandwidth::Rate::from_bps(350_000)),
Rung::new(480, moq_net::bandwidth::Rate::from_bps(1_200_000)),
Rung::new(360, moq_net::bandwidth::Rate::from_bps(600_000)),
])
.unwrap();
for (picture, expected) in [
(source(1920, 1080, None), vec![240, 360, 480, 720]),
(source(1280, 720, Some(1_000_000)), vec![240, 360]),
(source(320, 180, None), vec![]),
(source(1920, 1080, Some(6_000_000)), vec![240, 360, 480, 720]),
] {
let resolved = resolve_rungs(&ladder, "video", &picture).unwrap();
assert_eq!(resolved.iter().map(|rung| rung.height).collect::<Vec<_>>(), expected);
assert!(resolved.windows(2).all(|pair| pair[0].bitrate < pair[1].bitrate));
}
}
#[test]
fn same_height_needs_lower_bitrate() {
let rungs = Ladder::new([Rung::new(480, moq_net::bandwidth::Rate::from_bps(1_200_000))]).unwrap();
assert!(
resolve_rungs(&rungs, "video", &source(854, 480, None))
.unwrap()
.is_empty()
);
assert!(
resolve_rungs(&rungs, "video", &source(854, 480, Some(1_000_000)))
.unwrap()
.is_empty()
);
}
#[test]
fn rung_geometry_follows_source_aspect() {
let resolved = resolve_rungs(
&Ladder::new([Rung::new(360, moq_net::bandwidth::Rate::from_bps(600_000))]).unwrap(),
"video",
&source(1920, 1080, Some(6_000_000)),
)
.unwrap();
assert_eq!(resolved.len(), 1);
assert_eq!(resolved[0].size, moq_video::Size::new(640, 360));
let resolved = resolve_rungs(
&Ladder::new([Rung::new(360, moq_net::bandwidth::Rate::from_bps(600_000))]).unwrap(),
"video",
&source(1080, 1920, Some(6_000_000)),
)
.unwrap();
assert_eq!(resolved[0].size, moq_video::Size::new(202, 360));
}
#[test]
fn dimensionless_rendition_is_not_a_source_yet() {
let mut provisional = source(0, 0, None);
provisional.coded_width = None;
provisional.coded_height = None;
let mut video = Video::default();
video.renditions.insert("video".to_string(), provisional.clone());
assert!(matches!(choose_source(&video), Err(Error::NoSource)));
video.renditions.insert("video".to_string(), source(1920, 1080, None));
let (name, chosen) = choose_source(&video).unwrap();
assert_eq!(name, "video");
assert_eq!(chosen.coded_width, Some(1920));
}
#[test]
fn source_needs_dimensions() {
let mut config = source(0, 0, Some(1_000_000));
config.coded_width = None;
config.coded_height = None;
assert!(matches!(
resolve_rungs(
&Ladder::new([Rung::new(360, moq_net::bandwidth::Rate::from_bps(600_000))]).unwrap(),
"video",
&config
),
Err(Error::SourceDimensions(_))
));
}
#[test]
fn a_new_description_is_a_new_decode_stream() {
let mut before = source(1920, 1080, Some(6_000_000));
let VideoCodec::H264(h264) = &mut before.codec else {
unreachable!()
};
h264.inline = false;
before.description = Some(bytes::Bytes::from_static(b"old avcC"));
let mut after = before.clone();
after.description = Some(bytes::Bytes::from_static(b"new avcC"));
assert!(!same_stream(&before, &after));
}
#[tokio::test]
async fn rung_entry_describes_the_rung() {
let rung = Resolved {
name: "video/360p".to_string(),
height: 360,
size: moq_video::Size::new(640, 360),
bitrate: moq_net::bandwidth::Rate::from_bps(600_000),
framerate: Some(moq_video::Rate::new(30, 1).unwrap()),
};
let mut source = source(1920, 1080, Some(6_000_000));
source.optimize_for_latency = Some(true);
let entry = rung_entry(&rung, &source, &moq_video::encode::Kind::Software)
.await
.unwrap();
let hang::catalog::VideoCodec::H264(h264) = &entry.codec else {
panic!("expected H.264, got {}", entry.codec)
};
assert!(h264.inline, "an avc3 rung carries its parameter sets in band");
assert_eq!(entry.coded_width, Some(640));
assert_eq!(entry.coded_height, Some(360));
assert_eq!(entry.bitrate, Some(600_000));
assert_eq!(entry.framerate, Some(30.0));
assert_eq!(entry.optimize_for_latency, Some(true));
}
#[tokio::test]
async fn an_unknown_rate_stays_unknown_after_probe() {
let rung = Resolved {
name: "video/360p".to_string(),
height: 360,
size: moq_video::Size::new(640, 360),
bitrate: moq_net::bandwidth::Rate::from_bps(600_000),
framerate: None,
};
let entry = rung_entry(
&rung,
&source(1920, 1080, Some(6_000_000)),
&moq_video::encode::Kind::Software,
)
.await
.unwrap();
assert_eq!(entry.framerate, None);
}
#[test]
fn names_are_never_reused() {
let mut names = Names::default();
assert_eq!(names.mint(360), "video/360p");
assert_eq!(names.mint(120), "video/120p");
assert_eq!(names.mint(360), "video/360p.2");
assert_eq!(names.mint(360), "video/360p.3");
assert_eq!(names.mint(120), "video/120p.2");
}
#[test]
fn chooses_highest_local_rendition() {
let mut video = Video::default();
video.insert("low", source(640, 360, None)).unwrap();
video.insert("high", source(1920, 1080, None)).unwrap();
let mut remote = source(3840, 2160, None);
remote.broadcast = Some(RelativeOwned::from("./other".to_string()));
video.insert("remote", remote).unwrap();
let (name, config) = choose_source(&video).unwrap();
assert_eq!(name, "high");
assert_eq!(config.coded_height, Some(1080));
}
#[test]
fn chooses_av1_source() {
let mut video = Video::default();
let mut av1 = VideoConfig::new(hang::catalog::AV1::default());
av1.coded_width = Some(1920);
av1.coded_height = Some(1080);
video.insert("av1", av1).unwrap();
let (name, config) = choose_source(&video).unwrap();
assert_eq!(name, "av1");
assert!(matches!(config.codec, VideoCodec::AV1(_)));
}
#[test]
fn skips_unsupported_av1_source() {
let mut video = Video::default();
let mut av1 = VideoConfig::new(hang::catalog::AV1 {
bitdepth: 10,
..hang::catalog::AV1::default()
});
av1.coded_width = Some(3840);
av1.coded_height = Some(2160);
video.insert("av1", av1).unwrap();
video.insert("h264", source(1920, 1080, None)).unwrap();
let (name, config) = choose_source(&video).unwrap();
assert_eq!(name, "h264");
assert!(matches!(config.codec, VideoCodec::H264(_)));
}
#[test]
fn populate_preserves_the_child_archive() {
let mut child = moq_mux::catalog::hang::Catalog::<()>::default();
child
.video
.insert("video", source(1920, 1080, Some(6_000_000)))
.unwrap();
let mut archive = hang::catalog::Archive::new("timeline.z");
archive.replay = Some(RelativeOwned::from("./recordings/clip".to_string()));
archive.version = Some(hang::catalog::Archive::VERSION);
child.archive = Some(archive.clone());
let clock = hang::catalog::Clock::new(moq_net::Timestamp::from_micros(1_751_846_400_000_000).unwrap()).unwrap();
child.clock = Some(clock);
let mut out = moq_mux::catalog::hang::Catalog::<()>::default();
populate(&mut out, &child, &[], None).unwrap();
assert_eq!(out.archive, Some(archive), "a derivative keeps the child's archive");
assert_eq!(out.clock, Some(clock), "a derivative keeps the child's clock");
}
}