use std::marker::PhantomData;
use std::time::Duration;
use moq_net::Timestamp;
use super::Producer;
use super::hang::{Catalog, CatalogExt};
use crate::container::jitter::Metrics;
#[derive(Clone, Default, Debug, PartialEq)]
#[non_exhaustive]
pub struct Estimate {
pub jitter: Option<Duration>,
pub bitrate: Option<u64>,
}
impl Estimate {
pub fn with_jitter(mut self, jitter: impl Into<Option<Duration>>) -> Self {
self.jitter = jitter.into();
self
}
pub fn with_bitrate(mut self, bitrate: impl Into<Option<u64>>) -> Self {
self.bitrate = bitrate.into();
self
}
}
pub trait RenditionConfig<E: CatalogExt>: Sized + 'static {
fn insert(self, catalog: &mut Catalog<E>, name: &str);
fn get_mut<'a>(catalog: &'a mut Catalog<E>, name: &str) -> Option<&'a mut Self>;
fn remove(catalog: &mut Catalog<E>, name: &str);
fn estimate(&self) -> Estimate {
Estimate::default()
}
fn set_estimate(&mut self, _estimate: Estimate) {}
}
#[derive(Clone, Default, Debug, PartialEq)]
#[non_exhaustive]
pub struct VideoHint {
pub codec: Option<hang::catalog::VideoCodec>,
pub coded_width: Option<u32>,
pub coded_height: Option<u32>,
pub display_aspect_width: Option<u32>,
pub display_aspect_height: Option<u32>,
pub bitrate: Option<u64>,
pub framerate: Option<f64>,
pub optimize_for_latency: Option<bool>,
pub jitter: Option<Duration>,
}
fn fill<T>(slot: &mut Option<T>, value: Option<T>) {
if slot.is_none() {
*slot = value;
}
}
impl VideoHint {
pub fn apply(&self, config: &mut hang::catalog::VideoConfig) {
fill(&mut config.coded_width, self.coded_width);
fill(&mut config.coded_height, self.coded_height);
fill(&mut config.display_aspect_width, self.display_aspect_width);
fill(&mut config.display_aspect_height, self.display_aspect_height);
fill(&mut config.bitrate, self.bitrate);
fill(&mut config.framerate, self.framerate);
fill(&mut config.optimize_for_latency, self.optimize_for_latency);
fill(&mut config.jitter, self.jitter);
}
pub fn to_config(&self) -> Option<hang::catalog::VideoConfig> {
let codec = self.codec.clone()?;
let mut config = hang::catalog::VideoConfig::new(codec);
config.container = hang::catalog::Container::Legacy;
self.apply(&mut config);
Some(config)
}
}
impl<E: CatalogExt> RenditionConfig<E> for hang::catalog::VideoConfig {
fn insert(self, catalog: &mut Catalog<E>, name: &str) {
catalog.video.renditions.insert(name.to_string(), self);
}
fn get_mut<'a>(catalog: &'a mut Catalog<E>, name: &str) -> Option<&'a mut Self> {
catalog.video.renditions.get_mut(name)
}
fn remove(catalog: &mut Catalog<E>, name: &str) {
catalog.video.renditions.remove(name);
}
fn estimate(&self) -> Estimate {
Estimate::default().with_jitter(self.jitter).with_bitrate(self.bitrate)
}
fn set_estimate(&mut self, estimate: Estimate) {
self.jitter = estimate.jitter;
self.bitrate = estimate.bitrate;
}
}
impl<E: CatalogExt> RenditionConfig<E> for hang::catalog::AudioConfig {
fn insert(self, catalog: &mut Catalog<E>, name: &str) {
catalog.audio.renditions.insert(name.to_string(), self);
}
fn get_mut<'a>(catalog: &'a mut Catalog<E>, name: &str) -> Option<&'a mut Self> {
catalog.audio.renditions.get_mut(name)
}
fn remove(catalog: &mut Catalog<E>, name: &str) {
catalog.audio.renditions.remove(name);
}
fn estimate(&self) -> Estimate {
Estimate::default().with_jitter(self.jitter).with_bitrate(self.bitrate)
}
fn set_estimate(&mut self, estimate: Estimate) {
self.jitter = estimate.jitter;
self.bitrate = estimate.bitrate;
}
}
pub struct Reserved<E: CatalogExt = ()> {
catalog: Producer<E>,
}
impl<E: CatalogExt> Reserved<E> {
pub(super) fn new(catalog: Producer<E>) -> Self {
catalog.add_reserver();
Self { catalog }
}
pub fn init<C: RenditionConfig<E>>(&self, name: impl Into<String>) -> Rendition<E, C> {
Rendition::new(self.clone(), name)
}
pub fn video(&self, name: impl Into<String>) -> VideoTrack<E> {
self.init(name)
}
pub fn audio(&self, name: impl Into<String>) -> AudioTrack<E> {
self.init(name)
}
pub fn timestamp(&self, hint: Option<moq_net::Timestamp>) -> crate::Result<moq_net::Timestamp> {
self.catalog.timestamp(hint)
}
pub fn producer(&self) -> Producer<E> {
self.catalog.clone()
}
}
impl<E: CatalogExt> Clone for Reserved<E> {
fn clone(&self) -> Self {
self.catalog.add_reserver();
Self {
catalog: self.catalog.clone(),
}
}
}
impl<E: CatalogExt> Drop for Reserved<E> {
fn drop(&mut self) {
self.catalog.release_reserver();
}
}
pub struct Rendition<E: CatalogExt, C: RenditionConfig<E>> {
catalog: Producer<E>,
name: String,
gate: Option<Reserved<E>>,
present: bool,
metrics: Metrics,
supplied: Estimate,
_config: PhantomData<fn() -> C>,
}
pub type VideoTrack<E = ()> = Rendition<E, hang::catalog::VideoConfig>;
pub type AudioTrack<E = ()> = Rendition<E, hang::catalog::AudioConfig>;
impl<E: CatalogExt, C: RenditionConfig<E>> Rendition<E, C> {
fn new(reserved: Reserved<E>, name: impl Into<String>) -> Self {
Self {
catalog: reserved.catalog.clone(),
gate: Some(reserved),
name: name.into(),
present: false,
metrics: Metrics::new(),
supplied: Estimate::default(),
_config: PhantomData,
}
}
pub fn name(&self) -> &str {
&self.name
}
pub fn timestamp(&self, hint: Option<moq_net::Timestamp>) -> crate::Result<moq_net::Timestamp> {
self.catalog.timestamp(hint)
}
pub fn set(&mut self, mut config: C) {
self.supplied = config.estimate();
config.set_estimate(self.resolved());
{
let mut guard = self.catalog.lock();
config.insert(&mut guard, &self.name);
}
self.present = true;
self.gate = None;
}
fn resolved(&self) -> Estimate {
let mut estimate = self.supplied.clone();
if estimate.jitter.is_none() {
estimate.jitter = self.metrics.jitter();
}
if estimate.bitrate.is_none() {
estimate.bitrate = self.metrics.bitrate();
}
estimate
}
fn refresh(&mut self) {
let estimate = self.resolved();
self.update(|config| config.set_estimate(estimate));
}
pub fn update(&mut self, f: impl FnOnce(&mut C)) {
if !self.present {
return;
}
let mut guard = self.catalog.lock();
if let Some(config) = C::get_mut(&mut guard, &self.name) {
f(config);
}
}
pub fn record_frame(&mut self, ts: Timestamp, bytes: usize) {
if self.metrics.record_frame(ts, bytes).is_some() && self.supplied.jitter.is_none() {
self.refresh();
}
}
pub fn record_reorder(&mut self, reorder: Timestamp) {
if self.metrics.record_reorder(reorder).is_some() && self.supplied.jitter.is_none() {
self.refresh();
}
}
pub fn record_group_end(&mut self, next: Option<Timestamp>) {
if self.metrics.finish_group(next).is_some() && self.supplied.bitrate.is_none() {
self.refresh();
}
}
}
impl<E: CatalogExt, C: RenditionConfig<E>> Drop for Rendition<E, C> {
fn drop(&mut self) {
if self.present {
let mut guard = self.catalog.lock();
C::remove(&mut guard, &self.name);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn video_track() -> (moq_net::broadcast::Producer, super::super::Producer, VideoTrack) {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = super::super::Producer::new(&mut broadcast).unwrap();
let reserved = catalog.reserve();
let rendition = reserved.video("v");
drop(reserved);
(broadcast, catalog, rendition)
}
fn config(bitrate: Option<u64>, jitter: Option<Duration>) -> hang::catalog::VideoConfig {
let mut config = hang::catalog::VideoConfig::new(hang::catalog::VideoCodec::VP8);
config.bitrate = bitrate;
config.jitter = jitter;
config
}
fn ts(micros: u64) -> Timestamp {
Timestamp::from_micros(micros).unwrap()
}
fn feed<E: CatalogExt, C: RenditionConfig<E>>(rendition: &mut Rendition<E, C>) {
for i in 0..60u64 {
let t = ts(i * 40_000);
rendition.record_group_end(Some(t));
rendition.record_frame(t, 100_000);
}
rendition.record_group_end(None);
}
#[test]
fn detects_absent_jitter_and_bitrate() {
let (_broadcast, catalog, mut rendition) = video_track();
rendition.set(config(None, None));
feed(&mut rendition);
let snapshot = catalog.snapshot();
let config = snapshot.video.renditions.get("v").unwrap();
assert!(config.jitter.is_some(), "absent jitter should be auto-detected");
assert!(config.bitrate.is_some(), "absent bitrate should be auto-detected");
}
#[test]
fn keeps_provided_jitter_and_bitrate() {
let (_broadcast, catalog, mut rendition) = video_track();
rendition.set(config(Some(123), Some(Duration::from_millis(50))));
feed(&mut rendition);
let snapshot = catalog.snapshot();
let config = snapshot.video.renditions.get("v").unwrap();
assert_eq!(config.bitrate, Some(123), "a provided bitrate must not be overwritten");
assert_eq!(
config.jitter,
Some(Duration::from_millis(50)),
"a provided jitter must not be overwritten"
);
}
#[test]
fn hinted_fields_count_as_supplied() {
let (_broadcast, catalog, mut rendition) = video_track();
let hint = VideoHint {
bitrate: Some(456),
..Default::default()
};
let mut config = config(None, None);
hint.apply(&mut config);
rendition.set(config);
feed(&mut rendition);
let snapshot = catalog.snapshot();
let config = snapshot.video.renditions.get("v").unwrap();
assert_eq!(config.bitrate, Some(456), "a hinted bitrate must not be overwritten");
assert!(config.jitter.is_some(), "the unhinted jitter should still be detected");
}
#[test]
fn resetting_recaptures_supplied() {
let (_broadcast, catalog, mut rendition) = video_track();
rendition.set(config(None, None));
feed(&mut rendition);
assert!(catalog.snapshot().video.renditions.get("v").unwrap().bitrate.is_some());
rendition.set(config(Some(789), None));
feed(&mut rendition);
let snapshot = catalog.snapshot();
let config = snapshot.video.renditions.get("v").unwrap();
assert_eq!(config.bitrate, Some(789), "the re-set bitrate is now authoritative");
}
#[test]
fn renditions_share_a_timeline() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = super::super::Producer::new(&mut broadcast).unwrap();
let reserved = catalog.reserve();
let shared = catalog.timeline("video").unwrap();
let mut source = reserved.video("video0");
let mut rung = reserved.video("video1");
drop(reserved);
for rendition in [&mut source, &mut rung] {
let mut config = config(None, None);
config.timeline = Some(shared.section());
rendition.set(config);
}
let snapshot = catalog.snapshot();
let source_tl = snapshot
.video
.renditions
.get("video0")
.unwrap()
.timeline
.as_ref()
.unwrap();
let rung_tl = snapshot
.video
.renditions
.get("video1")
.unwrap()
.timeline
.as_ref()
.unwrap();
assert_eq!(source_tl.track, "video.timeline.z");
assert_eq!(
rung_tl.track, source_tl.track,
"both renditions index off one shared timeline track"
);
}
mod custom {
use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use super::*;
#[derive(Serialize, Deserialize, Clone, Default, Debug, PartialEq)]
struct TelemetryExt {
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
telemetry: BTreeMap<String, Telemetry>,
}
impl CatalogExt for TelemetryExt {}
#[derive(Serialize, Deserialize, Clone, Default, Debug, PartialEq)]
struct Telemetry {
schema: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
bitrate: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
timeline: Option<hang::catalog::Timeline>,
}
impl RenditionConfig<TelemetryExt> for Telemetry {
fn insert(self, catalog: &mut Catalog<TelemetryExt>, name: &str) {
catalog.telemetry.insert(name.to_string(), self);
}
fn get_mut<'a>(catalog: &'a mut Catalog<TelemetryExt>, name: &str) -> Option<&'a mut Self> {
catalog.telemetry.get_mut(name)
}
fn remove(catalog: &mut Catalog<TelemetryExt>, name: &str) {
catalog.telemetry.remove(name);
}
fn estimate(&self) -> Estimate {
Estimate::default().with_bitrate(self.bitrate)
}
fn set_estimate(&mut self, estimate: Estimate) {
self.bitrate = estimate.bitrate;
}
}
fn telemetry(bitrate: Option<u64>) -> Telemetry {
Telemetry {
schema: "gps/v1".to_string(),
bitrate,
timeline: None,
}
}
fn produce() -> (moq_net::broadcast::Producer, crate::catalog::Producer<TelemetryExt>) {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = crate::catalog::Producer::with_catalog(&mut broadcast, Catalog::default()).unwrap();
(broadcast, catalog)
}
#[test]
fn detects_and_advertises() {
let (_broadcast, catalog) = produce();
let reserved = catalog.reserve();
let mut rendition = reserved.init::<Telemetry>("gps");
drop(reserved);
let mut config = telemetry(None);
config.timeline = Some(catalog.timeline("gps").unwrap().section());
rendition.set(config);
feed(&mut rendition);
let snapshot = catalog.snapshot();
let config = snapshot.telemetry.get("gps").unwrap();
assert!(config.bitrate.is_some(), "absent bitrate should be auto-detected");
assert_eq!(
config.timeline.as_ref().map(|t| t.track.as_str()),
Some("gps.timeline.z"),
"the advertised timeline names the companion track"
);
drop(rendition);
assert!(
!catalog.snapshot().telemetry.contains_key("gps"),
"the rendition should be removed on drop"
);
}
#[test]
fn keeps_supplied_bitrate() {
let (_broadcast, catalog) = produce();
let reserved = catalog.reserve();
let mut rendition = reserved.init::<Telemetry>("gps");
drop(reserved);
rendition.set(telemetry(Some(4_200)));
feed(&mut rendition);
let snapshot = catalog.snapshot();
assert_eq!(snapshot.telemetry.get("gps").unwrap().bitrate, Some(4_200));
}
#[test]
fn gates_with_media_tracks() {
let (_broadcast, catalog) = produce();
let mut consumer = catalog.consume().unwrap();
let reserved = catalog.reserve();
let mut video = reserved.video("v");
let mut gps = reserved.init::<Telemetry>("gps");
drop(reserved);
let waiter = kio::Waiter::noop();
video.set(config(None, None));
assert!(
matches!(consumer.poll_next(&waiter), std::task::Poll::Pending),
"the catalog stays withheld while the telemetry rendition is unresolved"
);
gps.set(telemetry(None));
let mut latest = None;
while let std::task::Poll::Ready(Ok(Some(catalog))) = consumer.poll_next(&waiter) {
latest = Some(catalog);
}
let published = latest.expect("catalog published");
assert!(published.video.renditions.contains_key("v"));
assert!(published.telemetry.contains_key("gps"));
}
}
}