use std::collections::BTreeMap;
use std::ops::{Deref, DerefMut};
use std::sync::{Arc, Mutex, MutexGuard};
use base64::Engine;
use super::hang::{Catalog, CatalogExt, Consumer, Extra};
#[derive(Default)]
struct Reservations {
reservers: usize,
pending: bool,
published: bool,
}
pub struct Producer<E: CatalogExt = ()> {
hang: moq_json::snapshot::Producer<Catalog<E>>,
hangz: moq_json::snapshot::Producer<Catalog<E>>,
msf_track: moq_net::track::Producer,
current: Arc<Mutex<Catalog<E>>>,
reservations: Arc<Mutex<Reservations>>,
clock: crate::Clock,
broadcast: moq_net::broadcast::Producer,
timelines: Arc<Mutex<BTreeMap<String, crate::timeline::Producer>>>,
}
impl<E: CatalogExt> Clone for Producer<E> {
fn clone(&self) -> Self {
Self {
hang: self.hang.clone(),
hangz: self.hangz.clone(),
msf_track: self.msf_track.clone(),
current: self.current.clone(),
reservations: self.reservations.clone(),
clock: self.clock,
broadcast: self.broadcast.clone(),
timelines: self.timelines.clone(),
}
}
}
impl Producer<()> {
pub fn new(broadcast: &mut moq_net::broadcast::Producer) -> Result<Self, moq_net::Error> {
Self::with_catalog(broadcast, Catalog::default())
}
}
impl<E: CatalogExt> Producer<E> {
pub fn with_catalog(
broadcast: &mut moq_net::broadcast::Producer,
catalog: Catalog<E>,
) -> Result<Self, moq_net::Error> {
let hang_track = broadcast.create_track(hang::Catalog::DEFAULT_NAME, hang::Catalog::default_track_info())?;
let hangz_track =
broadcast.create_track(hang::Catalog::COMPRESSED_NAME, hang::Catalog::default_track_info())?;
let msf_track = broadcast.create_track(moq_msf::DEFAULT_NAME, None)?;
let mut json_config = moq_json::snapshot::ProducerConfig::default();
json_config.delta_ratio = 0;
let hang = moq_json::snapshot::Producer::new(hang_track, json_config.clone());
json_config.compression = true;
let hangz = moq_json::snapshot::Producer::new(hangz_track, json_config);
Ok(Self {
hang,
hangz,
msf_track,
current: Arc::new(Mutex::new(catalog)),
reservations: Arc::new(Mutex::new(Reservations::default())),
clock: crate::Clock::new(),
broadcast: broadcast.clone(),
timelines: Arc::new(Mutex::new(BTreeMap::new())),
})
}
pub fn timestamp(&self, hint: Option<moq_net::Timestamp>) -> crate::Result<moq_net::Timestamp> {
match hint {
Some(pts) => Ok(pts),
None => Ok(moq_net::Timestamp::from_micros(self.clock.micros())?),
}
}
pub fn lock(&mut self) -> Guard<'_, E> {
Guard {
catalog: self.current.lock().unwrap(),
hang: &mut self.hang,
hangz: &mut self.hangz,
msf_track: &mut self.msf_track,
reservations: &self.reservations,
updated: false,
}
}
pub fn snapshot(&self) -> Catalog<E> {
self.current.lock().unwrap().clone()
}
pub fn reserve(&self) -> super::Reserved<E> {
super::Reserved::new(self.clone())
}
pub(super) fn add_reserver(&self) {
self.reservations.lock().unwrap().reservers += 1;
}
pub(super) fn release_reserver(&mut self) {
{
let mut r = self.reservations.lock().unwrap();
r.reservers = r.reservers.saturating_sub(1);
}
self.flush_if_ready();
}
pub(super) fn flush_if_ready(&mut self) {
{
let mut r = self.reservations.lock().unwrap();
if r.published || r.reservers != 0 {
return;
}
if !r.pending {
return;
}
r.pending = false;
r.published = true;
}
let catalog = self.current.lock().unwrap().clone();
if let Err(err) = emit(&mut self.hang, &mut self.hangz, &mut self.msf_track, &catalog) {
tracing::warn!(%err, "failed to publish the catalog");
}
}
pub fn media_producer<C: crate::container::Container>(
&self,
track: moq_net::track::Producer,
container: C,
) -> crate::Result<crate::container::Producer<C>> {
let recorder = self.timeline(track.name())?.recorder();
Ok(crate::container::Producer::new(track, container).with_recorder(recorder))
}
pub fn timeline(&self, name: &str) -> crate::Result<crate::timeline::Producer> {
let mut timelines = self.timelines.lock().unwrap();
if let Some(timeline) = timelines.get(name) {
return Ok(timeline.clone());
}
let timeline = crate::timeline::Producer::new(&mut self.broadcast.clone(), name)?;
timelines.insert(name.to_string(), timeline.clone());
Ok(timeline)
}
pub fn consume(&self) -> Result<Consumer<E>, moq_net::Error> {
Ok(Consumer::new(self.hang.consume()))
}
pub fn finish(&mut self) -> crate::Result<()> {
self.hang.finish()?;
self.hangz.finish()?;
self.msf_track.finish()?;
for timeline in self.timelines.lock().unwrap().values_mut() {
timeline.finish()?;
}
Ok(())
}
}
pub struct Guard<'a, E: CatalogExt = ()> {
catalog: MutexGuard<'a, Catalog<E>>,
hang: &'a mut moq_json::snapshot::Producer<Catalog<E>>,
hangz: &'a mut moq_json::snapshot::Producer<Catalog<E>>,
msf_track: &'a mut moq_net::track::Producer,
reservations: &'a Mutex<Reservations>,
updated: bool,
}
impl<E: CatalogExt> Guard<'_, E> {
pub fn commit(mut self) -> crate::Result<()> {
self.publish()
}
fn publish(&mut self) -> crate::Result<()> {
if !self.updated {
return Ok(());
}
self.updated = false;
{
let mut r = self.reservations.lock().unwrap();
if !r.published && r.reservers != 0 {
r.pending = true;
return Ok(());
}
r.pending = false;
r.published = true;
}
emit(self.hang, self.hangz, self.msf_track, &self.catalog)
}
}
impl<E: CatalogExt> Deref for Guard<'_, E> {
type Target = Catalog<E>;
fn deref(&self) -> &Self::Target {
&self.catalog
}
}
impl<E: CatalogExt> DerefMut for Guard<'_, E> {
fn deref_mut(&mut self) -> &mut Self::Target {
self.updated = true;
&mut self.catalog
}
}
impl Guard<'_, Extra> {
pub fn set_section(&mut self, name: impl Into<String>, value: serde_json::Value) -> crate::Result<()> {
self.catalog.ext.set(name, value)?;
self.updated = true;
Ok(())
}
pub fn remove_section(&mut self, name: &str) -> Option<serde_json::Value> {
let removed = self.catalog.ext.remove(name);
if removed.is_some() {
self.updated = true;
}
removed
}
}
impl<E: CatalogExt> Drop for Guard<'_, E> {
fn drop(&mut self) {
if let Err(err) = self.publish() {
tracing::warn!(%err, "failed to publish the catalog on guard drop");
}
}
}
fn emit<E: CatalogExt>(
hang: &mut moq_json::snapshot::Producer<Catalog<E>>,
hangz: &mut moq_json::snapshot::Producer<Catalog<E>>,
msf_track: &mut moq_net::track::Producer,
catalog: &Catalog<E>,
) -> crate::Result<()> {
hang.update(catalog)?;
hangz.update(catalog)?;
let msf = to_msf(&catalog.media()).to_json().map_err(moq_json::Error::from)?;
let mut group = msf_track.append_group()?;
group.write_frame(moq_net::Timestamp::now(), msf)?;
group.finish()?;
Ok(())
}
fn video_sap_type(codec: &hang::catalog::VideoCodec) -> Option<u8> {
use hang::catalog::VideoCodec;
match codec {
VideoCodec::VP8 => Some(1),
VideoCodec::H264(_) | VideoCodec::AV1(_) | VideoCodec::VP9(_) => Some(2),
_ => None,
}
}
fn to_msf(catalog: &hang::Catalog) -> moq_msf::Catalog {
let mut tracks = Vec::new();
let has_multiple_video = catalog.video.renditions.len() > 1;
for (name, config) in &catalog.video.renditions {
let packaging = match &config.container {
hang::catalog::Container::Cmaf { .. } => moq_msf::Packaging::Cmaf,
_ => moq_msf::Packaging::Legacy,
};
let init_data = match &config.container {
hang::catalog::Container::Cmaf { init, .. } => Some(base64::engine::general_purpose::STANDARD.encode(init)),
_ => config
.description
.as_ref()
.map(|d| base64::engine::general_purpose::STANDARD.encode(d.as_ref())),
};
let sap_type = video_sap_type(&config.codec);
let mut track = moq_msf::Track::new(name.clone(), packaging);
track.is_live = true;
track.role = Some(moq_msf::Role::Video);
track.codec = Some(config.codec.to_string());
track.width = config.coded_width;
track.height = config.coded_height;
track.framerate = config.framerate;
track.bitrate = config.bitrate;
track.init_data = init_data;
track.render_group = Some(1);
track.alt_group = if has_multiple_video { Some(1) } else { None };
track.max_grp_sap_starting_type = sap_type;
track.max_obj_sap_starting_type = sap_type;
track.jitter = config.jitter;
tracks.push(track);
}
let has_multiple_audio = catalog.audio.renditions.len() > 1;
for (name, config) in &catalog.audio.renditions {
let packaging = match &config.container {
hang::catalog::Container::Cmaf { .. } => moq_msf::Packaging::Cmaf,
_ => moq_msf::Packaging::Legacy,
};
let init_data = match &config.container {
hang::catalog::Container::Cmaf { init, .. } => Some(base64::engine::general_purpose::STANDARD.encode(init)),
_ => config
.description
.as_ref()
.map(|d| base64::engine::general_purpose::STANDARD.encode(d.as_ref())),
};
let mut track = moq_msf::Track::new(name.clone(), packaging);
track.is_live = true;
track.role = Some(moq_msf::Role::Audio);
track.codec = Some(config.codec.to_string());
track.samplerate = Some(config.sample_rate);
track.channel_config = Some(config.channel_count.to_string());
track.bitrate = config.bitrate;
track.init_data = init_data;
track.render_group = Some(1);
track.alt_group = if has_multiple_audio { Some(1) } else { None };
track.max_grp_sap_starting_type = Some(1);
track.max_obj_sap_starting_type = Some(1);
track.jitter = config.jitter;
tracks.push(track);
}
moq_msf::Catalog::new(tracks)
}
#[cfg(test)]
mod test {
use std::collections::BTreeMap;
use std::task::Poll;
use bytes::Bytes;
use hang::catalog::{AudioCodec, AudioConfig, Container, H264, VideoConfig};
use super::*;
#[test]
fn publishes_plain_and_compressed_tracks() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let mut catalog = Producer::new(&mut broadcast).unwrap();
let mut plain = Consumer::new(catalog.hang.consume());
let mut compressed = Consumer::compressed(catalog.hangz.consume());
{
let mut guard = catalog.lock();
guard
.audio
.renditions
.insert("audio0".to_string(), AudioConfig::new(AudioCodec::Opus, 48_000, 2));
}
let expected = catalog.snapshot();
let waiter = kio::Waiter::noop();
let got_plain = match plain.poll_next(&waiter) {
Poll::Ready(Ok(Some(c))) => c,
other => panic!("expected plain catalog, got {other:?}"),
};
let got_compressed = match compressed.poll_next(&waiter) {
Poll::Ready(Ok(Some(c))) => c,
other => panic!("expected compressed catalog, got {other:?}"),
};
assert_eq!(got_plain, expected);
assert_eq!(got_compressed, expected);
}
#[test]
fn commit_reports_a_publish_failure() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let mut catalog = Producer::new(&mut broadcast).unwrap();
catalog.finish().unwrap();
let mut guard = catalog.lock();
guard
.audio
.renditions
.insert("audio0".to_string(), AudioConfig::new(AudioCodec::Opus, 48_000, 2));
assert!(guard.commit().is_err());
}
#[test]
fn commit_publishes_once() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let mut catalog = Producer::new(&mut broadcast).unwrap();
let track = catalog.hang.consume();
let mut guard = catalog.lock();
guard
.audio
.renditions
.insert("audio0".to_string(), AudioConfig::new(AudioCodec::Opus, 48_000, 2));
guard.commit().unwrap();
catalog.finish().unwrap();
assert_eq!(track.latest(), Some(0));
}
#[test]
fn timeline_reports_a_track_collision() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = Producer::new(&mut broadcast).unwrap();
let _taken = broadcast.create_track("video0.timeline.z", None).unwrap();
assert!(catalog.timeline("video0").is_err());
}
fn h264_config() -> VideoConfig {
let mut config = VideoConfig::new(H264 {
profile: 0x64,
constraints: 0x00,
level: 0x1f,
inline: true,
});
config.container = Container::Legacy;
config
}
#[test]
fn reservation_gates_until_all_renditions_resolve() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = Producer::new(&mut broadcast).unwrap();
let mut consumer: Consumer = Consumer::new(catalog.hang.consume());
let waiter = kio::Waiter::noop();
let reserved = catalog.reserve();
let mut audio = reserved.audio("audio0");
let mut video = reserved.video("video0");
drop(reserved);
audio.set(AudioConfig::new(AudioCodec::Opus, 48_000, 2));
assert!(
matches!(consumer.poll_next(&waiter), Poll::Pending),
"an audio-only catalog must not publish while video is unresolved"
);
video.set(h264_config());
let snapshot = match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(c))) => c,
other => panic!("expected the complete catalog, got {other:?}"),
};
assert!(snapshot.audio.renditions.contains_key("audio0"));
assert!(snapshot.video.renditions.contains_key("video0"));
assert!(
matches!(consumer.poll_next(&waiter), Poll::Pending),
"only the one complete snapshot should be published"
);
}
#[test]
fn reservation_gate_opens_when_unresolved_reservation_is_dropped() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = Producer::new(&mut broadcast).unwrap();
let mut consumer: Consumer = Consumer::new(catalog.hang.consume());
let waiter = kio::Waiter::noop();
let reserved = catalog.reserve();
let mut audio = reserved.audio("audio0");
let video = reserved.video("video0");
drop(reserved);
audio.set(AudioConfig::new(AudioCodec::Opus, 48_000, 2));
assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
drop(video);
let snapshot = match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(c))) => c,
other => panic!("expected the audio catalog, got {other:?}"),
};
assert!(snapshot.audio.renditions.contains_key("audio0"));
assert!(!snapshot.video.renditions.contains_key("video0"));
}
#[test]
fn staged_change_waits_for_a_held_reservation() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = Producer::new(&mut broadcast).unwrap();
let mut consumer: Consumer = Consumer::new(catalog.hang.consume());
let waiter = kio::Waiter::noop();
let deferred = catalog.reserve();
let early = catalog.reserve();
let mut a0 = early.audio("audio0");
a0.set(AudioConfig::new(AudioCodec::Opus, 48_000, 2));
drop(early);
assert!(
matches!(consumer.poll_next(&waiter), Poll::Pending),
"the staged rendition must wait for the deferred importer's held reservation"
);
let mut late = deferred.audio("audio1");
drop(deferred); assert!(matches!(consumer.poll_next(&waiter), Poll::Pending));
late.set(AudioConfig::new(AudioCodec::Opus, 48_000, 1));
let snapshot = match consumer.poll_next(&waiter) {
Poll::Ready(Ok(Some(c))) => c,
other => panic!("expected one complete catalog, got {other:?}"),
};
assert!(snapshot.audio.renditions.contains_key("audio0"));
assert!(snapshot.audio.renditions.contains_key("audio1"));
assert!(
matches!(consumer.poll_next(&waiter), Poll::Pending),
"only the one complete snapshot should be published"
);
}
#[test]
fn convert_simple() {
let mut video_config = VideoConfig::new(H264 {
profile: 0x64,
constraints: 0x00,
level: 0x1f,
inline: true,
});
video_config.coded_width = Some(1280);
video_config.coded_height = Some(720);
video_config.bitrate = Some(6_000_000);
video_config.framerate = Some(30.0);
video_config.container = Container::Legacy;
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0.avc3".to_string(), video_config);
let mut audio_config = AudioConfig::new(AudioCodec::Opus, 48_000, 2);
audio_config.bitrate = Some(128_000);
audio_config.container = Container::Legacy;
let mut audio_renditions = BTreeMap::new();
audio_renditions.insert("audio0".to_string(), audio_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
catalog.audio.renditions = audio_renditions;
let msf = to_msf(&catalog);
assert_eq!(msf.tracks.len(), 2);
let video = &msf.tracks[0];
assert_eq!(video.name, "video0.avc3");
assert_eq!(video.role, Some(moq_msf::Role::Video));
assert_eq!(video.packaging, moq_msf::Packaging::Legacy);
assert_eq!(video.codec, Some("avc3.64001f".to_string()));
assert_eq!(video.width, Some(1280));
assert_eq!(video.height, Some(720));
assert_eq!(video.framerate, Some(30.0));
assert_eq!(video.bitrate, Some(6_000_000));
assert!(video.init_data.is_none());
assert_eq!(video.max_grp_sap_starting_type, Some(2));
assert_eq!(video.max_obj_sap_starting_type, Some(2));
assert_eq!(video.jitter, None);
let audio = &msf.tracks[1];
assert_eq!(audio.name, "audio0");
assert_eq!(audio.role, Some(moq_msf::Role::Audio));
assert_eq!(audio.packaging, moq_msf::Packaging::Legacy);
assert_eq!(audio.codec, Some("opus".to_string()));
assert_eq!(audio.samplerate, Some(48_000));
assert_eq!(audio.channel_config, Some("2".to_string()));
assert_eq!(audio.bitrate, Some(128_000));
assert_eq!(audio.max_grp_sap_starting_type, Some(1));
assert_eq!(audio.max_obj_sap_starting_type, Some(1));
assert_eq!(audio.jitter, None);
}
#[test]
fn convert_with_description() {
let mut video_config = VideoConfig::new(H264 {
profile: 0x64,
constraints: 0x00,
level: 0x1f,
inline: false,
});
video_config.description = Some(Bytes::from_static(&[0x01, 0x02, 0x03]));
video_config.coded_width = Some(1920);
video_config.coded_height = Some(1080);
video_config.container = Container::Legacy;
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0.m4s".to_string(), video_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
let msf = to_msf(&catalog);
let video = &msf.tracks[0];
assert_eq!(video.init_data, Some("AQID".to_string()));
}
#[test]
fn convert_empty() {
let catalog = hang::Catalog::default();
let msf = to_msf(&catalog);
assert!(msf.tracks.is_empty());
}
#[test]
fn convert_cmaf_packaging() {
let mut video_config = VideoConfig::new(H264 {
profile: 0x64,
constraints: 0x00,
level: 0x28,
inline: false,
});
video_config.coded_width = Some(1920);
video_config.coded_height = Some(1080);
video_config.container = Container::Cmaf {
init: base64::engine::general_purpose::STANDARD
.decode("AAAYZ2Z0eXA=")
.unwrap()
.into(),
};
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0.m4s".to_string(), video_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
let msf = to_msf(&catalog);
let video = &msf.tracks[0];
assert_eq!(video.packaging, moq_msf::Packaging::Cmaf);
assert_eq!(video.init_data, Some("AAAYZ2Z0eXA=".to_string()));
}
#[test]
fn convert_sap_h264_with_jitter() {
let mut video_config = VideoConfig::new(H264 {
profile: 0x64,
constraints: 0x00,
level: 0x1f,
inline: true,
});
video_config.coded_width = Some(1280);
video_config.coded_height = Some(720);
video_config.framerate = Some(30.0);
video_config.container = Container::Legacy;
video_config.jitter = Some(std::time::Duration::from_millis(100));
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0".to_string(), video_config);
let mut audio_config = AudioConfig::new(AudioCodec::Opus, 48_000, 2);
audio_config.container = Container::Legacy;
audio_config.jitter = Some(std::time::Duration::from_millis(40));
let mut audio_renditions = BTreeMap::new();
audio_renditions.insert("audio0".to_string(), audio_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
catalog.audio.renditions = audio_renditions;
let msf = to_msf(&catalog);
let video = &msf.tracks[0];
assert_eq!(video.role, Some(moq_msf::Role::Video));
assert_eq!(video.max_grp_sap_starting_type, Some(2));
assert_eq!(video.max_obj_sap_starting_type, Some(2));
assert_eq!(video.jitter, Some(std::time::Duration::from_millis(100)));
let audio = &msf.tracks[1];
assert_eq!(audio.role, Some(moq_msf::Role::Audio));
assert_eq!(audio.max_grp_sap_starting_type, Some(1));
assert_eq!(audio.max_obj_sap_starting_type, Some(1));
assert_eq!(audio.jitter, Some(std::time::Duration::from_millis(40)));
}
#[test]
fn convert_sap_h265() {
use hang::catalog::H265;
let mut video_config = VideoConfig::new(H265 {
in_band: false,
profile_space: 0,
profile_idc: 1,
profile_compatibility_flags: [0, 0, 0, 0],
tier_flag: false,
level_idc: 93,
constraint_flags: [0, 0, 0, 0, 0, 0],
});
video_config.coded_width = Some(1920);
video_config.coded_height = Some(1080);
video_config.framerate = Some(60.0);
video_config.container = Container::Legacy;
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0".to_string(), video_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
let msf = to_msf(&catalog);
let video = &msf.tracks[0];
assert_eq!(video.max_grp_sap_starting_type, None);
assert_eq!(video.max_obj_sap_starting_type, None);
assert_eq!(video.jitter, None);
}
#[test]
fn convert_sap_unknown_codec() {
use hang::catalog::VideoCodec;
let mut video_config = VideoConfig::new(VideoCodec::Unknown("future-codec.01".to_string()));
video_config.container = Container::Legacy;
let mut video_renditions = BTreeMap::new();
video_renditions.insert("video0".to_string(), video_config);
let mut catalog = hang::Catalog::default();
catalog.video.renditions = video_renditions;
let msf = to_msf(&catalog);
let video = &msf.tracks[0];
assert_eq!(video.max_grp_sap_starting_type, None);
assert_eq!(video.max_obj_sap_starting_type, None);
}
}