1use moq_net::AsPath;
13
14#[derive(Clone)]
22pub struct Source {
23 origin: moq_net::origin::Consumer,
24 path: moq_net::PathOwned,
25}
26
27impl Source {
28 pub fn new(origin: moq_net::origin::Consumer, path: impl AsPath) -> Self {
36 Self {
37 origin,
38 path: path.as_path().to_owned(),
39 }
40 }
41
42 pub async fn broadcast(&self) -> crate::Result<moq_net::broadcast::Consumer> {
44 Ok(self.origin.request_broadcast(&self.path).await?)
45 }
46
47 pub(crate) fn request(&self, rel: Option<&moq_net::PathRelative<'_>>) -> kio::Pending<moq_net::origin::Requesting> {
56 let target = match rel.filter(|rel| !rel.is_empty()) {
57 Some(rel) => match self.path.resolve(rel) {
60 resolved if resolved.is_empty() => self.path.clone(),
61 resolved => resolved,
62 },
63 None => self.path.clone(),
64 };
65
66 self.origin.request_broadcast(&target)
67 }
68
69 pub async fn resolve(
76 &self,
77 rel: Option<&moq_net::PathRelative<'_>>,
78 ) -> crate::Result<moq_net::broadcast::Consumer> {
79 Ok(self.request(rel).await?)
80 }
81
82 pub async fn subscribe_track(
93 &self,
94 rel: Option<&moq_net::PathRelative<'_>>,
95 name: &str,
96 ) -> crate::Result<moq_net::track::Subscriber> {
97 let broadcast = self.request(rel).await?;
98 Ok(broadcast.track(name)?.subscribe(None).await?)
99 }
100}
101
102#[cfg(test)]
107pub(crate) fn announced(broadcast: &moq_net::broadcast::Consumer) -> Source {
108 let origin = moq_net::Origin::random().produce();
109 let mut dynamic = origin.dynamic();
110 let served = broadcast.clone();
111 tokio::spawn(async move {
112 while let Ok(request) = dynamic.requested_broadcast().await {
113 request.accept(served.clone());
114 }
115 });
116 let source = Source::new(origin.consume(), "test");
117 Box::leak(Box::new(origin));
118 source
119}
120
121#[cfg(test)]
122mod tests {
123 use super::*;
124 use moq_net::{Origin, PathRelative};
125
126 async fn settle() {
129 for _ in 0..10 {
130 tokio::task::yield_now().await;
131 }
132 }
133
134 #[tokio::test]
135 async fn no_override_targets_catalog_broadcast() {
136 let origin = Origin::random().produce();
137 let _producer = origin
138 .create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
139 .unwrap();
140 settle().await;
141
142 let source = Source::new(origin.consume(), "a/pub");
143
144 source.request(None).await.expect("catalog broadcast should resolve");
146 let empty = PathRelative::empty();
147 source
148 .request(Some(&empty))
149 .await
150 .expect("empty reference should resolve to the catalog broadcast");
151 }
152
153 #[tokio::test]
154 async fn subscribe_track_resolves_catalog_broadcast() {
155 let origin = Origin::random().produce();
156 let mut producer = origin
157 .create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
158 .unwrap();
159 let _video = producer.create_track("video", None).unwrap();
161 settle().await;
162
163 let source = Source::new(origin.consume(), "a/pub");
164 source
165 .subscribe_track(None, "video")
166 .await
167 .expect("catalog track should resolve");
168 }
169
170 #[tokio::test]
171 async fn self_reference_targets_catalog_broadcast() {
172 let origin = Origin::random().produce();
173 let mut producer = origin
174 .create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
175 .unwrap();
176 let _video = producer.create_track("video", None).unwrap();
177 settle().await;
178
179 let source = Source::new(origin.consume(), "a/pub");
180
181 let rel = PathRelative::new("../pub");
183 source
184 .subscribe_track(Some(&rel), "video")
185 .await
186 .expect("self-reference should resolve to the catalog broadcast");
187
188 let rel = PathRelative::new("../../..");
190 source
191 .subscribe_track(Some(&rel), "video")
192 .await
193 .expect("excess `..` should resolve to the catalog broadcast");
194 }
195
196 #[tokio::test]
197 async fn subscribe_track_resolves_referenced_broadcast() {
198 let origin = Origin::random().produce();
199
200 let _catalog = origin
201 .create_broadcast("a/pub", moq_net::broadcast::Route::new().with_announce(true))
202 .unwrap();
203
204 let mut referenced = origin
205 .create_broadcast("a/source", moq_net::broadcast::Route::new().with_announce(true))
206 .unwrap();
207 let _video = referenced.create_track("video", None).unwrap();
208 settle().await;
209
210 let source = Source::new(origin.consume(), "a/pub");
211
212 let rel = PathRelative::new("../source");
214 source
215 .subscribe_track(Some(&rel), "video")
216 .await
217 .expect("referenced track should resolve");
218 }
219}