1use crate::{cache, stats, track};
10use std::{
11 collections::{HashMap, VecDeque},
12 sync::Arc,
13 task::{Poll, ready},
14};
15
16use crate::Error;
17use crate::origin::Route;
18
19use super::origin_impl::Announcer;
20use super::{Requests, WeakCache};
21
22#[derive(Clone, Debug)]
27#[non_exhaustive]
28pub struct Info {
29 pub pool: cache::Pool,
31
32 pub cache_duration: std::time::Duration,
34
35 pub path: crate::PathOwned,
49}
50
51impl Default for Info {
52 fn default() -> Self {
53 Self {
54 pool: cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY)),
55 cache_duration: std::time::Duration::MAX,
56 path: crate::PathOwned::default(),
57 }
58 }
59}
60
61impl Info {
62 pub fn new() -> Self {
64 Self::default()
65 }
66
67 pub fn produce(self) -> Producer {
72 Producer::new(self)
73 }
74}
75
76#[derive(Default)]
77struct BroadcastState {
78 tracks: WeakCache<Arc<str>, track::TrackWeak>,
83
84 unique: u64,
86
87 requests: Requests<Arc<str>, track::Request>,
92
93 spliced: Option<SplicedState>,
96
97 closing: bool,
100
101 finished: bool,
104
105 abort: Option<Error>,
108}
109
110#[derive(Default)]
113struct SplicedState {
114 tracks: HashMap<Arc<str>, super::resume::Producer>,
117
118 pending: VecDeque<Arc<str>>,
120}
121
122impl BroadcastState {
123 fn insert_track(&mut self, weak: track::TrackWeak) -> Result<(), Error> {
126 match self.tracks.insert(weak.name().clone(), weak) {
127 Some(_) => Err(Error::Duplicate),
128 None => Ok(()),
129 }
130 }
131
132 fn reject_unserved(&mut self, err: Error) {
140 for request in self.requests.drain_queued() {
141 request.reject(err.clone());
142 }
143 for track in self.tracks.iter() {
144 track.reject(err.clone());
145 }
146 }
147
148 fn is_used(&self) -> bool {
151 if let Some(spliced) = &self.spliced {
152 return spliced.tracks.values().any(|track| track.is_used());
153 }
154 !self.requests.is_empty() || self.tracks.iter().any(|track| track.is_used())
155 }
156
157 fn register_demand(&self, waiter: &kio::Waiter, want: bool) {
162 if let Some(spliced) = &self.spliced {
163 for track in spliced.tracks.values() {
164 let _ = match want {
165 true => track.poll_used(waiter),
166 false => track.poll_unused(waiter),
167 };
168 }
169 return;
170 }
171 for track in self.tracks.iter() {
172 match want {
173 true => track.poll_used(waiter),
174 false => track.poll_unused(waiter),
175 }
176 }
177 }
178}
179
180#[derive(Clone)]
200pub struct Producer {
201 info: Arc<Info>,
204
205 alive: Arc<Alive>,
208
209 state: kio::Shared<BroadcastState>,
212
213 stats: stats::Scope,
217}
218
219impl Producer {
220 pub fn new(info: Info) -> Self {
222 let state = kio::Shared::<BroadcastState>::default();
223 Self {
224 info: Arc::new(info),
225 alive: Alive::new(state.clone()),
226 state,
227 stats: stats::Scope::default(),
228 }
229 }
230
231 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
234 self.stats = scope;
235 self
236 }
237
238 pub(crate) fn with_announcer(self, announcer: Announcer) -> Self {
241 *self.alive.announcer.lock() = Some(announcer);
242 self
243 }
244
245 pub fn announce(&self, route: Route) -> Result<(), Error> {
260 let mut announcer = self.alive.announcer.lock();
261 let announcer = announcer.as_mut().ok_or(Error::Closed)?;
262 announcer.announce(route)
263 }
264
265 pub fn unannounce(&self) {
271 self.alive.unannounce();
272 }
273
274 pub(crate) fn new_spliced(info: Info) -> Self {
278 let state = kio::Shared::new(BroadcastState {
279 spliced: Some(SplicedState::default()),
280 ..Default::default()
281 });
282 Self {
283 info: Arc::new(info),
284 alive: Alive::new(state.clone()),
285 state,
286 stats: stats::Scope::default(),
289 }
290 }
291
292 pub fn info(&self) -> &Info {
294 &self.info
295 }
296
297 pub fn demand(&self) -> Demand {
299 Demand {
300 alive: self.alive.token.consume().weak(),
301 state: self.state.clone(),
302 }
303 }
304
305 pub fn create_track(
310 &self,
311 name: impl Into<Arc<str>>,
312 info: impl Into<Option<track::Info>>,
313 ) -> Result<track::Producer, Error> {
314 let name = name.into();
315 let info = info.into().unwrap_or_default();
316 let mut state = self.state.lock();
317
318 if let Some(request) = state.requests.take(name.as_ref()) {
324 let track = request.with_stats(self.stats.clone()).accept(info);
325 let _ = state.tracks.insert(name, track.weak());
329 return Ok(track);
330 }
331
332 let track = track::Producer::new(self.info.clone(), name, info).with_stats(self.stats.clone());
333 state.insert_track(track.weak())?;
334 Ok(track)
335 }
336
337 pub fn reserve_track(&self, name: impl Into<Arc<str>>) -> Result<track::Request, Error> {
349 let request = track::Request::new(self.info.clone(), name).with_stats(self.stats.clone());
350 self.state.lock().insert_track(request.weak())?;
351 Ok(request)
352 }
353
354 pub fn unique_track(&self, suffix: &str, info: impl Into<Option<track::Info>>) -> Result<track::Producer, Error> {
358 let name = self.unique_name(suffix);
359 self.create_track(name, info)
360 }
361
362 pub fn unique_name(&self, suffix: &str) -> String {
374 let mut state = self.state.lock();
375 let separator = if suffix.starts_with(|c: char| c.is_ascii_digit()) {
376 "-"
377 } else {
378 ""
379 };
380 loop {
381 let id = state.unique;
382 state.unique = id.checked_add(1).expect("unique track IDs exhausted");
383 let name = format!("{id}{separator}{suffix}");
384 if !state.tracks.contains_key(name.as_str()) {
385 return name;
386 }
387 }
388 }
389
390 pub fn dynamic(&self) -> Dynamic {
392 Dynamic::new(
393 self.info.clone(),
394 self.alive.clone(),
395 self.state.clone(),
396 self.stats.clone(),
397 )
398 }
399
400 pub(crate) fn poll_spliced_assigned(&self, waiter: &kio::Waiter) -> Poll<(Arc<str>, super::resume::Producer)> {
403 let mut state = ready!(self.state.poll(waiter, |state| {
404 match &state.spliced {
405 Some(spliced) if !spliced.pending.is_empty() => Poll::Ready(()),
406 _ => Poll::Pending,
407 }
408 }));
409
410 let spliced = state.spliced.as_mut().expect("predicate guaranteed spliced");
411 let name = spliced.pending.pop_front().expect("predicate guaranteed a request");
412 let producer = spliced.tracks.get(&name).expect("pending name without a track").clone();
413 Poll::Ready((name, producer))
414 }
415
416 pub(crate) fn release_spliced(&self, err: Error) {
420 let mut state = self.state.lock();
421 if let Some(spliced) = state.spliced.as_mut() {
422 for name in std::mem::take(&mut spliced.pending) {
423 if let Some(producer) = spliced.tracks.get_mut(&name) {
424 let _ = producer.abort(err.clone());
425 }
426 }
427 spliced.tracks.clear();
428 }
429 }
430
431 pub fn consume(&self) -> Consumer {
433 Consumer {
434 info: self.info.clone(),
435 alive: self.alive.token.consume(),
436 state: self.state.clone(),
437 stats: stats::Scope::default(),
438 }
439 }
440
441 pub fn finish(&self) {
459 {
460 let mut state = self.state.lock();
461 state.closing = true;
462 state.finished = true;
463 state.reject_unserved(Error::NotFound);
467 }
468 let _ = self.alive.token.close();
471 self.alive.retire();
472 }
473
474 pub fn abort(self, err: Error) -> Result<(), Error> {
487 {
488 let mut state = self.state.lock();
489 if state.closing {
490 return Err(Error::Closed);
491 }
492 state.closing = true;
493 state.abort = Some(err.clone());
494 state.reject_unserved(err);
497 }
498 let _ = self.alive.token.close();
499 self.alive.retire();
500 Ok(())
501 }
502
503 pub fn is_clone(&self, other: &Self) -> bool {
505 self.state.same_channel(&other.state)
506 }
507}
508
509struct Alive {
515 token: kio::Producer<()>,
516 state: kio::Shared<BroadcastState>,
517 announcer: kio::Lock<Option<Announcer>>,
521}
522
523impl Alive {
524 fn new(state: kio::Shared<BroadcastState>) -> Arc<Self> {
525 Arc::new(Self {
526 token: kio::Producer::default(),
527 state,
528 announcer: kio::Lock::new(None),
529 })
530 }
531
532 fn unannounce(&self) {
534 if let Some(announcer) = self.announcer.lock().as_mut() {
535 announcer.withdraw();
536 }
537 }
538
539 fn retire(&self) {
542 let announcer = self.announcer.lock().take();
543 drop(announcer);
546 }
547}
548
549impl Drop for Alive {
550 fn drop(&mut self) {
551 if !self.state.read().closing {
555 tracing::warn!(
556 "broadcast::Producer dropped without finish(). Keep the producer alive while publishing, then call finish()."
557 );
558 }
559 self.retire();
560 }
561}
562
563#[cfg(test)]
564#[allow(missing_docs)] impl Producer {
566 pub fn assert_create_track(
567 &mut self,
568 name: impl Into<Arc<str>>,
569 info: impl Into<Option<track::Info>>,
570 ) -> track::Producer {
571 self.create_track(name, info).expect("should not have errored")
572 }
573}
574
575pub(crate) struct SourceGuard {
581 producer: Option<Producer>,
583}
584
585impl SourceGuard {
586 pub fn new(producer: Producer) -> Self {
587 Self {
588 producer: Some(producer),
589 }
590 }
591
592 pub fn finish(mut self) {
595 if let Some(producer) = self.producer.take() {
596 producer.finish();
597 }
598 }
599}
600
601impl Drop for SourceGuard {
602 fn drop(&mut self) {
603 if let Some(producer) = self.producer.take() {
604 let _ = producer.abort(Error::Dropped);
605 }
606 }
607}
608
609pub struct Dynamic {
617 info: Arc<Info>,
618 alive: Arc<Alive>,
620 state: kio::Shared<BroadcastState>,
621 stats: stats::Scope,
624}
625
626impl Clone for Dynamic {
627 fn clone(&self) -> Self {
628 self.state.lock().requests.add_handler();
632
633 Self {
634 info: self.info.clone(),
635 alive: self.alive.clone(),
636 state: self.state.clone(),
637 stats: self.stats.clone(),
638 }
639 }
640}
641
642impl Dynamic {
643 fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
644 state.lock().requests.add_handler();
645
646 Self {
647 info,
648 alive,
649 state,
650 stats,
651 }
652 }
653
654 pub fn info(&self) -> &Info {
656 &self.info
657 }
658
659 pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
665 let mut state = ready!(self.state.poll(waiter, |state| {
666 if state.requests.has_queued() || state.closing {
667 Poll::Ready(())
668 } else {
669 Poll::Pending
670 }
671 }));
672
673 if state.closing && !state.requests.has_queued() {
674 return Poll::Ready(Err(Error::Closed));
675 }
676
677 let name = state.requests.pop().expect("predicate guaranteed a request");
678 let pending = state.requests.remove(&name).expect("popped key must be pending");
679 let _ = state.tracks.insert(name, pending.weak());
682 Poll::Ready(Ok(pending.with_stats(self.stats.clone())))
684 }
685
686 pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
688 kio::wait(|waiter| self.poll_requested_track(waiter)).await
689 }
690
691 pub fn consume(&self) -> Consumer {
693 Consumer {
694 info: self.info.clone(),
695 alive: self.alive.token.consume(),
696 state: self.state.clone(),
697 stats: stats::Scope::default(),
698 }
699 }
700
701 pub async fn closed(&self) -> Error {
704 kio::wait(|waiter| self.poll_closed(waiter)).await
705 }
706
707 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
711 ready!(self.alive.token.poll_closed(waiter));
712 Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
713 }
714
715 pub fn is_clone(&self, other: &Self) -> bool {
717 self.state.same_channel(&other.state)
718 }
719}
720
721impl Drop for Dynamic {
722 fn drop(&mut self) {
723 let mut state = self.state.lock();
726 if state.requests.remove_handler() {
727 for request in state.requests.drain_queued() {
730 request.reject(Error::Dropped);
731 }
732 }
733 }
734}
735
736#[cfg(test)]
737use futures::FutureExt;
738
739#[cfg(test)]
740#[allow(missing_docs)] impl Dynamic {
742 pub fn assert_request(&mut self) -> track::Request {
743 self.requested_track()
744 .now_or_never()
745 .expect("should not have blocked")
746 .expect("should not have errored")
747 }
748
749 pub fn assert_no_request(&mut self) {
750 assert!(self.requested_track().now_or_never().is_none(), "should have blocked");
751 }
752}
753
754pub struct Consumer {
756 info: Arc<Info>,
757 alive: kio::Consumer<()>,
759 state: kio::Shared<BroadcastState>,
761 stats: stats::Scope,
765}
766
767impl Clone for Consumer {
768 fn clone(&self) -> Self {
769 Self {
770 info: self.info.clone(),
771 alive: self.alive.clone(),
772 state: self.state.clone(),
773 stats: self.stats.clone(),
774 }
775 }
776}
777
778impl Consumer {
779 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
782 self.stats = scope;
783 self
784 }
785
786 pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
794 if self.info.path != path {
795 let mut info = (*self.info).clone();
796 info.path = path;
797 self.info = Arc::new(info);
798 }
799 self
800 }
801
802 pub fn info(&self) -> &Info {
804 &self.info
805 }
806
807 pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
809 self.track_inner(name)
814 .map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
815 }
816
817 fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
818 if self.is_closed() {
820 return Err(Error::Dropped);
821 }
822
823 let mut state = self.state.lock();
824
825 let closing = state.closing;
828 if let Some(spliced) = state.spliced.as_mut() {
829 if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
843 spliced.tracks.remove(name);
844 }
845 if let Some(producer) = spliced.tracks.get(name) {
846 return Ok(track::Consumer::spliced(
847 name.into(),
848 self.info.clone(),
849 producer.consume(),
850 ));
851 }
852 if closing {
855 return Err(Error::NotFound);
856 }
857 let name: Arc<str> = name.into();
858 let producer = super::resume::Producer::new();
859 let consumer = producer.consume();
860 spliced.tracks.insert(name.clone(), producer);
861 spliced.pending.push_back(name.clone());
862 return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
863 }
864
865 if let Some(weak) = state.tracks.get(name) {
868 match weak.try_consume() {
869 Some(consumer) => return Ok(consumer),
870 None => {
874 state.tracks.remove(name);
875 }
876 }
877 }
878
879 if let Some(pending) = state.requests.join(name) {
880 return Ok(pending.consume());
882 }
883
884 if state.closing {
887 return Err(Error::NotFound);
888 }
889
890 let name: Arc<str> = name.into();
894 let request = track::Request::new(self.info.clone(), name.clone());
895 let consumer = request.consume();
896
897 if state.requests.insert(name, request).is_err() {
900 return Err(Error::NotFound);
901 }
902
903 Ok(consumer)
904 }
905
906 pub fn demand(&self) -> Demand {
915 Demand {
916 alive: self.alive.weak(),
917 state: self.state.clone(),
918 }
919 }
920
921 pub async fn closed(&self) -> Error {
928 self.alive.closed().await;
929 self.state.read().abort.clone().unwrap_or(Error::Dropped)
930 }
931
932 pub fn is_closed(&self) -> bool {
934 self.alive.is_closed()
935 }
936
937 pub(crate) fn is_closing(&self) -> bool {
942 self.is_closed() || self.state.read().closing
943 }
944
945 pub fn is_finished(&self) -> bool {
948 self.state.read().finished
949 }
950
951 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
957 self.alive.poll_closed(waiter)
958 }
959
960 pub fn is_clone(&self, other: &Self) -> bool {
962 self.state.same_channel(&other.state)
963 }
964
965 pub(crate) fn weak(&self) -> WeakConsumer {
970 WeakConsumer {
971 info: self.info.clone(),
972 alive: self.alive.weak(),
973 state: self.state.clone(),
974 }
975 }
976}
977
978#[derive(Clone)]
985pub(crate) struct WeakConsumer {
986 info: Arc<Info>,
987 alive: kio::ConsumerWeak<()>,
988 state: kio::Shared<BroadcastState>,
989}
990
991impl WeakConsumer {
992 pub fn consume(&self) -> Consumer {
994 Consumer {
995 info: self.info.clone(),
996 alive: self.alive.consume(),
997 state: self.state.clone(),
998 stats: stats::Scope::default(),
999 }
1000 }
1001}
1002
1003impl super::WeakEntry for WeakConsumer {
1004 fn is_closed(&self) -> bool {
1005 self.alive.is_closed()
1006 }
1007
1008 fn same_channel(&self, other: &Self) -> bool {
1009 self.state.same_channel(&other.state)
1010 }
1011}
1012
1013#[derive(Clone)]
1026pub struct Demand {
1027 alive: kio::ConsumerWeak<()>,
1028 state: kio::Shared<BroadcastState>,
1029}
1030
1031impl Demand {
1032 pub fn is_used(&self) -> bool {
1037 self.state.read().is_used()
1038 }
1039
1040 pub async fn used(&self) -> Result<(), Error> {
1043 kio::wait(|waiter| self.poll_used(waiter)).await
1044 }
1045
1046 pub async fn unused(&self) -> Result<(), Error> {
1049 kio::wait(|waiter| self.poll_unused(waiter)).await
1050 }
1051
1052 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1054 self.poll_demand(waiter, true)
1055 }
1056
1057 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1059 self.poll_demand(waiter, false)
1060 }
1061
1062 fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
1063 if self.alive.poll_closed(waiter).is_ready() {
1066 return Poll::Ready(Err(Error::Dropped));
1067 }
1068 let ready = self.state.poll(waiter, |state| {
1069 state.register_demand(waiter, want);
1073 match state.is_used() == want {
1074 true => Poll::Ready(()),
1075 false => Poll::Pending,
1076 }
1077 });
1078 match ready {
1079 Poll::Ready(_) => Poll::Ready(Ok(())),
1080 Poll::Pending => Poll::Pending,
1081 }
1082 }
1083}
1084
1085#[cfg(test)]
1086#[allow(missing_docs)] impl Consumer {
1088 pub fn assert_not_closed(&self) {
1089 assert!(self.closed().now_or_never().is_none(), "should not be closed");
1090 }
1091
1092 pub fn assert_closed(&self) {
1093 assert!(self.closed().now_or_never().is_some(), "should be closed");
1094 }
1095}
1096
1097#[cfg(test)]
1098mod test {
1099 use super::*;
1100 use std::time::Duration;
1101
1102 #[test]
1103 fn unique_names_are_never_reused() {
1104 let producer = Info::new().produce();
1105 let name = producer.unique_name(".opus");
1106 assert_eq!(name, "0.opus");
1107 let track = producer.create_track(name.clone(), None).unwrap();
1108 assert_eq!(producer.unique_name(".opus"), "1.opus");
1109 drop(track);
1110 }
1111
1112 #[test]
1113 fn unique_names_survive_closed_track_pruning() {
1114 let producer = Info::new().produce();
1115 let consumer = producer.consume();
1116 let track = producer.unique_track(".opus", None).unwrap();
1117 assert_eq!(track.name(), "0.opus");
1118 drop(track);
1119 assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
1120 assert_eq!(producer.unique_name(".opus"), "1.opus");
1121 }
1122
1123 #[test]
1124 fn unique_name_skips_a_live_collision() {
1125 let producer = Info::new().produce();
1126 let track = producer.create_track("0.opus", None).unwrap();
1127 assert_eq!(producer.unique_name(".opus"), "1.opus");
1128 drop(track);
1129 assert_eq!(producer.unique_name(".opus"), "2.opus");
1130 }
1131
1132 #[test]
1133 fn unique_names_share_a_counter() {
1134 let producer = Info::new().produce();
1135 assert_eq!(producer.unique_name("-video"), "0-video");
1136 assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
1137 assert_eq!(producer.unique_name("-video"), "2-video");
1138 }
1139
1140 #[test]
1141 fn unique_names_separate_numeric_suffixes() {
1142 let producer = Info::new().produce();
1143 assert_eq!(producer.unique_name(""), "0");
1144 let name = producer.unique_name("2");
1145 assert_eq!(name, "1-2");
1146 for _ in 2..12 {
1147 producer.unique_name("");
1148 }
1149 assert_eq!(producer.unique_name(""), "12");
1150 }
1151
1152 async fn expect<T>(fut: impl Future<Output = T>) -> T {
1155 tokio::time::timeout(Duration::from_secs(1), fut)
1156 .await
1157 .expect("timed out waiting for a demand edge")
1158 }
1159
1160 #[tokio::test]
1164 async fn demand_ordinary() {
1165 tokio::time::pause();
1166
1167 let producer = Info::new().produce();
1168 let consumer = producer.consume();
1169 let demand = producer.demand();
1170
1171 assert!(!demand.is_used());
1173 demand.unused().await.unwrap();
1174
1175 let _track = producer.create_track("a", None).unwrap();
1177 assert!(!demand.is_used());
1178
1179 let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
1181 used.unwrap();
1182 assert!(demand.is_used());
1183
1184 let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
1186 unused.unwrap();
1187 assert!(!demand.is_used());
1188
1189 producer.finish();
1191 assert!(matches!(demand.used().await, Err(Error::Dropped)));
1192 assert!(matches!(demand.unused().await, Err(Error::Dropped)));
1193 }
1194
1195 #[tokio::test]
1198 async fn demand_spliced() {
1199 tokio::time::pause();
1200
1201 let producer = Producer::new_spliced(Info::new());
1202 let consumer = producer.consume();
1203 let demand = producer.demand();
1204 let watched = consumer.demand();
1205
1206 assert!(!demand.is_used());
1207 assert!(!watched.is_used());
1208 let track = consumer.track("video").unwrap();
1209 assert!(demand.is_used());
1210 assert!(watched.is_used());
1211
1212 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1215 unused.unwrap();
1216 assert!(!demand.is_used());
1217 assert!(!watched.is_used());
1218
1219 let _track = consumer.track("video").unwrap();
1221 assert!(demand.is_used());
1222 }
1223
1224 #[tokio::test]
1226 async fn consumer_demand_reports_dropped_producer() {
1227 let producer = Producer::new_spliced(Info::new());
1228 let consumer = producer.consume();
1229 let watched = consumer.demand();
1230
1231 let track = consumer.track("video").unwrap();
1232 assert!(watched.is_used());
1233
1234 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1235 unused.unwrap();
1236
1237 drop(producer);
1238 assert!(matches!(watched.used().await, Err(Error::Dropped)));
1239 assert!(matches!(watched.unused().await, Err(Error::Dropped)));
1240 }
1241
1242 macro_rules! subscribe_pending {
1245 ($consumer:expr, $name:expr) => {{
1246 let pending = $consumer.track($name).unwrap().subscribe(None);
1247 assert!(
1248 pending.poll_ok(&kio::Waiter::noop()).is_pending(),
1249 "subscribe should stay pending until the request is accepted"
1250 );
1251 pending
1252 }};
1253 }
1254
1255 #[tokio::test]
1256 async fn insert() {
1257 let mut producer = Info::new().produce();
1258
1259 let track1 = producer.assert_create_track("track1", None);
1261 track1.append_group().unwrap();
1262
1263 let consumer = producer.consume();
1264
1265 let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1267 track1_sub.assert_group();
1268
1269 let track2 = producer.assert_create_track("track2", None);
1270
1271 let consumer2 = producer.consume();
1272 let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
1273 track2_consumer.assert_no_group();
1274
1275 track2.append_group().unwrap();
1276
1277 track2_consumer.assert_group();
1278 }
1279
1280 #[tokio::test]
1281 async fn closed() {
1282 let mut producer = Info::new().produce();
1283 let dynamic = producer.dynamic();
1284
1285 let consumer = producer.consume();
1286 consumer.assert_not_closed();
1287
1288 let track1 = producer.assert_create_track("track1", None);
1290 let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1291
1292 let track2_fut = subscribe_pending!(consumer, "track2");
1294
1295 drop(dynamic);
1298
1299 assert!(track2_fut.await.is_err());
1301
1302 assert!(!track1.is_closed());
1304 track1c.assert_not_closed();
1305 }
1306
1307 #[tokio::test]
1310 async fn closed_cause() {
1311 let producer = Info::new().produce();
1313 let consumer = producer.consume();
1314 producer.abort(Error::Timeout).unwrap();
1315 assert!(matches!(consumer.closed().await, Error::Timeout));
1316 assert!(!consumer.is_finished());
1317
1318 let producer = Info::new().produce();
1320 let consumer = producer.consume();
1321 producer.finish();
1322 assert!(matches!(consumer.closed().await, Error::Dropped));
1323 assert!(consumer.is_finished());
1324
1325 let producer = Info::new().produce();
1327 let consumer = producer.consume();
1328 drop(producer);
1330 assert!(matches!(consumer.closed().await, Error::Dropped));
1331 assert!(!consumer.is_finished());
1332 }
1333
1334 #[tokio::test]
1335 async fn requests() {
1336 let mut producer = Info::new().produce().dynamic();
1337
1338 let consumer = producer.consume();
1339 let consumer2 = consumer.clone();
1340
1341 let track1_fut = subscribe_pending!(consumer, "track1");
1343 let track2_fut = subscribe_pending!(consumer2, "track1");
1344
1345 let request = producer.assert_request();
1347 producer.assert_no_request();
1348 assert_eq!(request.name(), "track1");
1349
1350 let track3 = request.accept(None);
1352 let mut track1 = track1_fut.await.unwrap();
1353 let mut track2 = track2_fut.await.unwrap();
1354
1355 track1.assert_not_closed();
1356 track1.assert_is_clone(&track2);
1357 track3.subscribe(None).assert_is_clone(&track1);
1358
1359 track3.append_group().unwrap();
1361 track1.assert_group();
1362 track2.assert_group();
1363
1364 let track4_fut = subscribe_pending!(consumer, "track2");
1366 drop(producer);
1367 assert!(track4_fut.await.is_err());
1368
1369 let track5 = consumer2.track("track3");
1371 assert!(track5.is_err(), "should have errored");
1372 }
1373
1374 #[tokio::test]
1375 async fn stale_producer() {
1376 let mut broadcast = Info::new().produce().dynamic();
1377 let consumer = broadcast.consume();
1378
1379 let track1_fut = subscribe_pending!(consumer, "track1");
1381 let producer1 = broadcast.assert_request().accept(None);
1382 let mut track1 = track1_fut.await.unwrap();
1383
1384 producer1.append_group().unwrap();
1386 producer1.finish().unwrap();
1387 drop(producer1);
1388
1389 track1.assert_closed();
1391
1392 let track2_fut = subscribe_pending!(consumer, "track1");
1394 let producer2 = broadcast.assert_request().accept(None);
1395 let mut track2 = track2_fut.await.unwrap();
1396 track2.assert_not_closed();
1397 track2.assert_not_clone(&track1);
1398
1399 producer2.append_group().unwrap();
1401 track2.assert_group();
1402 }
1403
1404 #[tokio::test(start_paused = true)]
1405 async fn requested_unused() {
1406 let mut broadcast = Info::new().produce().dynamic();
1407 let bc = broadcast.consume();
1408
1409 let c1_fut = subscribe_pending!(bc, "unknown_track");
1411 let producer1 = broadcast.assert_request().accept(None);
1412 let consumer1 = c1_fut.await.unwrap();
1413
1414 assert!(
1416 producer1.unused().now_or_never().is_none(),
1417 "track producer should be used"
1418 );
1419
1420 let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1422 consumer2.assert_is_clone(&consumer1);
1423
1424 drop(consumer1);
1425 assert!(
1426 producer1.unused().now_or_never().is_none(),
1427 "track producer should be used"
1428 );
1429
1430 drop(consumer2);
1431 assert!(
1432 producer1.unused().now_or_never().is_some(),
1433 "track producer should be unused after all consumers are dropped"
1434 );
1435
1436 let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1440 consumer3.assert_is_clone(&producer1.subscribe(None));
1441 broadcast.assert_no_request();
1442 drop(consumer3);
1443
1444 producer1.abort(Error::Cancel).unwrap();
1447
1448 let c4_fut = subscribe_pending!(bc, "unknown_track");
1449 let producer2 = broadcast.assert_request().accept(None);
1450 let consumer4 = c4_fut.await.unwrap();
1451 drop(consumer4);
1452 assert!(
1453 producer2.unused().now_or_never().is_some(),
1454 "new track producer should be unused after its consumer is dropped"
1455 );
1456 }
1457
1458 #[tokio::test]
1464 async fn create_track_fulfills_queued_request() {
1465 let producer = Info::new().produce();
1466 let mut dynamic = producer.dynamic();
1467 let bc = dynamic.consume();
1468
1469 let subscribing = subscribe_pending!(bc, "video");
1471
1472 let track = producer.create_track("video", None).unwrap();
1474 let mut sub = subscribing.await.expect("fulfilled by create_track");
1475
1476 track.append_group().unwrap();
1478 sub.recv_group().await.expect("recv").expect("group");
1479
1480 dynamic.assert_no_request();
1482 let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
1483 again.assert_is_clone(&track.subscribe(None));
1484 }
1485
1486 #[tokio::test]
1492 async fn dynamic_clone_keeps_alive() {
1493 let broadcast = Info::new().produce().dynamic();
1494 let consumer = broadcast.consume();
1495
1496 let clone = broadcast.clone();
1497 drop(clone);
1498
1499 let _fut = subscribe_pending!(consumer, "track1");
1502 }
1503
1504 #[tokio::test]
1509 async fn finish_resolves_a_reserved_name() {
1510 let producer = Info::new().produce();
1511 let consumer = producer.consume();
1512
1513 let _request = producer.reserve_track("track1").unwrap();
1514 let pending = subscribe_pending!(consumer, "track1");
1515
1516 producer.finish();
1517 assert!(matches!(pending.await, Err(Error::NotFound)));
1518 }
1519
1520 #[tokio::test]
1523 async fn abort_resolves_a_reserved_name_with_its_reason() {
1524 let producer = Info::new().produce();
1525 let consumer = producer.consume();
1526
1527 let request = producer.reserve_track("track1").unwrap();
1528 let pending = subscribe_pending!(consumer, "track1");
1529
1530 producer.abort(Error::Cancel).unwrap();
1531 assert!(matches!(pending.await, Err(Error::Cancel)));
1532
1533 let track = request.accept(None);
1534 let mut subscriber = track.subscribe(None);
1535 assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
1536 }
1537
1538 #[tokio::test]
1541 async fn finish_resolves_a_queued_request() {
1542 let producer = Info::new().produce();
1543 let dynamic = producer.dynamic();
1544 let consumer = dynamic.consume();
1545
1546 let pending = subscribe_pending!(consumer, "track1");
1547
1548 producer.finish();
1549 assert!(matches!(pending.await, Err(Error::NotFound)));
1550 drop(dynamic);
1551 }
1552
1553 #[tokio::test]
1556 async fn finish_resolves_a_request_a_handler_never_answered() {
1557 let producer = Info::new().produce();
1558 let mut dynamic = producer.dynamic();
1559 let consumer = dynamic.consume();
1560
1561 let pending = subscribe_pending!(consumer, "track1");
1562 let _request = dynamic.requested_track().await.unwrap();
1563
1564 producer.finish();
1565 assert!(matches!(pending.await, Err(Error::NotFound)));
1566 drop(dynamic);
1567 }
1568
1569 #[tokio::test]
1574 async fn finish_resolves_an_unaccepted_track_with_fetched_info() {
1575 let producer = Info::new().produce();
1576 let consumer = producer.consume();
1577
1578 let request = producer.reserve_track("track1").unwrap();
1579 let dynamic = request.dynamic();
1580 let track = consumer.track("track1").unwrap();
1581 let pending_fetch = track.fetch_group(0, None);
1582 let fetch = dynamic.requested_group().await.unwrap();
1583 let group = fetch.accept(None).unwrap();
1584 group.finish().unwrap();
1585 pending_fetch.await.unwrap();
1586
1587 let mut subscriber = track.subscribe(None).await.unwrap();
1588 producer.finish();
1589 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1590
1591 let stale = request.accept(None);
1592 assert!(stale.append_group().is_err());
1593 }
1594
1595 #[tokio::test]
1598 async fn finish_spares_a_served_track() {
1599 let producer = Info::new().produce();
1600 let consumer = producer.consume();
1601
1602 let track = producer.create_track("track1", None).unwrap();
1603 let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1604
1605 producer.finish();
1606
1607 track.append_group().unwrap();
1608 subscriber.assert_group();
1609 track.finish().unwrap();
1610 }
1611
1612 #[tokio::test]
1616 async fn finish_leaves_a_stale_reservation_inert() {
1617 let producer = Info::new().produce();
1618 let consumer = producer.consume();
1619
1620 let request = producer.reserve_track("track1").unwrap();
1621 let pending = subscribe_pending!(consumer, "track1");
1622
1623 producer.finish();
1624 assert!(matches!(pending.await, Err(Error::NotFound)));
1625
1626 let track = request.accept(None);
1627 assert!(track.append_group().is_err());
1628 let mut subscriber = track.subscribe(None);
1629 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1630 assert!(consumer.track("track1").is_err());
1631 }
1632
1633 #[tokio::test]
1637 async fn dropping_a_reserved_request_resolves_dropped() {
1638 let producer = Info::new().produce();
1639 let consumer = producer.consume();
1640
1641 let request = producer.reserve_track("track1").unwrap();
1642 let pending = subscribe_pending!(consumer, "track1");
1643
1644 drop(request);
1645 assert!(matches!(pending.await, Err(Error::Dropped)));
1646 producer.finish();
1647 }
1648
1649 #[tokio::test]
1652 async fn rejecting_a_reserved_request_carries_the_reason() {
1653 let producer = Info::new().produce();
1654 let consumer = producer.consume();
1655
1656 let request = producer.reserve_track("track1").unwrap();
1657 let pending = subscribe_pending!(consumer, "track1");
1658
1659 request.reject(Error::NotFound);
1660 assert!(matches!(pending.await, Err(Error::NotFound)));
1661 producer.finish();
1662 }
1663
1664 #[tokio::test]
1670 async fn an_idle_teardown_yields_to_a_returning_viewer() {
1671 let producer = Info::new().produce();
1672 let consumer = producer.consume();
1673 let track = producer.create_track("video", None).unwrap();
1674
1675 assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
1677
1678 let viewer = consumer.track("video").unwrap();
1680 let track = track
1681 .abort_unused(Error::Cancel)
1682 .expect_err("viewer keeps the track alive");
1683
1684 assert!(!track.is_closed());
1686 let mut subscriber = viewer.subscribe(None).await.unwrap();
1687 subscriber.assert_no_group();
1688 track.append_group().unwrap();
1689 assert!(subscriber.recv_group().await.unwrap().is_some());
1690
1691 drop(subscriber);
1694 drop(viewer);
1695 assert!(track.abort_unused(Error::Cancel).is_ok());
1696 assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
1697
1698 producer.finish();
1699 }
1700
1701 #[test]
1702 fn abort_unused_accepts_an_already_closed_track_with_consumers() {
1703 let producer = Info::new().produce();
1704 let consumer = producer.consume();
1705 let track = producer.create_track("video", None).unwrap();
1706 let _viewer = consumer.track("video").unwrap();
1707 assert!(track.is_used());
1708 track.clone().abort(Error::Cancel).unwrap();
1709 assert!(!track.is_used());
1710 assert!(track.abort_unused(Error::Cancel).is_ok());
1711 producer.finish();
1712 }
1713}