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<'_>>) -> kio::Pending<moq_net::origin::Requesting> {
let target = match rel.filter(|rel| !rel.is_empty()) {
Some(rel) => match self.path.resolve(rel) {
resolved if resolved.is_empty() => self.path.clone(),
resolved => resolved,
},
None => self.path.clone(),
};
self.origin.request_broadcast(&target)
}
pub async fn resolve(
&self,
rel: Option<&moq_net::PathRelative<'_>>,
) -> crate::Result<moq_net::broadcast::Consumer> {
Ok(self.request(rel).await?)
}
pub async fn subscribe_track(
&self,
rel: Option<&moq_net::PathRelative<'_>>,
name: &str,
) -> crate::Result<moq_net::track::Subscriber> {
let broadcast = self.request(rel).await?;
Ok(broadcast.track(name)?.subscribe(None).await?)
}
}
#[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 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).await.expect("catalog broadcast should resolve");
let empty = PathRelative::empty();
source
.request(Some(&empty))
.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");
let rel = PathRelative::new("../../..");
source
.subscribe_track(Some(&rel), "video")
.await
.expect("excess `..` should resolve to the catalog broadcast");
}
#[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");
}
}