use hang::catalog::{AV1, Video, VideoCodec, VideoConfig};
use moq_net::path::RelativeOwned;
use moq_video::decode::Codec;
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) struct Decoders {
config: moq_video::decode::Config,
probed: Vec<(Codec, bool)>,
}
impl Decoders {
pub fn new(config: moq_video::decode::Config) -> Self {
Self {
config,
probed: Vec::new(),
}
}
#[cfg(test)]
fn assume(codecs: &[Codec]) -> Self {
let probed = [Codec::H264, Codec::H265, Codec::Av1]
.into_iter()
.map(|codec| (codec, codecs.contains(&codec)))
.collect();
Self {
config: moq_video::decode::Config::new(),
probed,
}
}
async fn probe(&mut self, codec: Codec, rendition: &VideoConfig) -> bool {
if let Some((_, decodes)) = self.probed.iter().find(|(probed, _)| *probed == codec) {
return *decodes;
}
let decodes = match self.open(rendition).await {
Ok(()) => true,
Err(err) => {
tracing::warn!(?codec, %err, "no decoder for this codec; its renditions will not be transcoded");
false
}
};
self.probed.push((codec, decodes));
decodes
}
async fn refusal(&mut self, rendition: &VideoConfig) -> Result<(), Error> {
match self.open(rendition).await {
Err(err) => Err(err.into()),
Ok(()) => {
if let Some(codec) = codec(rendition)
&& let Some((_, decodes)) = self.probed.iter_mut().find(|(probed, _)| *probed == codec)
{
*decodes = true;
}
Ok(())
}
}
}
async fn open(&self, rendition: &VideoConfig) -> Result<(), moq_video::Error> {
let mut config = rendition.clone();
config.description = None;
match &mut config.codec {
VideoCodec::H264(h264) => h264.inline = true,
VideoCodec::H265(h265) => h265.in_band = true,
_ => {}
}
moq_video::decode::Sink::open(&config, &self.config).await.map(drop)
}
}
pub(crate) async fn choose_source(video: &Video, decoders: &mut Decoders) -> Result<(String, VideoConfig), Error> {
let mut candidates: Vec<_> = video
.renditions
.iter()
.filter(|(_, config)| config.broadcast.is_none())
.filter_map(|(name, config)| Some((name, config, codec(config)?)))
.collect();
candidates.reverse();
candidates.sort_by_key(|(_, config, _)| {
std::cmp::Reverse((
dimensions(config).is_some(),
config.coded_height,
config.coded_width,
config.bitrate,
))
});
let mut refused = None;
for (name, config, codec) in candidates {
match decoders.probe(codec, config).await {
true if dimensions(config).is_some() => return Ok((name.clone(), config.clone())),
true => return Err(Error::NoSource),
false => {
refused.get_or_insert((name, config));
}
}
}
let Some((name, config)) = refused else {
return Err(Error::NoSource);
};
match decoders.refusal(config).await {
Err(err) => {
tracing::warn!(rendition = %name, "no source rendition can be decoded on this host");
Err(err)
}
Ok(()) if dimensions(config).is_some() => Ok((name.clone(), config.clone())),
Ok(()) => Err(Error::NoSource),
}
}
pub(crate) async fn follow_source(
video: &Video,
current: &str,
decoders: &mut Decoders,
) -> Result<(String, VideoConfig), Error> {
if let Some(config) = video.renditions.get(current)
&& config.broadcast.is_none()
&& dimensions(config).is_some()
&& let Some(codec) = codec(config)
&& decoders.probe(codec, config).await
{
return Ok((current.to_string(), config.clone()));
}
choose_source(video, decoders).await
}
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 codec(config: &VideoConfig) -> Option<Codec> {
match &config.codec {
VideoCodec::H264(_) => Some(Codec::H264),
VideoCodec::H265(_) => Some(Codec::H265),
VideoCodec::AV1(av1) if is_supported_av1(av1) => Some(Codec::Av1),
_ => None,
}
}
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));
}
fn decodes_all() -> Decoders {
Decoders::assume(&[Codec::H264, Codec::H265, Codec::Av1])
}
fn hevc(width: u32, height: u32) -> VideoConfig {
let mut config = 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],
});
config.coded_width = Some(width);
config.coded_height = Some(height);
config
}
#[tokio::test]
async 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, &mut decodes_all()).await,
Err(Error::NoSource)
));
video.renditions.insert("video".to_string(), source(1920, 1080, None));
let (name, chosen) = choose_source(&video, &mut decodes_all()).await.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");
}
#[tokio::test]
async 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, &mut decodes_all()).await.unwrap();
assert_eq!(name, "high");
assert_eq!(config.coded_height, Some(1080));
}
#[tokio::test]
async 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, &mut decodes_all()).await.unwrap();
assert_eq!(name, "av1");
assert!(matches!(config.codec, VideoCodec::AV1(_)));
}
#[tokio::test]
async 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, &mut decodes_all()).await.unwrap();
assert_eq!(name, "h264");
assert!(matches!(config.codec, VideoCodec::H264(_)));
}
#[tokio::test]
async fn skips_a_rendition_this_host_cannot_decode() {
let mut video = Video::default();
video.insert("hevc", hevc(1920, 1080)).unwrap();
video.insert("avc", source(640, 360, None)).unwrap();
let mut av1 = VideoConfig::new(hang::catalog::AV1::default());
av1.coded_width = Some(3840);
av1.coded_height = Some(2160);
video.insert("av1", av1).unwrap();
let mut decoders = Decoders::assume(&[Codec::H264]);
let (name, _) = choose_source(&video, &mut decoders).await.unwrap();
assert_eq!(name, "avc");
let (name, _) = choose_source(&video, &mut Decoders::assume(&[Codec::H264, Codec::H265]))
.await
.unwrap();
assert_eq!(name, "hevc");
}
#[tokio::test]
async fn refuses_when_no_rendition_decodes() {
let mut video = Video::default();
video.insert("small", hevc(640, 360)).unwrap();
video.insert("large", hevc(1920, 1080)).unwrap();
let mut config = moq_video::decode::Config::new();
config.kind = moq_video::decode::Kind::Named("missing".to_string());
match choose_source(&video, &mut Decoders::new(config)).await {
Err(Error::Video(moq_video::Error::UnknownDecoder { name, codec, .. })) => {
assert_eq!(name, "missing");
assert_eq!(codec, Codec::H265);
}
other => panic!("expected the decoder's refusal, got {other:?}"),
}
}
#[tokio::test]
async fn a_reopen_that_succeeds_is_the_source() {
let mut video = Video::default();
video.insert("avc", source(640, 360, None)).unwrap();
let mut decoders = Decoders::assume(&[]);
decoders.config.kind = moq_video::decode::Kind::Software;
let (name, _) = choose_source(&video, &mut decoders).await.unwrap();
assert_eq!(name, "avc");
assert!(decoders.probe(Codec::H264, &source(640, 360, None)).await);
}
#[tokio::test]
async fn a_pending_decodable_rendition_defers_the_refusal() {
let mut pending = source(0, 0, None);
pending.coded_width = None;
pending.coded_height = None;
let mut video = Video::default();
video.insert("hevc", hevc(1920, 1080)).unwrap();
video.insert("avc", pending).unwrap();
let mut decoders = Decoders::assume(&[Codec::H264]);
assert!(matches!(
choose_source(&video, &mut decoders).await,
Err(Error::NoSource)
));
video.renditions.insert("avc".to_string(), source(640, 360, None));
let (name, _) = choose_source(&video, &mut decoders).await.unwrap();
assert_eq!(name, "avc");
}
#[tokio::test]
async fn follow_keeps_a_valid_source() {
let mut video = Video::default();
video.insert("avc", source(640, 360, None)).unwrap();
let mut decoders = Decoders::assume(&[Codec::H264]);
video.insert("tall", source(1920, 1080, None)).unwrap();
video.insert("hevc", hevc(3840, 2160)).unwrap();
let (name, _) = follow_source(&video, "avc", &mut decoders).await.unwrap();
assert_eq!(name, "avc");
video.renditions.insert("avc".to_string(), hevc(640, 360));
let (name, _) = follow_source(&video, "avc", &mut decoders).await.unwrap();
assert_eq!(name, "tall");
}
#[tokio::test]
async fn each_codec_is_probed_once() {
let mut config = moq_video::decode::Config::new();
config.kind = moq_video::decode::Kind::Named("missing".to_string());
let mut decoders = Decoders::new(config);
for height in [360, 720, 1080] {
assert!(!decoders.probe(Codec::H265, &hevc(height * 16 / 9, height)).await);
}
assert_eq!(decoders.probed, [(Codec::H265, false)]);
}
#[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");
}
}