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> {
258 let mut announcer = self.alive.announcer.lock();
259 let announcer = announcer.as_mut().ok_or(Error::Closed)?;
260 announcer.announce(route)
261 }
262
263 pub fn unannounce(&self) {
266 self.alive.unannounce();
267 }
268
269 pub(crate) fn new_spliced(info: Info) -> Self {
273 let state = kio::Shared::new(BroadcastState {
274 spliced: Some(SplicedState::default()),
275 ..Default::default()
276 });
277 Self {
278 info: Arc::new(info),
279 alive: Alive::new(state.clone()),
280 state,
281 stats: stats::Scope::default(),
284 }
285 }
286
287 pub fn info(&self) -> &Info {
289 &self.info
290 }
291
292 pub fn demand(&self) -> Demand {
294 Demand {
295 alive: self.alive.token.consume().weak(),
296 state: self.state.clone(),
297 }
298 }
299
300 pub fn create_track(
305 &self,
306 name: impl Into<Arc<str>>,
307 info: impl Into<Option<track::Info>>,
308 ) -> Result<track::Producer, Error> {
309 let name = name.into();
310 let info = info.into().unwrap_or_default();
311 let mut state = self.state.lock();
312
313 if let Some(request) = state.requests.take(name.as_ref()) {
319 let track = request.with_stats(self.stats.clone()).accept(info);
320 let _ = state.tracks.insert(name, track.weak());
324 return Ok(track);
325 }
326
327 let track = track::Producer::new(self.info.clone(), name, info).with_stats(self.stats.clone());
328 state.insert_track(track.weak())?;
329 Ok(track)
330 }
331
332 pub fn reserve_track(&self, name: impl Into<Arc<str>>) -> Result<track::Request, Error> {
344 let request = track::Request::new(self.info.clone(), name).with_stats(self.stats.clone());
345 self.state.lock().insert_track(request.weak())?;
346 Ok(request)
347 }
348
349 pub fn unique_track(&self, suffix: &str, info: impl Into<Option<track::Info>>) -> Result<track::Producer, Error> {
353 let name = self.unique_name(suffix);
354 self.create_track(name, info)
355 }
356
357 pub fn unique_name(&self, suffix: &str) -> String {
369 let mut state = self.state.lock();
370 let separator = if suffix.starts_with(|c: char| c.is_ascii_digit()) {
371 "-"
372 } else {
373 ""
374 };
375 loop {
376 let id = state.unique;
377 state.unique = id.checked_add(1).expect("unique track IDs exhausted");
378 let name = format!("{id}{separator}{suffix}");
379 if !state.tracks.contains_key(name.as_str()) {
380 return name;
381 }
382 }
383 }
384
385 pub fn dynamic(&self) -> Dynamic {
387 Dynamic::new(
388 self.info.clone(),
389 self.alive.clone(),
390 self.state.clone(),
391 self.stats.clone(),
392 )
393 }
394
395 pub(crate) fn poll_spliced_assigned(&self, waiter: &kio::Waiter) -> Poll<(Arc<str>, super::resume::Producer)> {
398 let mut state = ready!(self.state.poll(waiter, |state| {
399 match &state.spliced {
400 Some(spliced) if !spliced.pending.is_empty() => Poll::Ready(()),
401 _ => Poll::Pending,
402 }
403 }));
404
405 let spliced = state.spliced.as_mut().expect("predicate guaranteed spliced");
406 let name = spliced.pending.pop_front().expect("predicate guaranteed a request");
407 let producer = spliced.tracks.get(&name).expect("pending name without a track").clone();
408 Poll::Ready((name, producer))
409 }
410
411 pub(crate) fn release_spliced(&self, err: Error) {
415 let mut state = self.state.lock();
416 if let Some(spliced) = state.spliced.as_mut() {
417 for name in std::mem::take(&mut spliced.pending) {
418 if let Some(producer) = spliced.tracks.get_mut(&name) {
419 let _ = producer.abort(err.clone());
420 }
421 }
422 spliced.tracks.clear();
423 }
424 }
425
426 pub fn consume(&self) -> Consumer {
428 Consumer {
429 info: self.info.clone(),
430 alive: self.alive.token.consume(),
431 state: self.state.clone(),
432 stats: stats::Scope::default(),
433 }
434 }
435
436 pub fn finish(&self) {
454 {
455 let mut state = self.state.lock();
456 state.closing = true;
457 state.finished = true;
458 state.reject_unserved(Error::NotFound);
462 }
463 let _ = self.alive.token.close();
466 self.alive.retire();
467 }
468
469 pub fn abort(self, err: Error) -> Result<(), Error> {
482 {
483 let mut state = self.state.lock();
484 if state.closing {
485 return Err(Error::Closed);
486 }
487 state.closing = true;
488 state.abort = Some(err.clone());
489 state.reject_unserved(err);
492 }
493 let _ = self.alive.token.close();
494 self.alive.retire();
495 Ok(())
496 }
497
498 pub fn is_clone(&self, other: &Self) -> bool {
500 self.state.same_channel(&other.state)
501 }
502}
503
504struct Alive {
510 token: kio::Producer<()>,
511 state: kio::Shared<BroadcastState>,
512 announcer: kio::Lock<Option<Announcer>>,
516}
517
518impl Alive {
519 fn new(state: kio::Shared<BroadcastState>) -> Arc<Self> {
520 Arc::new(Self {
521 token: kio::Producer::default(),
522 state,
523 announcer: kio::Lock::new(None),
524 })
525 }
526
527 fn unannounce(&self) {
529 if let Some(announcer) = self.announcer.lock().as_mut() {
530 announcer.withdraw();
531 }
532 }
533
534 fn retire(&self) {
537 let announcer = self.announcer.lock().take();
538 drop(announcer);
541 }
542}
543
544impl Drop for Alive {
545 fn drop(&mut self) {
546 if !self.state.read().closing {
550 tracing::warn!(
551 "broadcast::Producer dropped without finish(). Keep the producer alive while publishing, then call finish()."
552 );
553 }
554 self.retire();
555 }
556}
557
558#[cfg(test)]
559#[allow(missing_docs)] impl Producer {
561 pub fn assert_create_track(
562 &mut self,
563 name: impl Into<Arc<str>>,
564 info: impl Into<Option<track::Info>>,
565 ) -> track::Producer {
566 self.create_track(name, info).expect("should not have errored")
567 }
568}
569
570pub(crate) struct SourceGuard {
576 producer: Option<Producer>,
578}
579
580impl SourceGuard {
581 pub fn new(producer: Producer) -> Self {
582 Self {
583 producer: Some(producer),
584 }
585 }
586
587 pub fn finish(mut self) {
590 if let Some(producer) = self.producer.take() {
591 producer.finish();
592 }
593 }
594}
595
596impl Drop for SourceGuard {
597 fn drop(&mut self) {
598 if let Some(producer) = self.producer.take() {
599 let _ = producer.abort(Error::Dropped);
600 }
601 }
602}
603
604pub struct Dynamic {
612 info: Arc<Info>,
613 alive: Arc<Alive>,
615 state: kio::Shared<BroadcastState>,
616 stats: stats::Scope,
619}
620
621impl Clone for Dynamic {
622 fn clone(&self) -> Self {
623 self.state.lock().requests.add_handler();
627
628 Self {
629 info: self.info.clone(),
630 alive: self.alive.clone(),
631 state: self.state.clone(),
632 stats: self.stats.clone(),
633 }
634 }
635}
636
637impl Dynamic {
638 fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
639 state.lock().requests.add_handler();
640
641 Self {
642 info,
643 alive,
644 state,
645 stats,
646 }
647 }
648
649 pub fn info(&self) -> &Info {
651 &self.info
652 }
653
654 pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
660 let mut state = ready!(self.state.poll(waiter, |state| {
661 if state.requests.has_queued() || state.closing {
662 Poll::Ready(())
663 } else {
664 Poll::Pending
665 }
666 }));
667
668 if state.closing && !state.requests.has_queued() {
669 return Poll::Ready(Err(Error::Closed));
670 }
671
672 let name = state.requests.pop().expect("predicate guaranteed a request");
673 let pending = state.requests.remove(&name).expect("popped key must be pending");
674 let _ = state.tracks.insert(name, pending.weak());
677 Poll::Ready(Ok(pending.with_stats(self.stats.clone())))
679 }
680
681 pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
683 kio::wait(|waiter| self.poll_requested_track(waiter)).await
684 }
685
686 pub fn consume(&self) -> Consumer {
688 Consumer {
689 info: self.info.clone(),
690 alive: self.alive.token.consume(),
691 state: self.state.clone(),
692 stats: stats::Scope::default(),
693 }
694 }
695
696 pub async fn closed(&self) -> Error {
699 kio::wait(|waiter| self.poll_closed(waiter)).await
700 }
701
702 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
706 ready!(self.alive.token.poll_closed(waiter));
707 Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
708 }
709
710 pub fn is_clone(&self, other: &Self) -> bool {
712 self.state.same_channel(&other.state)
713 }
714}
715
716impl Drop for Dynamic {
717 fn drop(&mut self) {
718 let mut state = self.state.lock();
721 if state.requests.remove_handler() {
722 for request in state.requests.drain_queued() {
725 request.reject(Error::Dropped);
726 }
727 }
728 }
729}
730
731#[cfg(test)]
732use futures::FutureExt;
733
734#[cfg(test)]
735#[allow(missing_docs)] impl Dynamic {
737 pub fn assert_request(&mut self) -> track::Request {
738 self.requested_track()
739 .now_or_never()
740 .expect("should not have blocked")
741 .expect("should not have errored")
742 }
743
744 pub fn assert_no_request(&mut self) {
745 assert!(self.requested_track().now_or_never().is_none(), "should have blocked");
746 }
747}
748
749pub struct Consumer {
751 info: Arc<Info>,
752 alive: kio::Consumer<()>,
754 state: kio::Shared<BroadcastState>,
756 stats: stats::Scope,
760}
761
762impl Clone for Consumer {
763 fn clone(&self) -> Self {
764 Self {
765 info: self.info.clone(),
766 alive: self.alive.clone(),
767 state: self.state.clone(),
768 stats: self.stats.clone(),
769 }
770 }
771}
772
773impl Consumer {
774 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
777 self.stats = scope;
778 self
779 }
780
781 pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
789 if self.info.path != path {
790 let mut info = (*self.info).clone();
791 info.path = path;
792 self.info = Arc::new(info);
793 }
794 self
795 }
796
797 pub fn info(&self) -> &Info {
799 &self.info
800 }
801
802 pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
804 self.track_inner(name)
809 .map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
810 }
811
812 fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
813 if self.is_closed() {
815 return Err(Error::Dropped);
816 }
817
818 let mut state = self.state.lock();
819
820 let closing = state.closing;
823 if let Some(spliced) = state.spliced.as_mut() {
824 if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
838 spliced.tracks.remove(name);
839 }
840 if let Some(producer) = spliced.tracks.get(name) {
841 return Ok(track::Consumer::spliced(
842 name.into(),
843 self.info.clone(),
844 producer.consume(),
845 ));
846 }
847 if closing {
850 return Err(Error::NotFound);
851 }
852 let name: Arc<str> = name.into();
853 let producer = super::resume::Producer::new();
854 let consumer = producer.consume();
855 spliced.tracks.insert(name.clone(), producer);
856 spliced.pending.push_back(name.clone());
857 return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
858 }
859
860 if let Some(weak) = state.tracks.get(name) {
863 match weak.try_consume() {
864 Some(consumer) => return Ok(consumer),
865 None => {
869 state.tracks.remove(name);
870 }
871 }
872 }
873
874 if let Some(pending) = state.requests.join(name) {
875 return Ok(pending.consume());
877 }
878
879 if state.closing {
882 return Err(Error::NotFound);
883 }
884
885 let name: Arc<str> = name.into();
889 let request = track::Request::new(self.info.clone(), name.clone());
890 let consumer = request.consume();
891
892 if state.requests.insert(name, request).is_err() {
895 return Err(Error::NotFound);
896 }
897
898 Ok(consumer)
899 }
900
901 pub fn demand(&self) -> Demand {
910 Demand {
911 alive: self.alive.weak(),
912 state: self.state.clone(),
913 }
914 }
915
916 pub async fn closed(&self) -> Error {
923 self.alive.closed().await;
924 self.state.read().abort.clone().unwrap_or(Error::Dropped)
925 }
926
927 pub fn is_closed(&self) -> bool {
929 self.alive.is_closed()
930 }
931
932 pub(crate) fn is_closing(&self) -> bool {
937 self.is_closed() || self.state.read().closing
938 }
939
940 pub fn is_finished(&self) -> bool {
943 self.state.read().finished
944 }
945
946 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
952 self.alive.poll_closed(waiter)
953 }
954
955 pub fn is_clone(&self, other: &Self) -> bool {
957 self.state.same_channel(&other.state)
958 }
959
960 pub(crate) fn weak(&self) -> WeakConsumer {
965 WeakConsumer {
966 info: self.info.clone(),
967 alive: self.alive.weak(),
968 state: self.state.clone(),
969 }
970 }
971}
972
973#[derive(Clone)]
980pub(crate) struct WeakConsumer {
981 info: Arc<Info>,
982 alive: kio::ConsumerWeak<()>,
983 state: kio::Shared<BroadcastState>,
984}
985
986impl WeakConsumer {
987 pub fn consume(&self) -> Consumer {
989 Consumer {
990 info: self.info.clone(),
991 alive: self.alive.consume(),
992 state: self.state.clone(),
993 stats: stats::Scope::default(),
994 }
995 }
996}
997
998impl super::WeakEntry for WeakConsumer {
999 fn is_closed(&self) -> bool {
1000 self.alive.is_closed()
1001 }
1002
1003 fn same_channel(&self, other: &Self) -> bool {
1004 self.state.same_channel(&other.state)
1005 }
1006}
1007
1008#[derive(Clone)]
1021pub struct Demand {
1022 alive: kio::ConsumerWeak<()>,
1023 state: kio::Shared<BroadcastState>,
1024}
1025
1026impl Demand {
1027 pub fn is_used(&self) -> bool {
1032 self.state.read().is_used()
1033 }
1034
1035 pub async fn used(&self) -> Result<(), Error> {
1038 kio::wait(|waiter| self.poll_used(waiter)).await
1039 }
1040
1041 pub async fn unused(&self) -> Result<(), Error> {
1044 kio::wait(|waiter| self.poll_unused(waiter)).await
1045 }
1046
1047 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1049 self.poll_demand(waiter, true)
1050 }
1051
1052 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1054 self.poll_demand(waiter, false)
1055 }
1056
1057 fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
1058 if self.alive.poll_closed(waiter).is_ready() {
1061 return Poll::Ready(Err(Error::Dropped));
1062 }
1063 let ready = self.state.poll(waiter, |state| {
1064 state.register_demand(waiter, want);
1068 match state.is_used() == want {
1069 true => Poll::Ready(()),
1070 false => Poll::Pending,
1071 }
1072 });
1073 match ready {
1074 Poll::Ready(_) => Poll::Ready(Ok(())),
1075 Poll::Pending => Poll::Pending,
1076 }
1077 }
1078}
1079
1080#[cfg(test)]
1081#[allow(missing_docs)] impl Consumer {
1083 pub fn assert_not_closed(&self) {
1084 assert!(self.closed().now_or_never().is_none(), "should not be closed");
1085 }
1086
1087 pub fn assert_closed(&self) {
1088 assert!(self.closed().now_or_never().is_some(), "should be closed");
1089 }
1090}
1091
1092#[cfg(test)]
1093mod test {
1094 use super::*;
1095 use std::time::Duration;
1096
1097 #[test]
1098 fn unique_names_are_never_reused() {
1099 let producer = Info::new().produce();
1100 let name = producer.unique_name(".opus");
1101 assert_eq!(name, "0.opus");
1102 let track = producer.create_track(name.clone(), None).unwrap();
1103 assert_eq!(producer.unique_name(".opus"), "1.opus");
1104 drop(track);
1105 }
1106
1107 #[test]
1108 fn unique_names_survive_closed_track_pruning() {
1109 let producer = Info::new().produce();
1110 let consumer = producer.consume();
1111 let track = producer.unique_track(".opus", None).unwrap();
1112 assert_eq!(track.name(), "0.opus");
1113 drop(track);
1114 assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
1115 assert_eq!(producer.unique_name(".opus"), "1.opus");
1116 }
1117
1118 #[test]
1119 fn unique_name_skips_a_live_collision() {
1120 let producer = Info::new().produce();
1121 let track = producer.create_track("0.opus", None).unwrap();
1122 assert_eq!(producer.unique_name(".opus"), "1.opus");
1123 drop(track);
1124 assert_eq!(producer.unique_name(".opus"), "2.opus");
1125 }
1126
1127 #[test]
1128 fn unique_names_share_a_counter() {
1129 let producer = Info::new().produce();
1130 assert_eq!(producer.unique_name("-video"), "0-video");
1131 assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
1132 assert_eq!(producer.unique_name("-video"), "2-video");
1133 }
1134
1135 #[test]
1136 fn unique_names_separate_numeric_suffixes() {
1137 let producer = Info::new().produce();
1138 assert_eq!(producer.unique_name(""), "0");
1139 let name = producer.unique_name("2");
1140 assert_eq!(name, "1-2");
1141 for _ in 2..12 {
1142 producer.unique_name("");
1143 }
1144 assert_eq!(producer.unique_name(""), "12");
1145 }
1146
1147 async fn expect<T>(fut: impl Future<Output = T>) -> T {
1150 tokio::time::timeout(Duration::from_secs(1), fut)
1151 .await
1152 .expect("timed out waiting for a demand edge")
1153 }
1154
1155 #[tokio::test]
1159 async fn demand_ordinary() {
1160 tokio::time::pause();
1161
1162 let producer = Info::new().produce();
1163 let consumer = producer.consume();
1164 let demand = producer.demand();
1165
1166 assert!(!demand.is_used());
1168 demand.unused().await.unwrap();
1169
1170 let _track = producer.create_track("a", None).unwrap();
1172 assert!(!demand.is_used());
1173
1174 let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
1176 used.unwrap();
1177 assert!(demand.is_used());
1178
1179 let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
1181 unused.unwrap();
1182 assert!(!demand.is_used());
1183
1184 producer.finish();
1186 assert!(matches!(demand.used().await, Err(Error::Dropped)));
1187 assert!(matches!(demand.unused().await, Err(Error::Dropped)));
1188 }
1189
1190 #[tokio::test]
1193 async fn demand_spliced() {
1194 tokio::time::pause();
1195
1196 let producer = Producer::new_spliced(Info::new());
1197 let consumer = producer.consume();
1198 let demand = producer.demand();
1199 let watched = consumer.demand();
1200
1201 assert!(!demand.is_used());
1202 assert!(!watched.is_used());
1203 let track = consumer.track("video").unwrap();
1204 assert!(demand.is_used());
1205 assert!(watched.is_used());
1206
1207 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1210 unused.unwrap();
1211 assert!(!demand.is_used());
1212 assert!(!watched.is_used());
1213
1214 let _track = consumer.track("video").unwrap();
1216 assert!(demand.is_used());
1217 }
1218
1219 #[tokio::test]
1221 async fn consumer_demand_reports_dropped_producer() {
1222 let producer = Producer::new_spliced(Info::new());
1223 let consumer = producer.consume();
1224 let watched = consumer.demand();
1225
1226 let track = consumer.track("video").unwrap();
1227 assert!(watched.is_used());
1228
1229 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1230 unused.unwrap();
1231
1232 drop(producer);
1233 assert!(matches!(watched.used().await, Err(Error::Dropped)));
1234 assert!(matches!(watched.unused().await, Err(Error::Dropped)));
1235 }
1236
1237 macro_rules! subscribe_pending {
1240 ($consumer:expr, $name:expr) => {{
1241 let pending = $consumer.track($name).unwrap().subscribe(None);
1242 assert!(
1243 pending.poll_ok(&kio::Waiter::noop()).is_pending(),
1244 "subscribe should stay pending until the request is accepted"
1245 );
1246 pending
1247 }};
1248 }
1249
1250 #[tokio::test]
1251 async fn insert() {
1252 let mut producer = Info::new().produce();
1253
1254 let track1 = producer.assert_create_track("track1", None);
1256 track1.append_group().unwrap();
1257
1258 let consumer = producer.consume();
1259
1260 let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1262 track1_sub.assert_group();
1263
1264 let track2 = producer.assert_create_track("track2", None);
1265
1266 let consumer2 = producer.consume();
1267 let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
1268 track2_consumer.assert_no_group();
1269
1270 track2.append_group().unwrap();
1271
1272 track2_consumer.assert_group();
1273 }
1274
1275 #[tokio::test]
1276 async fn closed() {
1277 let mut producer = Info::new().produce();
1278 let dynamic = producer.dynamic();
1279
1280 let consumer = producer.consume();
1281 consumer.assert_not_closed();
1282
1283 let track1 = producer.assert_create_track("track1", None);
1285 let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1286
1287 let track2_fut = subscribe_pending!(consumer, "track2");
1289
1290 drop(dynamic);
1293
1294 assert!(track2_fut.await.is_err());
1296
1297 assert!(!track1.is_closed());
1299 track1c.assert_not_closed();
1300 }
1301
1302 #[tokio::test]
1305 async fn closed_cause() {
1306 let producer = Info::new().produce();
1308 let consumer = producer.consume();
1309 producer.abort(Error::Timeout).unwrap();
1310 assert!(matches!(consumer.closed().await, Error::Timeout));
1311 assert!(!consumer.is_finished());
1312
1313 let producer = Info::new().produce();
1315 let consumer = producer.consume();
1316 producer.finish();
1317 assert!(matches!(consumer.closed().await, Error::Dropped));
1318 assert!(consumer.is_finished());
1319
1320 let producer = Info::new().produce();
1322 let consumer = producer.consume();
1323 drop(producer);
1325 assert!(matches!(consumer.closed().await, Error::Dropped));
1326 assert!(!consumer.is_finished());
1327 }
1328
1329 #[tokio::test]
1330 async fn requests() {
1331 let mut producer = Info::new().produce().dynamic();
1332
1333 let consumer = producer.consume();
1334 let consumer2 = consumer.clone();
1335
1336 let track1_fut = subscribe_pending!(consumer, "track1");
1338 let track2_fut = subscribe_pending!(consumer2, "track1");
1339
1340 let request = producer.assert_request();
1342 producer.assert_no_request();
1343 assert_eq!(request.name(), "track1");
1344
1345 let track3 = request.accept(None);
1347 let mut track1 = track1_fut.await.unwrap();
1348 let mut track2 = track2_fut.await.unwrap();
1349
1350 track1.assert_not_closed();
1351 track1.assert_is_clone(&track2);
1352 track3.subscribe(None).assert_is_clone(&track1);
1353
1354 track3.append_group().unwrap();
1356 track1.assert_group();
1357 track2.assert_group();
1358
1359 let track4_fut = subscribe_pending!(consumer, "track2");
1361 drop(producer);
1362 assert!(track4_fut.await.is_err());
1363
1364 let track5 = consumer2.track("track3");
1366 assert!(track5.is_err(), "should have errored");
1367 }
1368
1369 #[tokio::test]
1370 async fn stale_producer() {
1371 let mut broadcast = Info::new().produce().dynamic();
1372 let consumer = broadcast.consume();
1373
1374 let track1_fut = subscribe_pending!(consumer, "track1");
1376 let producer1 = broadcast.assert_request().accept(None);
1377 let mut track1 = track1_fut.await.unwrap();
1378
1379 producer1.append_group().unwrap();
1381 producer1.finish().unwrap();
1382 drop(producer1);
1383
1384 track1.assert_closed();
1386
1387 let track2_fut = subscribe_pending!(consumer, "track1");
1389 let producer2 = broadcast.assert_request().accept(None);
1390 let mut track2 = track2_fut.await.unwrap();
1391 track2.assert_not_closed();
1392 track2.assert_not_clone(&track1);
1393
1394 producer2.append_group().unwrap();
1396 track2.assert_group();
1397 }
1398
1399 #[tokio::test(start_paused = true)]
1400 async fn requested_unused() {
1401 let mut broadcast = Info::new().produce().dynamic();
1402 let bc = broadcast.consume();
1403
1404 let c1_fut = subscribe_pending!(bc, "unknown_track");
1406 let producer1 = broadcast.assert_request().accept(None);
1407 let consumer1 = c1_fut.await.unwrap();
1408
1409 assert!(
1411 producer1.unused().now_or_never().is_none(),
1412 "track producer should be used"
1413 );
1414
1415 let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1417 consumer2.assert_is_clone(&consumer1);
1418
1419 drop(consumer1);
1420 assert!(
1421 producer1.unused().now_or_never().is_none(),
1422 "track producer should be used"
1423 );
1424
1425 drop(consumer2);
1426 assert!(
1427 producer1.unused().now_or_never().is_some(),
1428 "track producer should be unused after all consumers are dropped"
1429 );
1430
1431 let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1435 consumer3.assert_is_clone(&producer1.subscribe(None));
1436 broadcast.assert_no_request();
1437 drop(consumer3);
1438
1439 producer1.abort(Error::Cancel).unwrap();
1442
1443 let c4_fut = subscribe_pending!(bc, "unknown_track");
1444 let producer2 = broadcast.assert_request().accept(None);
1445 let consumer4 = c4_fut.await.unwrap();
1446 drop(consumer4);
1447 assert!(
1448 producer2.unused().now_or_never().is_some(),
1449 "new track producer should be unused after its consumer is dropped"
1450 );
1451 }
1452
1453 #[tokio::test]
1459 async fn create_track_fulfills_queued_request() {
1460 let producer = Info::new().produce();
1461 let mut dynamic = producer.dynamic();
1462 let bc = dynamic.consume();
1463
1464 let subscribing = subscribe_pending!(bc, "video");
1466
1467 let track = producer.create_track("video", None).unwrap();
1469 let mut sub = subscribing.await.expect("fulfilled by create_track");
1470
1471 track.append_group().unwrap();
1473 sub.recv_group().await.expect("recv").expect("group");
1474
1475 dynamic.assert_no_request();
1477 let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
1478 again.assert_is_clone(&track.subscribe(None));
1479 }
1480
1481 #[tokio::test]
1487 async fn dynamic_clone_keeps_alive() {
1488 let broadcast = Info::new().produce().dynamic();
1489 let consumer = broadcast.consume();
1490
1491 let clone = broadcast.clone();
1492 drop(clone);
1493
1494 let _fut = subscribe_pending!(consumer, "track1");
1497 }
1498
1499 #[tokio::test]
1504 async fn finish_resolves_a_reserved_name() {
1505 let producer = Info::new().produce();
1506 let consumer = producer.consume();
1507
1508 let _request = producer.reserve_track("track1").unwrap();
1509 let pending = subscribe_pending!(consumer, "track1");
1510
1511 producer.finish();
1512 assert!(matches!(pending.await, Err(Error::NotFound)));
1513 }
1514
1515 #[tokio::test]
1518 async fn abort_resolves_a_reserved_name_with_its_reason() {
1519 let producer = Info::new().produce();
1520 let consumer = producer.consume();
1521
1522 let request = producer.reserve_track("track1").unwrap();
1523 let pending = subscribe_pending!(consumer, "track1");
1524
1525 producer.abort(Error::Cancel).unwrap();
1526 assert!(matches!(pending.await, Err(Error::Cancel)));
1527
1528 let track = request.accept(None);
1529 let mut subscriber = track.subscribe(None);
1530 assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
1531 }
1532
1533 #[tokio::test]
1536 async fn finish_resolves_a_queued_request() {
1537 let producer = Info::new().produce();
1538 let dynamic = producer.dynamic();
1539 let consumer = dynamic.consume();
1540
1541 let pending = subscribe_pending!(consumer, "track1");
1542
1543 producer.finish();
1544 assert!(matches!(pending.await, Err(Error::NotFound)));
1545 drop(dynamic);
1546 }
1547
1548 #[tokio::test]
1551 async fn finish_resolves_a_request_a_handler_never_answered() {
1552 let producer = Info::new().produce();
1553 let mut dynamic = producer.dynamic();
1554 let consumer = dynamic.consume();
1555
1556 let pending = subscribe_pending!(consumer, "track1");
1557 let _request = dynamic.requested_track().await.unwrap();
1558
1559 producer.finish();
1560 assert!(matches!(pending.await, Err(Error::NotFound)));
1561 drop(dynamic);
1562 }
1563
1564 #[tokio::test]
1569 async fn finish_resolves_an_unaccepted_track_with_fetched_info() {
1570 let producer = Info::new().produce();
1571 let consumer = producer.consume();
1572
1573 let request = producer.reserve_track("track1").unwrap();
1574 let dynamic = request.dynamic();
1575 let track = consumer.track("track1").unwrap();
1576 let pending_fetch = track.fetch_group(0, None);
1577 let fetch = dynamic.requested_group().await.unwrap();
1578 let group = fetch.accept(None).unwrap();
1579 group.finish().unwrap();
1580 pending_fetch.await.unwrap();
1581
1582 let mut subscriber = track.subscribe(None).await.unwrap();
1583 producer.finish();
1584 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1585
1586 let stale = request.accept(None);
1587 assert!(stale.append_group().is_err());
1588 }
1589
1590 #[tokio::test]
1593 async fn finish_spares_a_served_track() {
1594 let producer = Info::new().produce();
1595 let consumer = producer.consume();
1596
1597 let track = producer.create_track("track1", None).unwrap();
1598 let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1599
1600 producer.finish();
1601
1602 track.append_group().unwrap();
1603 subscriber.assert_group();
1604 track.finish().unwrap();
1605 }
1606
1607 #[tokio::test]
1611 async fn finish_leaves_a_stale_reservation_inert() {
1612 let producer = Info::new().produce();
1613 let consumer = producer.consume();
1614
1615 let request = producer.reserve_track("track1").unwrap();
1616 let pending = subscribe_pending!(consumer, "track1");
1617
1618 producer.finish();
1619 assert!(matches!(pending.await, Err(Error::NotFound)));
1620
1621 let track = request.accept(None);
1622 assert!(track.append_group().is_err());
1623 let mut subscriber = track.subscribe(None);
1624 assert!(matches!(subscriber.recv_group().await, Err(Error::NotFound)));
1625 assert!(consumer.track("track1").is_err());
1626 }
1627
1628 #[tokio::test]
1632 async fn dropping_a_reserved_request_resolves_dropped() {
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 drop(request);
1640 assert!(matches!(pending.await, Err(Error::Dropped)));
1641 producer.finish();
1642 }
1643
1644 #[tokio::test]
1647 async fn rejecting_a_reserved_request_carries_the_reason() {
1648 let producer = Info::new().produce();
1649 let consumer = producer.consume();
1650
1651 let request = producer.reserve_track("track1").unwrap();
1652 let pending = subscribe_pending!(consumer, "track1");
1653
1654 request.reject(Error::NotFound);
1655 assert!(matches!(pending.await, Err(Error::NotFound)));
1656 producer.finish();
1657 }
1658
1659 #[tokio::test]
1665 async fn an_idle_teardown_yields_to_a_returning_viewer() {
1666 let producer = Info::new().produce();
1667 let consumer = producer.consume();
1668 let track = producer.create_track("video", None).unwrap();
1669
1670 assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
1672
1673 let viewer = consumer.track("video").unwrap();
1675 let track = track
1676 .abort_unused(Error::Cancel)
1677 .expect_err("viewer keeps the track alive");
1678
1679 assert!(!track.is_closed());
1681 let mut subscriber = viewer.subscribe(None).await.unwrap();
1682 subscriber.assert_no_group();
1683 track.append_group().unwrap();
1684 assert!(subscriber.recv_group().await.unwrap().is_some());
1685
1686 drop(subscriber);
1689 drop(viewer);
1690 assert!(track.abort_unused(Error::Cancel).is_ok());
1691 assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
1692
1693 producer.finish();
1694 }
1695
1696 #[test]
1697 fn abort_unused_accepts_an_already_closed_track_with_consumers() {
1698 let producer = Info::new().produce();
1699 let consumer = producer.consume();
1700 let track = producer.create_track("video", None).unwrap();
1701 let _viewer = consumer.track("video").unwrap();
1702 assert!(track.is_used());
1703 track.clone().abort(Error::Cancel).unwrap();
1704 assert!(!track.is_used());
1705 assert!(track.abort_unused(Error::Cancel).is_ok());
1706 producer.finish();
1707 }
1708}