use moq_net::AsPath;
#[derive(Clone)]
pub struct Source {
origin: moq_net::origin::Consumer,
path: moq_net::PathOwned,
}
impl Source {
pub fn new(origin: moq_net::origin::Consumer, path: impl AsPath) -> Self {
Self {
origin,
path: path.as_path().to_owned(),
}
}
pub async fn broadcast(&self) -> crate::Result<moq_net::broadcast::Consumer> {
Ok(self.origin.request_broadcast(&self.path).await?)
}
pub(crate) fn request(
&self,
rel: Option<&moq_net::PathRelative<'_>>,
) -> Option<kio::Pending<moq_net::origin::Requesting>> {
let target = self.resolve_reference(rel)?;
Some(self.origin.request_broadcast(&target))
}
pub fn resolve_reference(&self, rel: Option<&moq_net::PathRelative<'_>>) -> Option<moq_net::PathOwned> {
match rel.filter(|rel| !rel.is_empty()) {
Some(rel) => self.path.try_resolve(rel),
None => Some(self.path.clone()),
}
}
pub(crate) fn retain_valid<E: crate::catalog::hang::CatalogExt>(
&self,
catalog: &mut crate::catalog::hang::Catalog<E>,
) {
self.retain_valid_references("video", &mut catalog.video.renditions);
self.retain_valid_references("audio", &mut catalog.audio.renditions);
}
pub(crate) fn retain_valid_media(&self, catalog: &mut hang::Catalog) {
self.retain_valid_references("video", &mut catalog.video.renditions);
self.retain_valid_references("audio", &mut catalog.audio.renditions);
}
fn retain_valid_references<C: BroadcastConfig>(
&self,
kind: &'static str,
renditions: &mut std::collections::BTreeMap<String, C>,
) {
renditions.retain(|name, config| {
let valid = self.resolve_reference(config.broadcast()).is_some();
if !valid {
tracing::warn!(
rendition = name,
kind,
"ignoring rendition whose broadcast escapes above the root"
);
}
valid
});
}
pub async fn resolve(
&self,
rel: Option<&moq_net::PathRelative<'_>>,
) -> crate::Result<moq_net::broadcast::Consumer> {
let request = self.request(rel).ok_or_else(|| invalid_broadcast_reference(rel))?;
Ok(request.await?)
}
pub async fn subscribe_track(
&self,
rel: Option<&moq_net::PathRelative<'_>>,
name: &str,
) -> crate::Result<moq_net::track::Subscriber> {
let request = self.request(rel).ok_or_else(|| invalid_broadcast_reference(rel))?;
let broadcast = request.await?;
Ok(broadcast.track(name)?.subscribe(None).await?)
}
}
trait BroadcastConfig {
fn broadcast(&self) -> Option<&moq_net::PathRelativeOwned>;
}
impl BroadcastConfig for hang::catalog::VideoConfig {
fn broadcast(&self) -> Option<&moq_net::PathRelativeOwned> {
self.broadcast.as_ref()
}
}
impl BroadcastConfig for hang::catalog::AudioConfig {
fn broadcast(&self) -> Option<&moq_net::PathRelativeOwned> {
self.broadcast.as_ref()
}
}
fn invalid_broadcast_reference(rel: Option<&moq_net::PathRelative<'_>>) -> crate::Error {
crate::Error::InvalidBroadcastReference(rel.map_or_else(String::new, |rel| rel.as_str().to_string()))
}
#[cfg(test)]
pub(crate) fn announced(broadcast: &moq_net::broadcast::Consumer) -> Source {
let origin = moq_net::Origin::random().produce();
let mut dynamic = origin.dynamic();
let served = broadcast.clone();
tokio::spawn(async move {
while let Ok(request) = dynamic.requested_broadcast().await {
request.accept(served.clone());
}
});
let source = Source::new(origin.consume(), "test");
Box::leak(Box::new(origin));
source
}
#[cfg(test)]
mod tests {
use super::*;
use hang::catalog::{H264, VideoConfig};
use moq_net::{Origin, PathRelative};
async fn settle() {
for _ in 0..10 {
tokio::task::yield_now().await;
}
}
#[tokio::test]
async fn no_override_targets_catalog_broadcast() {
let origin = Origin::random().produce();
let _producer = origin
.create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/pub");
source
.request(None)
.expect("catalog reference should be valid")
.await
.expect("catalog broadcast should resolve");
let empty = PathRelative::empty();
source
.request(Some(&empty))
.expect("empty reference should be valid")
.await
.expect("empty reference should resolve to the catalog broadcast");
}
#[tokio::test]
async fn subscribe_track_resolves_catalog_broadcast() {
let origin = Origin::random().produce();
let mut producer = origin
.create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = producer.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/pub");
source
.subscribe_track(None, "video")
.await
.expect("catalog track should resolve");
}
#[tokio::test]
async fn self_reference_targets_catalog_broadcast() {
let origin = Origin::random().produce();
let mut producer = origin
.create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = producer.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/pub");
let rel = PathRelative::new("./pub");
source
.subscribe_track(Some(&rel), "video")
.await
.expect("self-reference should resolve to the catalog broadcast");
}
#[tokio::test]
async fn escaping_reference_is_rejected_instead_of_using_the_catalog() {
let origin = Origin::random().produce();
let mut producer = origin
.create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = producer.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/pub");
let rel = PathRelative::new("../../source");
assert!(source.resolve_reference(Some(&rel)).is_none());
assert!(matches!(
source.subscribe_track(Some(&rel), "video").await,
Err(crate::Error::InvalidBroadcastReference(reference)) if reference == "../../source"
));
}
#[test]
fn escaping_rendition_is_removed_while_valid_sibling_remains() {
let origin = Origin::random().produce();
let source = Source::new(origin.consume(), "a/pub");
let mut escaped = VideoConfig::new(H264 {
profile: 0x42,
constraints: 0,
level: 0x1e,
inline: false,
});
escaped.broadcast = Some(PathRelative::new("../../source").to_owned());
let mut sibling = escaped.clone();
sibling.broadcast = Some(PathRelative::new("./source").to_owned());
let mut catalog = hang::Catalog::default();
catalog.video.renditions.insert("escaped".to_string(), escaped);
catalog.video.renditions.insert("sibling".to_string(), sibling);
source.retain_valid_media(&mut catalog);
assert!(!catalog.video.renditions.contains_key("escaped"));
assert!(catalog.video.renditions.contains_key("sibling"));
}
#[tokio::test]
async fn subscribe_track_resolves_referenced_broadcast() {
let origin = Origin::random().produce();
let _catalog = origin
.create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let mut referenced = origin
.create_broadcast("a/source", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = referenced.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/pub");
let rel = PathRelative::new("./source");
source
.subscribe_track(Some(&rel), "video")
.await
.expect("referenced track should resolve");
}
#[tokio::test]
async fn dot_resolves_output_parent() {
let origin = Origin::random().produce();
let _catalog = origin
.create_broadcast(
"a/source/transcode",
moq_net::broadcast::Route::new().with_announce(true),
)
.unwrap();
let mut referenced = origin
.create_broadcast("a/source", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = referenced.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "a/source/transcode");
let rel = PathRelative::new(".");
source
.subscribe_track(Some(&rel), "video")
.await
.expect("dot should resolve to the catalog broadcast's parent");
}
#[tokio::test]
async fn dot_resolves_one_segment_catalog_to_root() {
let origin = Origin::random().produce();
let _catalog = origin
.create_broadcast("top", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let mut root = origin
.create_broadcast("", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
let _video = root.create_track("video", None).unwrap();
settle().await;
let source = Source::new(origin.consume(), "top");
let rel = PathRelative::new(".");
source
.subscribe_track(Some(&rel), "video")
.await
.expect("dot should resolve to the empty root broadcast");
}
}