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) {
460 {
461 let mut state = self.state.lock();
462 state.closing = true;
463 state.finished = true;
464 state.reject_unserved(Error::NotFound);
468 }
469 let _ = self.alive.token.close();
472 self.alive.retire();
473 }
474
475 pub fn abort(self, err: Error) -> Result<(), Error> {
488 {
489 let mut state = self.state.lock();
490 if state.closing {
491 return Err(Error::Closed);
492 }
493 state.closing = true;
494 state.abort = Some(err.clone());
495 state.reject_unserved(err);
498 }
499 let _ = self.alive.token.close();
500 self.alive.retire();
501 Ok(())
502 }
503
504 pub fn is_clone(&self, other: &Self) -> bool {
506 self.state.same_channel(&other.state)
507 }
508}
509
510struct Alive {
516 token: kio::Producer<()>,
517 state: kio::Shared<BroadcastState>,
518 announcer: kio::Lock<Option<Announcer>>,
522}
523
524impl Alive {
525 fn new(state: kio::Shared<BroadcastState>) -> Arc<Self> {
526 Arc::new(Self {
527 token: kio::Producer::default(),
528 state,
529 announcer: kio::Lock::new(None),
530 })
531 }
532
533 fn unannounce(&self) {
535 if let Some(announcer) = self.announcer.lock().as_mut() {
536 announcer.withdraw();
537 }
538 }
539
540 fn retire(&self) {
543 let announcer = self.announcer.lock().take();
544 drop(announcer);
547 }
548}
549
550impl Drop for Alive {
551 fn drop(&mut self) {
552 if !self.state.read().closing {
556 tracing::warn!(
557 "broadcast::Producer dropped without finish(). Keep the producer alive while publishing, then call finish()."
558 );
559 }
560 self.retire();
561 }
562}
563
564#[cfg(test)]
565#[allow(missing_docs)] impl Producer {
567 pub fn assert_create_track(
568 &mut self,
569 name: impl Into<Arc<str>>,
570 info: impl Into<Option<track::Info>>,
571 ) -> track::Producer {
572 self.create_track(name, info).expect("should not have errored")
573 }
574}
575
576pub(crate) struct SourceGuard {
582 producer: Option<Producer>,
584}
585
586impl SourceGuard {
587 pub fn new(producer: Producer) -> Self {
588 Self {
589 producer: Some(producer),
590 }
591 }
592
593 pub fn finish(mut self) {
596 if let Some(producer) = self.producer.take() {
597 producer.finish();
598 }
599 }
600}
601
602impl Drop for SourceGuard {
603 fn drop(&mut self) {
604 if let Some(producer) = self.producer.take() {
605 let _ = producer.abort(Error::Dropped);
606 }
607 }
608}
609
610pub struct Dynamic {
618 info: Arc<Info>,
619 alive: Arc<Alive>,
621 state: kio::Shared<BroadcastState>,
622 stats: stats::Scope,
625}
626
627impl Clone for Dynamic {
628 fn clone(&self) -> Self {
629 self.state.lock().requests.add_handler();
633
634 Self {
635 info: self.info.clone(),
636 alive: self.alive.clone(),
637 state: self.state.clone(),
638 stats: self.stats.clone(),
639 }
640 }
641}
642
643impl Dynamic {
644 fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
645 state.lock().requests.add_handler();
646
647 Self {
648 info,
649 alive,
650 state,
651 stats,
652 }
653 }
654
655 pub fn info(&self) -> &Info {
657 &self.info
658 }
659
660 pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
666 let mut state = ready!(self.state.poll(waiter, |state| {
667 if state.requests.has_queued() || state.closing {
668 Poll::Ready(())
669 } else {
670 Poll::Pending
671 }
672 }));
673
674 if state.closing && !state.requests.has_queued() {
675 return Poll::Ready(Err(Error::Closed));
676 }
677
678 let name = state.requests.pop().expect("predicate guaranteed a request");
679 let pending = state.requests.remove(&name).expect("popped key must be pending");
680 let _ = state.tracks.insert(name, pending.weak());
683 Poll::Ready(Ok(pending.claim().with_stats(self.stats.clone())))
685 }
686
687 pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
689 kio::wait(|waiter| self.poll_requested_track(waiter)).await
690 }
691
692 pub fn consume(&self) -> Consumer {
694 Consumer {
695 info: self.info.clone(),
696 alive: self.alive.token.consume(),
697 state: self.state.clone(),
698 stats: stats::Scope::default(),
699 }
700 }
701
702 pub async fn closed(&self) -> Error {
705 kio::wait(|waiter| self.poll_closed(waiter)).await
706 }
707
708 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
712 ready!(self.alive.token.poll_closed(waiter));
713 Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
714 }
715
716 pub fn is_clone(&self, other: &Self) -> bool {
718 self.state.same_channel(&other.state)
719 }
720}
721
722impl Drop for Dynamic {
723 fn drop(&mut self) {
724 let mut state = self.state.lock();
727 if state.requests.remove_handler() {
728 for request in state.requests.drain_queued() {
731 request.reject(Error::Dropped);
732 }
733 }
734 }
735}
736
737#[cfg(test)]
738use futures::FutureExt;
739
740#[cfg(test)]
741#[allow(missing_docs)] impl Dynamic {
743 pub fn assert_request(&mut self) -> track::Request {
744 self.requested_track()
745 .now_or_never()
746 .expect("should not have blocked")
747 .expect("should not have errored")
748 }
749
750 pub fn assert_no_request(&mut self) {
751 assert!(self.requested_track().now_or_never().is_none(), "should have blocked");
752 }
753}
754
755pub struct Consumer {
757 info: Arc<Info>,
758 alive: kio::Consumer<()>,
760 state: kio::Shared<BroadcastState>,
762 stats: stats::Scope,
766}
767
768impl Clone for Consumer {
769 fn clone(&self) -> Self {
770 Self {
771 info: self.info.clone(),
772 alive: self.alive.clone(),
773 state: self.state.clone(),
774 stats: self.stats.clone(),
775 }
776 }
777}
778
779impl Consumer {
780 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
783 self.stats = scope;
784 self
785 }
786
787 pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
795 if self.info.path != path {
796 let mut info = (*self.info).clone();
797 info.path = path;
798 self.info = Arc::new(info);
799 }
800 self
801 }
802
803 pub fn info(&self) -> &Info {
805 &self.info
806 }
807
808 pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
810 self.track_inner(name)
815 .map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
816 }
817
818 fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
819 if self.is_closed() {
821 return Err(Error::Dropped);
822 }
823
824 let mut state = self.state.lock();
825
826 let closing = state.closing;
829 if let Some(spliced) = state.spliced.as_mut() {
830 if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
844 spliced.tracks.remove(name);
845 }
846 if let Some(producer) = spliced.tracks.get(name) {
847 return Ok(track::Consumer::spliced(
848 name.into(),
849 self.info.clone(),
850 producer.consume(),
851 ));
852 }
853 if closing {
856 return Err(Error::NotFound);
857 }
858 let name: Arc<str> = name.into();
859 let producer = super::resume::Producer::new();
860 let consumer = producer.consume();
861 spliced.tracks.insert(name.clone(), producer);
862 spliced.pending.push_back(name.clone());
863 return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
864 }
865
866 if let Some(weak) = state.tracks.get(name) {
869 match weak.try_consume() {
870 Some(consumer) => return Ok(consumer),
871 None => {
875 state.tracks.remove(name);
876 }
877 }
878 }
879
880 if let Some(pending) = state.requests.join(name) {
881 return Ok(pending.consume());
883 }
884
885 if state.closing {
888 return Err(Error::NotFound);
889 }
890
891 let name: Arc<str> = name.into();
895 let request = track::Request::new(self.info.clone(), name.clone());
896 let consumer = request.consume();
897
898 if state.requests.insert(name, request).is_err() {
901 return Err(Error::NotFound);
902 }
903
904 Ok(consumer)
905 }
906
907 pub fn demand(&self) -> Demand {
916 Demand {
917 alive: self.alive.weak(),
918 state: self.state.clone(),
919 }
920 }
921
922 pub async fn closed(&self) -> Error {
929 self.alive.closed().await;
930 self.state.read().abort.clone().unwrap_or(Error::Dropped)
931 }
932
933 pub fn is_closed(&self) -> bool {
935 self.alive.is_closed()
936 }
937
938 pub(crate) fn is_closing(&self) -> bool {
943 self.is_closed() || self.state.read().closing
944 }
945
946 pub fn is_finished(&self) -> bool {
949 self.state.read().finished
950 }
951
952 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
958 self.alive.poll_closed(waiter)
959 }
960
961 pub fn is_clone(&self, other: &Self) -> bool {
963 self.state.same_channel(&other.state)
964 }
965
966 pub(crate) fn weak(&self) -> WeakConsumer {
971 WeakConsumer {
972 info: self.info.clone(),
973 alive: self.alive.weak(),
974 state: self.state.clone(),
975 }
976 }
977}
978
979#[derive(Clone)]
986pub(crate) struct WeakConsumer {
987 info: Arc<Info>,
988 alive: kio::ConsumerWeak<()>,
989 state: kio::Shared<BroadcastState>,
990}
991
992impl WeakConsumer {
993 pub fn consume(&self) -> Consumer {
995 Consumer {
996 info: self.info.clone(),
997 alive: self.alive.consume(),
998 state: self.state.clone(),
999 stats: stats::Scope::default(),
1000 }
1001 }
1002}
1003
1004impl super::WeakEntry for WeakConsumer {
1005 fn is_closed(&self) -> bool {
1006 self.alive.is_closed()
1007 }
1008
1009 fn same_channel(&self, other: &Self) -> bool {
1010 self.state.same_channel(&other.state)
1011 }
1012}
1013
1014#[derive(Clone)]
1027pub struct Demand {
1028 alive: kio::ConsumerWeak<()>,
1029 state: kio::Shared<BroadcastState>,
1030}
1031
1032impl Demand {
1033 pub fn is_used(&self) -> bool {
1038 self.state.read().is_used()
1039 }
1040
1041 pub async fn used(&self) -> Result<(), Error> {
1044 kio::wait(|waiter| self.poll_used(waiter)).await
1045 }
1046
1047 pub async fn unused(&self) -> Result<(), Error> {
1050 kio::wait(|waiter| self.poll_unused(waiter)).await
1051 }
1052
1053 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1055 self.poll_demand(waiter, true)
1056 }
1057
1058 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1060 self.poll_demand(waiter, false)
1061 }
1062
1063 fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
1064 if self.alive.poll_closed(waiter).is_ready() {
1067 return Poll::Ready(Err(Error::Dropped));
1068 }
1069 let ready = self.state.poll(waiter, |state| {
1070 state.register_demand(waiter, want);
1074 match state.is_used() == want {
1075 true => Poll::Ready(()),
1076 false => Poll::Pending,
1077 }
1078 });
1079 match ready {
1080 Poll::Ready(_) => Poll::Ready(Ok(())),
1081 Poll::Pending => Poll::Pending,
1082 }
1083 }
1084}
1085
1086#[cfg(test)]
1087#[allow(missing_docs)] impl Consumer {
1089 pub fn assert_not_closed(&self) {
1090 assert!(self.closed().now_or_never().is_none(), "should not be closed");
1091 }
1092
1093 pub fn assert_closed(&self) {
1094 assert!(self.closed().now_or_never().is_some(), "should be closed");
1095 }
1096}
1097
1098#[cfg(test)]
1099mod test {
1100 use super::*;
1101 use std::time::Duration;
1102
1103 #[test]
1104 fn unique_names_are_never_reused() {
1105 let producer = Info::new().produce();
1106 let name = producer.unique_name(".opus");
1107 assert_eq!(name, "0.opus");
1108 let track = producer.create_track(name.clone(), None).unwrap();
1109 assert_eq!(producer.unique_name(".opus"), "1.opus");
1110 drop(track);
1111 }
1112
1113 #[test]
1114 fn unique_names_survive_closed_track_pruning() {
1115 let producer = Info::new().produce();
1116 let consumer = producer.consume();
1117 let track = producer.unique_track(".opus", None).unwrap();
1118 assert_eq!(track.name(), "0.opus");
1119 drop(track);
1120 assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
1121 assert_eq!(producer.unique_name(".opus"), "1.opus");
1122 }
1123
1124 #[test]
1125 fn unique_name_skips_a_live_collision() {
1126 let producer = Info::new().produce();
1127 let track = producer.create_track("0.opus", None).unwrap();
1128 assert_eq!(producer.unique_name(".opus"), "1.opus");
1129 drop(track);
1130 assert_eq!(producer.unique_name(".opus"), "2.opus");
1131 }
1132
1133 #[test]
1134 fn unique_names_share_a_counter() {
1135 let producer = Info::new().produce();
1136 assert_eq!(producer.unique_name("-video"), "0-video");
1137 assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
1138 assert_eq!(producer.unique_name("-video"), "2-video");
1139 }
1140
1141 #[test]
1142 fn unique_names_separate_numeric_suffixes() {
1143 let producer = Info::new().produce();
1144 assert_eq!(producer.unique_name(""), "0");
1145 let name = producer.unique_name("2");
1146 assert_eq!(name, "1-2");
1147 for _ in 2..12 {
1148 producer.unique_name("");
1149 }
1150 assert_eq!(producer.unique_name(""), "12");
1151 }
1152
1153 async fn expect<T>(fut: impl Future<Output = T>) -> T {
1156 tokio::time::timeout(Duration::from_secs(1), fut)
1157 .await
1158 .expect("timed out waiting for a demand edge")
1159 }
1160
1161 #[tokio::test]
1165 async fn demand_ordinary() {
1166 tokio::time::pause();
1167
1168 let producer = Info::new().produce();
1169 let consumer = producer.consume();
1170 let demand = producer.demand();
1171
1172 assert!(!demand.is_used());
1174 demand.unused().await.unwrap();
1175
1176 let _track = producer.create_track("a", None).unwrap();
1178 assert!(!demand.is_used());
1179
1180 let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
1182 used.unwrap();
1183 assert!(demand.is_used());
1184
1185 let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
1187 unused.unwrap();
1188 assert!(!demand.is_used());
1189
1190 producer.finish();
1192 assert!(matches!(demand.used().await, Err(Error::Dropped)));
1193 assert!(matches!(demand.unused().await, Err(Error::Dropped)));
1194 }
1195
1196 #[tokio::test]
1199 async fn demand_spliced() {
1200 tokio::time::pause();
1201
1202 let producer = Producer::new_spliced(Info::new());
1203 let consumer = producer.consume();
1204 let demand = producer.demand();
1205 let watched = consumer.demand();
1206
1207 assert!(!demand.is_used());
1208 assert!(!watched.is_used());
1209 let track = consumer.track("video").unwrap();
1210 assert!(demand.is_used());
1211 assert!(watched.is_used());
1212
1213 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1216 unused.unwrap();
1217 assert!(!demand.is_used());
1218 assert!(!watched.is_used());
1219
1220 let _track = consumer.track("video").unwrap();
1222 assert!(demand.is_used());
1223 }
1224
1225 #[tokio::test]
1227 async fn consumer_demand_reports_dropped_producer() {
1228 let producer = Producer::new_spliced(Info::new());
1229 let consumer = producer.consume();
1230 let watched = consumer.demand();
1231
1232 let track = consumer.track("video").unwrap();
1233 assert!(watched.is_used());
1234
1235 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1236 unused.unwrap();
1237
1238 drop(producer);
1239 assert!(matches!(watched.used().await, Err(Error::Dropped)));
1240 assert!(matches!(watched.unused().await, Err(Error::Dropped)));
1241 }
1242
1243 macro_rules! subscribe_pending {
1246 ($consumer:expr, $name:expr) => {{
1247 let pending = $consumer.track($name).unwrap().subscribe(None);
1248 assert!(
1249 pending.poll_ok(&kio::Waiter::noop()).is_pending(),
1250 "subscribe should stay pending until the request is accepted"
1251 );
1252 pending
1253 }};
1254 }
1255
1256 #[tokio::test]
1257 async fn insert() {
1258 let mut producer = Info::new().produce();
1259
1260 let track1 = producer.assert_create_track("track1", None);
1262 track1.append_group().unwrap();
1263
1264 let consumer = producer.consume();
1265
1266 let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1268 track1_sub.assert_group();
1269
1270 let track2 = producer.assert_create_track("track2", None);
1271
1272 let consumer2 = producer.consume();
1273 let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
1274 track2_consumer.assert_no_group();
1275
1276 track2.append_group().unwrap();
1277
1278 track2_consumer.assert_group();
1279 }
1280
1281 #[tokio::test]
1282 async fn closed() {
1283 let mut producer = Info::new().produce();
1284 let dynamic = producer.dynamic();
1285
1286 let consumer = producer.consume();
1287 consumer.assert_not_closed();
1288
1289 let track1 = producer.assert_create_track("track1", None);
1291 let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1292
1293 let track2_fut = subscribe_pending!(consumer, "track2");
1295
1296 drop(dynamic);
1299
1300 assert!(track2_fut.await.is_err());
1302
1303 assert!(!track1.is_closed());
1305 track1c.assert_not_closed();
1306 }
1307
1308 #[tokio::test]
1311 async fn closed_cause() {
1312 let producer = Info::new().produce();
1314 let consumer = producer.consume();
1315 producer.abort(Error::Timeout).unwrap();
1316 assert!(matches!(consumer.closed().await, Error::Timeout));
1317 assert!(!consumer.is_finished());
1318
1319 let producer = Info::new().produce();
1321 let consumer = producer.consume();
1322 producer.finish();
1323 assert!(matches!(consumer.closed().await, Error::Dropped));
1324 assert!(consumer.is_finished());
1325
1326 let producer = Info::new().produce();
1328 let consumer = producer.consume();
1329 drop(producer);
1331 assert!(matches!(consumer.closed().await, Error::Dropped));
1332 assert!(!consumer.is_finished());
1333 }
1334
1335 #[tokio::test]
1336 async fn requests() {
1337 let mut producer = Info::new().produce().dynamic();
1338
1339 let consumer = producer.consume();
1340 let consumer2 = consumer.clone();
1341
1342 let track1_fut = subscribe_pending!(consumer, "track1");
1344 let track2_fut = subscribe_pending!(consumer2, "track1");
1345
1346 let request = producer.assert_request();
1348 producer.assert_no_request();
1349 assert_eq!(request.name(), "track1");
1350
1351 let track3 = request.accept(None);
1353 let mut track1 = track1_fut.await.unwrap();
1354 let mut track2 = track2_fut.await.unwrap();
1355
1356 track1.assert_not_closed();
1357 track1.assert_is_clone(&track2);
1358 track3.subscribe(None).assert_is_clone(&track1);
1359
1360 track3.append_group().unwrap();
1362 track1.assert_group();
1363 track2.assert_group();
1364
1365 let track4_fut = subscribe_pending!(consumer, "track2");
1367 drop(producer);
1368 assert!(track4_fut.await.is_err());
1369
1370 let track5 = consumer2.track("track3");
1372 assert!(track5.is_err(), "should have errored");
1373 }
1374
1375 #[tokio::test]
1376 async fn stale_producer() {
1377 let mut broadcast = Info::new().produce().dynamic();
1378 let consumer = broadcast.consume();
1379
1380 let track1_fut = subscribe_pending!(consumer, "track1");
1382 let producer1 = broadcast.assert_request().accept(None);
1383 let mut track1 = track1_fut.await.unwrap();
1384
1385 producer1.append_group().unwrap();
1387 producer1.finish().unwrap();
1388 drop(producer1);
1389
1390 track1.assert_closed();
1392
1393 let track2_fut = subscribe_pending!(consumer, "track1");
1395 let producer2 = broadcast.assert_request().accept(None);
1396 let mut track2 = track2_fut.await.unwrap();
1397 track2.assert_not_closed();
1398 track2.assert_not_clone(&track1);
1399
1400 producer2.append_group().unwrap();
1402 track2.assert_group();
1403 }
1404
1405 #[tokio::test(start_paused = true)]
1406 async fn requested_unused() {
1407 let mut broadcast = Info::new().produce().dynamic();
1408 let bc = broadcast.consume();
1409
1410 let c1_fut = subscribe_pending!(bc, "unknown_track");
1412 let producer1 = broadcast.assert_request().accept(None);
1413 let consumer1 = c1_fut.await.unwrap();
1414
1415 assert!(
1417 producer1.unused().now_or_never().is_none(),
1418 "track producer should be used"
1419 );
1420
1421 let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1423 consumer2.assert_is_clone(&consumer1);
1424
1425 drop(consumer1);
1426 assert!(
1427 producer1.unused().now_or_never().is_none(),
1428 "track producer should be used"
1429 );
1430
1431 drop(consumer2);
1432 assert!(
1433 producer1.unused().now_or_never().is_some(),
1434 "track producer should be unused after all consumers are dropped"
1435 );
1436
1437 let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1441 consumer3.assert_is_clone(&producer1.subscribe(None));
1442 broadcast.assert_no_request();
1443 drop(consumer3);
1444
1445 producer1.abort(Error::Cancel).unwrap();
1448
1449 let c4_fut = subscribe_pending!(bc, "unknown_track");
1450 let producer2 = broadcast.assert_request().accept(None);
1451 let consumer4 = c4_fut.await.unwrap();
1452 drop(consumer4);
1453 assert!(
1454 producer2.unused().now_or_never().is_some(),
1455 "new track producer should be unused after its consumer is dropped"
1456 );
1457 }
1458
1459 #[tokio::test]
1465 async fn create_track_fulfills_queued_request() {
1466 let producer = Info::new().produce();
1467 let mut dynamic = producer.dynamic();
1468 let bc = dynamic.consume();
1469
1470 let subscribing = subscribe_pending!(bc, "video");
1472
1473 let track = producer.create_track("video", None).unwrap();
1475 let mut sub = subscribing.await.expect("fulfilled by create_track");
1476
1477 track.append_group().unwrap();
1479 sub.recv_group().await.expect("recv").expect("group");
1480
1481 dynamic.assert_no_request();
1483 let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
1484 again.assert_is_clone(&track.subscribe(None));
1485 }
1486
1487 #[tokio::test]
1493 async fn dynamic_clone_keeps_alive() {
1494 let broadcast = Info::new().produce().dynamic();
1495 let consumer = broadcast.consume();
1496
1497 let clone = broadcast.clone();
1498 drop(clone);
1499
1500 let _fut = subscribe_pending!(consumer, "track1");
1503 }
1504
1505 #[tokio::test]
1510 async fn finish_resolves_a_reserved_name() {
1511 let producer = Info::new().produce();
1512 let consumer = producer.consume();
1513
1514 let _request = producer.reserve_track("track1").unwrap();
1515 let pending = subscribe_pending!(consumer, "track1");
1516
1517 producer.finish();
1518 assert!(matches!(pending.await, Err(Error::NotFound)));
1519 }
1520
1521 #[tokio::test]
1524 async fn abort_resolves_a_reserved_name_with_its_reason() {
1525 let producer = Info::new().produce();
1526 let consumer = producer.consume();
1527
1528 let request = producer.reserve_track("track1").unwrap();
1529 let pending = subscribe_pending!(consumer, "track1");
1530
1531 producer.abort(Error::Cancel).unwrap();
1532 assert!(matches!(pending.await, Err(Error::Cancel)));
1533
1534 let track = request.accept(None);
1535 let mut subscriber = track.subscribe(None);
1536 assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
1537 }
1538
1539 #[tokio::test]
1542 async fn finish_resolves_a_queued_request() {
1543 let producer = Info::new().produce();
1544 let dynamic = producer.dynamic();
1545 let consumer = dynamic.consume();
1546
1547 let pending = subscribe_pending!(consumer, "track1");
1548
1549 producer.finish();
1550 assert!(matches!(pending.await, Err(Error::NotFound)));
1551 drop(dynamic);
1552 }
1553
1554 #[tokio::test]
1558 async fn finish_leaves_a_claimed_request_to_its_handler() {
1559 let producer = Info::new().produce();
1560 let mut dynamic = producer.dynamic();
1561 let consumer = dynamic.consume();
1562
1563 let accepted = subscribe_pending!(consumer, "track1");
1564 let request = dynamic.requested_track().await.unwrap();
1565 let dropped = subscribe_pending!(consumer, "track2");
1566 let abandoned = dynamic.requested_track().await.unwrap();
1567
1568 producer.finish();
1569 assert!(
1570 accepted.poll_ok(&kio::Waiter::noop()).is_pending(),
1571 "finish rejected a claimed request"
1572 );
1573 assert!(
1574 dropped.poll_ok(&kio::Waiter::noop()).is_pending(),
1575 "finish rejected a claimed request"
1576 );
1577
1578 let _track = request.accept(None);
1579 assert!(accepted.await.is_ok(), "the handler's accept reaches the consumer");
1580 drop(abandoned);
1581 assert!(dropped.await.is_err(), "the handler dropping it rejects the consumer");
1582 drop(dynamic);
1583 }
1584
1585 #[tokio::test]
1590 async fn finish_resolves_an_unaccepted_track_with_fetched_info() {
1591 let producer = Info::new().produce();
1592 let consumer = producer.consume();
1593
1594 let request = producer.reserve_track("track1").unwrap();
1595 let dynamic = request.dynamic();
1596 let track = consumer.track("track1").unwrap();
1597 let pending_fetch = track.fetch_group(0, None);
1598 let fetch = dynamic.requested_group().await.unwrap();
1599 let group = fetch.accept(None).unwrap();
1600 group.finish().unwrap();
1601 pending_fetch.await.unwrap();
1602
1603 let mut subscriber = track.subscribe(None).await.unwrap();
1604 producer.finish();
1605 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1606
1607 let stale = request.accept(None);
1608 assert!(stale.append_group().is_err());
1609 }
1610
1611 #[tokio::test]
1614 async fn finish_spares_a_served_track() {
1615 let producer = Info::new().produce();
1616 let consumer = producer.consume();
1617
1618 let track = producer.create_track("track1", None).unwrap();
1619 let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1620
1621 producer.finish();
1622
1623 track.append_group().unwrap();
1624 subscriber.assert_group();
1625 track.finish().unwrap();
1626 }
1627
1628 #[tokio::test]
1632 async fn finish_leaves_a_stale_reservation_inert() {
1633 let producer = Info::new().produce();
1634 let consumer = producer.consume();
1635
1636 let request = producer.reserve_track("track1").unwrap();
1637 let pending = subscribe_pending!(consumer, "track1");
1638
1639 producer.finish();
1640 assert!(matches!(pending.await, Err(Error::NotFound)));
1641
1642 let track = request.accept(None);
1643 assert!(track.append_group().is_err());
1644 let mut subscriber = track.subscribe(None);
1645 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1646 assert!(consumer.track("track1").is_err());
1647 }
1648
1649 #[tokio::test]
1653 async fn dropping_a_reserved_request_resolves_dropped() {
1654 let producer = Info::new().produce();
1655 let consumer = producer.consume();
1656
1657 let request = producer.reserve_track("track1").unwrap();
1658 let pending = subscribe_pending!(consumer, "track1");
1659
1660 drop(request);
1661 assert!(matches!(pending.await, Err(Error::Dropped)));
1662 producer.finish();
1663 }
1664
1665 #[tokio::test]
1668 async fn rejecting_a_reserved_request_carries_the_reason() {
1669 let producer = Info::new().produce();
1670 let consumer = producer.consume();
1671
1672 let request = producer.reserve_track("track1").unwrap();
1673 let pending = subscribe_pending!(consumer, "track1");
1674
1675 request.reject(Error::NotFound);
1676 assert!(matches!(pending.await, Err(Error::NotFound)));
1677 producer.finish();
1678 }
1679
1680 #[tokio::test]
1686 async fn an_idle_teardown_yields_to_a_returning_viewer() {
1687 let producer = Info::new().produce();
1688 let consumer = producer.consume();
1689 let track = producer.create_track("video", None).unwrap();
1690
1691 assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
1693
1694 let viewer = consumer.track("video").unwrap();
1696 let track = track
1697 .abort_unused(Error::Cancel)
1698 .expect_err("viewer keeps the track alive");
1699
1700 assert!(!track.is_closed());
1702 let mut subscriber = viewer.subscribe(None).await.unwrap();
1703 subscriber.assert_no_group();
1704 track.append_group().unwrap();
1705 assert!(subscriber.recv_group().await.unwrap().is_some());
1706
1707 drop(subscriber);
1710 drop(viewer);
1711 assert!(track.abort_unused(Error::Cancel).is_ok());
1712 assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
1713
1714 producer.finish();
1715 }
1716
1717 #[test]
1718 fn abort_unused_accepts_an_already_closed_track_with_consumers() {
1719 let producer = Info::new().produce();
1720 let consumer = producer.consume();
1721 let track = producer.create_track("video", None).unwrap();
1722 let _viewer = consumer.track("video").unwrap();
1723 assert!(track.is_used());
1724 track.clone().abort(Error::Cancel).unwrap();
1725 assert!(!track.is_used());
1726 assert!(track.abort_unused(Error::Cancel).is_ok());
1727 producer.finish();
1728 }
1729}