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,
103
104 abort: Option<Error>,
107}
108
109#[derive(Default)]
112struct SplicedState {
113 tracks: HashMap<Arc<str>, super::resume::Producer>,
116
117 pending: VecDeque<Arc<str>>,
119}
120
121impl BroadcastState {
122 fn insert_track(&mut self, weak: track::TrackWeak) -> Result<(), Error> {
125 match self.tracks.insert(weak.name().clone(), weak) {
126 Some(_) => Err(Error::Duplicate),
127 None => Ok(()),
128 }
129 }
130
131 fn reject_unserved(&mut self, err: Error) {
139 for request in self.requests.drain_queued() {
140 request.reject(err.clone());
141 }
142 for track in self.tracks.iter() {
143 track.reject(err.clone());
144 }
145 }
146
147 fn is_used(&self) -> bool {
150 if let Some(spliced) = &self.spliced {
151 return spliced.tracks.values().any(|track| track.is_used());
152 }
153 !self.requests.is_empty() || self.tracks.iter().any(|track| track.is_used())
154 }
155
156 fn register_demand(&self, waiter: &kio::Waiter, want: bool) {
161 if let Some(spliced) = &self.spliced {
162 for track in spliced.tracks.values() {
163 let _ = match want {
164 true => track.poll_used(waiter),
165 false => track.poll_unused(waiter),
166 };
167 }
168 return;
169 }
170 for track in self.tracks.iter() {
171 match want {
172 true => track.poll_used(waiter),
173 false => track.poll_unused(waiter),
174 }
175 }
176 }
177}
178
179#[derive(Clone)]
192pub struct Producer {
193 info: Arc<Info>,
196
197 alive: Arc<Alive>,
200
201 state: kio::Shared<BroadcastState>,
204
205 stats: stats::Scope,
209}
210
211impl Producer {
212 pub fn new(info: Info) -> Self {
214 let state = kio::Shared::<BroadcastState>::default();
215 Self {
216 info: Arc::new(info),
217 alive: Alive::new(state.clone()),
218 state,
219 stats: stats::Scope::default(),
220 }
221 }
222
223 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
226 self.stats = scope;
227 self
228 }
229
230 pub(crate) fn with_announcer(self, announcer: Announcer) -> Self {
233 *self.alive.announcer.lock() = Some(announcer);
234 self
235 }
236
237 pub fn announce(&self, route: Route) -> Result<(), Error> {
251 let mut announcer = self.alive.announcer.lock();
252 let announcer = announcer.as_mut().ok_or(Error::Closed)?;
253 announcer.announce(route)
254 }
255
256 pub fn unannounce(&self) {
262 self.alive.unannounce();
263 }
264
265 pub(crate) fn new_spliced(info: Info) -> Self {
269 let state = kio::Shared::new(BroadcastState {
270 spliced: Some(SplicedState::default()),
271 ..Default::default()
272 });
273 Self {
274 info: Arc::new(info),
275 alive: Alive::new(state.clone()),
276 state,
277 stats: stats::Scope::default(),
280 }
281 }
282
283 pub fn info(&self) -> &Info {
285 &self.info
286 }
287
288 pub fn demand(&self) -> Demand {
290 Demand {
291 alive: self.alive.token.consume().weak(),
292 state: self.state.clone(),
293 }
294 }
295
296 pub fn create_track(
301 &self,
302 name: impl Into<Arc<str>>,
303 info: impl Into<Option<track::Info>>,
304 ) -> Result<track::Producer, Error> {
305 let name = name.into();
306 let info = info.into().unwrap_or_default();
307 let mut state = self.state.lock();
308
309 if let Some(request) = state.requests.take(name.as_ref()) {
315 let track = request.with_stats(self.stats.clone()).accept(info);
316 let _ = state.tracks.insert(name, track.weak());
320 return Ok(track);
321 }
322
323 let track = track::Producer::new(self.info.clone(), name, info).with_stats(self.stats.clone());
324 state.insert_track(track.weak())?;
325 Ok(track)
326 }
327
328 pub fn reserve_track(&self, name: impl Into<Arc<str>>) -> Result<track::Request, Error> {
340 let request = track::Request::new(self.info.clone(), name).with_stats(self.stats.clone());
341 self.state.lock().insert_track(request.weak())?;
342 Ok(request)
343 }
344
345 pub fn unique_track(&self, suffix: &str, info: impl Into<Option<track::Info>>) -> Result<track::Producer, Error> {
349 let name = self.unique_name(suffix);
350 self.create_track(name, info)
351 }
352
353 pub fn unique_name(&self, suffix: &str) -> String {
365 let mut state = self.state.lock();
366 let separator = if suffix.starts_with(|c: char| c.is_ascii_digit()) {
367 "-"
368 } else {
369 ""
370 };
371 loop {
372 let id = state.unique;
373 state.unique = id.checked_add(1).expect("unique track IDs exhausted");
374 let name = format!("{id}{separator}{suffix}");
375 if !state.tracks.contains_key(name.as_str()) {
376 return name;
377 }
378 }
379 }
380
381 pub fn dynamic(&self) -> Dynamic {
383 Dynamic::new(
384 self.info.clone(),
385 self.alive.clone(),
386 self.state.clone(),
387 self.stats.clone(),
388 )
389 }
390
391 pub(crate) fn poll_spliced_assigned(&self, waiter: &kio::Waiter) -> Poll<(Arc<str>, super::resume::Producer)> {
394 let mut state = ready!(self.state.poll(waiter, |state| {
395 match &state.spliced {
396 Some(spliced) if !spliced.pending.is_empty() => Poll::Ready(()),
397 _ => Poll::Pending,
398 }
399 }));
400
401 let spliced = state.spliced.as_mut().expect("predicate guaranteed spliced");
402 let name = spliced.pending.pop_front().expect("predicate guaranteed a request");
403 let producer = spliced.tracks.get(&name).expect("pending name without a track").clone();
404 Poll::Ready((name, producer))
405 }
406
407 pub(crate) fn release_spliced(&self, err: Error) {
411 let mut state = self.state.lock();
412 if let Some(spliced) = state.spliced.as_mut() {
413 for name in std::mem::take(&mut spliced.pending) {
414 if let Some(producer) = spliced.tracks.get_mut(&name) {
415 let _ = producer.abort(err.clone());
416 }
417 }
418 spliced.tracks.clear();
419 }
420 }
421
422 pub(crate) fn forget_spliced(&self, name: &str, producer: &super::resume::Producer) -> bool {
429 let mut state = self.state.lock();
430 let Some(spliced) = state.spliced.as_mut() else {
431 return true;
432 };
433 match spliced.tracks.get(name) {
434 Some(current) if current.is_clone(producer) => {
435 if current.is_used() {
436 return false;
437 }
438 spliced.tracks.remove(name);
439 true
440 }
441 _ => true,
442 }
443 }
444
445 pub fn consume(&self) -> Consumer {
451 Consumer {
452 info: self.info.clone(),
453 alive: self.alive.token.consume(),
454 state: self.state.clone(),
455 stats: stats::Scope::default(),
456 }
457 }
458
459 pub fn close(&self) {
470 self.alive.close();
471 }
472
473 #[doc(hidden)]
474 #[deprecated(note = "use close(); a broadcast end carries no cause")]
475 pub fn finish(&self) {
476 self.alive.end(true);
477 }
478
479 #[doc(hidden)]
480 #[deprecated(note = "use close(); a broadcast end carries no cause")]
481 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 close(&self) {
536 self.end(false);
537 }
538
539 fn end(&self, finished: bool) {
542 {
543 let mut state = self.state.lock();
544 if std::mem::replace(&mut state.closing, true) {
545 return;
546 }
547 state.finished = finished;
548 state.reject_unserved(Error::Unroutable);
552 }
553 let _ = self.token.close();
554 self.retire();
555 }
556
557 fn retire(&self) {
560 let announcer = self.announcer.lock().take();
561 drop(announcer);
564 }
565}
566
567impl Drop for Alive {
568 fn drop(&mut self) {
569 self.close();
570 }
571}
572
573#[cfg(test)]
574#[allow(missing_docs)] impl Producer {
576 pub fn assert_create_track(
577 &mut self,
578 name: impl Into<Arc<str>>,
579 info: impl Into<Option<track::Info>>,
580 ) -> track::Producer {
581 self.create_track(name, info).expect("should not have errored")
582 }
583}
584
585pub(crate) struct SourceGuard(Producer);
592
593impl SourceGuard {
594 pub fn new(producer: Producer) -> Self {
595 Self(producer)
596 }
597}
598
599impl Drop for SourceGuard {
600 fn drop(&mut self) {
601 self.0.close();
602 }
603}
604
605#[derive(Clone)]
613pub struct Dynamic {
614 info: Arc<Info>,
615 alive: Arc<Alive>,
617 state: kio::Shared<BroadcastState>,
618 stats: stats::Scope,
621 _handler: Handler,
625}
626
627struct Handler(kio::Shared<BroadcastState>);
629
630impl Handler {
631 fn new(state: kio::Shared<BroadcastState>) -> Self {
632 state.lock().requests.add_handler();
633 Self(state)
634 }
635}
636
637impl Clone for Handler {
638 fn clone(&self) -> Self {
639 Self::new(self.0.clone())
642 }
643}
644
645impl Drop for Handler {
646 fn drop(&mut self) {
647 let mut state = self.0.lock();
650 if state.requests.remove_handler() {
651 for request in state.requests.drain_queued() {
654 request.reject(Error::Dropped);
655 }
656 }
657 }
658}
659
660impl Dynamic {
661 fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
662 Self {
663 info,
664 alive,
665 _handler: Handler::new(state.clone()),
666 state,
667 stats,
668 }
669 }
670
671 pub fn info(&self) -> &Info {
673 &self.info
674 }
675
676 pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
681 let mut state = ready!(self.state.poll(waiter, |state| {
682 if state.requests.has_queued() || state.closing {
683 Poll::Ready(())
684 } else {
685 Poll::Pending
686 }
687 }));
688
689 if state.closing && !state.requests.has_queued() {
690 return Poll::Ready(Err(Error::Closed));
691 }
692
693 let name = state.requests.pop().expect("predicate guaranteed a request");
694 let pending = state.requests.remove(&name).expect("popped key must be pending");
695 let _ = state.tracks.insert(name, pending.weak());
698 Poll::Ready(Ok(pending.claim().with_stats(self.stats.clone())))
700 }
701
702 pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
704 kio::wait(|waiter| self.poll_requested_track(waiter)).await
705 }
706
707 pub fn consume(&self) -> Consumer {
709 Consumer {
710 info: self.info.clone(),
711 alive: self.alive.token.consume(),
712 state: self.state.clone(),
713 stats: stats::Scope::default(),
714 }
715 }
716
717 pub async fn closed(&self) -> Error {
721 kio::wait(|waiter| self.poll_closed(waiter)).await
722 }
723
724 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
726 ready!(self.alive.token.poll_closed(waiter));
727 Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
728 }
729
730 pub fn is_clone(&self, other: &Self) -> bool {
732 self.state.same_channel(&other.state)
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 {
760 info: Arc<Info>,
761 alive: kio::Consumer<()>,
763 state: kio::Shared<BroadcastState>,
765 stats: stats::Scope,
769}
770
771impl Clone for Consumer {
772 fn clone(&self) -> Self {
773 Self {
774 info: self.info.clone(),
775 alive: self.alive.clone(),
776 state: self.state.clone(),
777 stats: self.stats.clone(),
778 }
779 }
780}
781
782impl Consumer {
783 pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
786 self.stats = scope;
787 self
788 }
789
790 pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
798 if self.info.path != path {
799 let mut info = (*self.info).clone();
800 info.path = path;
801 self.info = Arc::new(info);
802 }
803 self
804 }
805
806 pub fn info(&self) -> &Info {
808 &self.info
809 }
810
811 pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
815 self.track_inner(name)
820 .map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
821 }
822
823 fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
824 let mut state = self.state.lock();
825
826 if state.closing {
830 return Err(Error::Unroutable);
831 }
832
833 if let Some(spliced) = state.spliced.as_mut() {
836 if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
851 spliced.tracks.remove(name);
852 }
853 if let Some(producer) = spliced.tracks.get(name) {
854 return Ok(track::Consumer::spliced(
855 name.into(),
856 self.info.clone(),
857 producer.consume(),
858 ));
859 }
860 let name: Arc<str> = name.into();
861 let producer = super::resume::Producer::new();
862 let consumer = producer.consume();
863 spliced.tracks.insert(name.clone(), producer);
864 spliced.pending.push_back(name.clone());
865 return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
866 }
867
868 if let Some(weak) = state.tracks.get(name) {
871 match weak.try_consume() {
872 Some(consumer) => return Ok(consumer),
873 None => {
877 state.tracks.remove(name);
878 }
879 }
880 }
881
882 if let Some(pending) = state.requests.join(name) {
883 return Ok(pending.consume());
885 }
886
887 let name: Arc<str> = name.into();
891 let request = track::Request::new(self.info.clone(), name.clone());
892 let consumer = request.consume();
893
894 if state.requests.insert(name, request).is_err() {
897 return Err(Error::NotFound);
898 }
899
900 Ok(consumer)
901 }
902
903 pub fn demand(&self) -> Demand {
912 Demand {
913 alive: self.alive.weak(),
914 state: self.state.clone(),
915 }
916 }
917
918 pub async fn closed(&self) -> Error {
922 self.alive.closed().await;
923 self.state.read().abort.clone().unwrap_or(Error::Dropped)
924 }
925
926 pub fn is_closed(&self) -> bool {
928 self.alive.is_closed()
929 }
930
931 pub(crate) fn is_closing(&self) -> bool {
935 self.state.read().closing
936 }
937
938 #[doc(hidden)]
939 #[deprecated(note = "a broadcast end carries no cause")]
940 pub fn is_finished(&self) -> bool {
941 self.state.read().finished
942 }
943
944 pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
950 self.alive.poll_closed(waiter)
951 }
952
953 pub fn is_clone(&self, other: &Self) -> bool {
955 self.state.same_channel(&other.state)
956 }
957
958 pub(crate) fn weak(&self) -> WeakConsumer {
963 WeakConsumer {
964 info: self.info.clone(),
965 alive: self.alive.weak(),
966 state: self.state.clone(),
967 }
968 }
969}
970
971#[derive(Clone)]
978pub(crate) struct WeakConsumer {
979 info: Arc<Info>,
980 alive: kio::ConsumerWeak<()>,
981 state: kio::Shared<BroadcastState>,
982}
983
984impl WeakConsumer {
985 pub fn consume(&self) -> Consumer {
987 Consumer {
988 info: self.info.clone(),
989 alive: self.alive.consume(),
990 state: self.state.clone(),
991 stats: stats::Scope::default(),
992 }
993 }
994}
995
996impl super::WeakEntry for WeakConsumer {
997 fn is_closed(&self) -> bool {
998 self.alive.is_closed()
999 }
1000
1001 fn same_channel(&self, other: &Self) -> bool {
1002 self.state.same_channel(&other.state)
1003 }
1004}
1005
1006#[derive(Clone)]
1019pub struct Demand {
1020 alive: kio::ConsumerWeak<()>,
1021 state: kio::Shared<BroadcastState>,
1022}
1023
1024impl Demand {
1025 pub fn is_used(&self) -> bool {
1030 self.state.read().is_used()
1031 }
1032
1033 pub async fn used(&self) -> Result<(), Error> {
1036 kio::wait(|waiter| self.poll_used(waiter)).await
1037 }
1038
1039 pub async fn unused(&self) -> Result<(), Error> {
1042 kio::wait(|waiter| self.poll_unused(waiter)).await
1043 }
1044
1045 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1047 self.poll_demand(waiter, true)
1048 }
1049
1050 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1052 self.poll_demand(waiter, false)
1053 }
1054
1055 fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
1056 if self.alive.poll_closed(waiter).is_ready() {
1059 return Poll::Ready(Err(Error::Dropped));
1060 }
1061 let ready = self.state.poll(waiter, |state| {
1062 state.register_demand(waiter, want);
1066 match state.is_used() == want {
1067 true => Poll::Ready(()),
1068 false => Poll::Pending,
1069 }
1070 });
1071 match ready {
1072 Poll::Ready(_) => Poll::Ready(Ok(())),
1073 Poll::Pending => Poll::Pending,
1074 }
1075 }
1076}
1077
1078#[cfg(test)]
1079#[allow(missing_docs)] impl Consumer {
1081 pub fn assert_not_closed(&self) {
1082 assert!(self.closed().now_or_never().is_none(), "should not be closed");
1083 }
1084
1085 pub fn assert_closed(&self) {
1086 assert!(self.closed().now_or_never().is_some(), "should be closed");
1087 }
1088}
1089
1090#[cfg(test)]
1091mod test {
1092 use super::*;
1093 use std::time::Duration;
1094
1095 #[test]
1096 fn unique_names_are_never_reused() {
1097 let producer = Info::new().produce();
1098 let name = producer.unique_name(".opus");
1099 assert_eq!(name, "0.opus");
1100 let track = producer.create_track(name.clone(), None).unwrap();
1101 assert_eq!(producer.unique_name(".opus"), "1.opus");
1102 drop(track);
1103 }
1104
1105 #[test]
1106 fn unique_names_survive_closed_track_pruning() {
1107 let producer = Info::new().produce();
1108 let consumer = producer.consume();
1109 let track = producer.unique_track(".opus", None).unwrap();
1110 assert_eq!(track.name(), "0.opus");
1111 drop(track);
1112 assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
1113 assert_eq!(producer.unique_name(".opus"), "1.opus");
1114 }
1115
1116 #[test]
1117 fn unique_name_skips_a_live_collision() {
1118 let producer = Info::new().produce();
1119 let track = producer.create_track("0.opus", None).unwrap();
1120 assert_eq!(producer.unique_name(".opus"), "1.opus");
1121 drop(track);
1122 assert_eq!(producer.unique_name(".opus"), "2.opus");
1123 }
1124
1125 #[test]
1126 fn unique_names_share_a_counter() {
1127 let producer = Info::new().produce();
1128 assert_eq!(producer.unique_name("-video"), "0-video");
1129 assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
1130 assert_eq!(producer.unique_name("-video"), "2-video");
1131 }
1132
1133 #[test]
1134 fn unique_names_separate_numeric_suffixes() {
1135 let producer = Info::new().produce();
1136 assert_eq!(producer.unique_name(""), "0");
1137 let name = producer.unique_name("2");
1138 assert_eq!(name, "1-2");
1139 for _ in 2..12 {
1140 producer.unique_name("");
1141 }
1142 assert_eq!(producer.unique_name(""), "12");
1143 }
1144
1145 async fn expect<T>(fut: impl Future<Output = T>) -> T {
1148 tokio::time::timeout(Duration::from_secs(1), fut)
1149 .await
1150 .expect("timed out waiting for a demand edge")
1151 }
1152
1153 #[tokio::test]
1157 async fn demand_ordinary() {
1158 tokio::time::pause();
1159
1160 let producer = Info::new().produce();
1161 let consumer = producer.consume();
1162 let demand = producer.demand();
1163
1164 assert!(!demand.is_used());
1166 demand.unused().await.unwrap();
1167
1168 let _track = producer.create_track("a", None).unwrap();
1170 assert!(!demand.is_used());
1171
1172 let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
1174 used.unwrap();
1175 assert!(demand.is_used());
1176
1177 let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
1179 unused.unwrap();
1180 assert!(!demand.is_used());
1181
1182 producer.close();
1184 assert!(matches!(demand.used().await, Err(Error::Dropped)));
1185 assert!(matches!(demand.unused().await, Err(Error::Dropped)));
1186 }
1187
1188 #[tokio::test]
1191 async fn demand_spliced() {
1192 tokio::time::pause();
1193
1194 let producer = Producer::new_spliced(Info::new());
1195 let consumer = producer.consume();
1196 let demand = producer.demand();
1197 let watched = consumer.demand();
1198
1199 assert!(!demand.is_used());
1200 assert!(!watched.is_used());
1201 let track = consumer.track("video").unwrap();
1202 assert!(demand.is_used());
1203 assert!(watched.is_used());
1204
1205 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1208 unused.unwrap();
1209 assert!(!demand.is_used());
1210 assert!(!watched.is_used());
1211
1212 let _track = consumer.track("video").unwrap();
1214 assert!(demand.is_used());
1215 }
1216
1217 #[tokio::test]
1219 async fn consumer_demand_reports_dropped_producer() {
1220 let producer = Producer::new_spliced(Info::new());
1221 let consumer = producer.consume();
1222 let watched = consumer.demand();
1223
1224 let track = consumer.track("video").unwrap();
1225 assert!(watched.is_used());
1226
1227 let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1228 unused.unwrap();
1229
1230 drop(producer);
1231 assert!(matches!(watched.used().await, Err(Error::Dropped)));
1232 assert!(matches!(watched.unused().await, Err(Error::Dropped)));
1233 }
1234
1235 macro_rules! subscribe_pending {
1238 ($consumer:expr, $name:expr) => {{
1239 let pending = $consumer.track($name).unwrap().subscribe(None);
1240 assert!(
1241 pending.poll_ok(&kio::Waiter::noop()).is_pending(),
1242 "subscribe should stay pending until the request is accepted"
1243 );
1244 pending
1245 }};
1246 }
1247
1248 #[tokio::test]
1249 async fn insert() {
1250 let mut producer = Info::new().produce();
1251
1252 let track1 = producer.assert_create_track("track1", None);
1254 track1.append_group().unwrap();
1255
1256 let consumer = producer.consume();
1257
1258 let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1260 track1_sub.assert_group();
1261
1262 let track2 = producer.assert_create_track("track2", None);
1263
1264 let consumer2 = producer.consume();
1265 let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
1266 track2_consumer.assert_no_group();
1267
1268 track2.append_group().unwrap();
1269
1270 track2_consumer.assert_group();
1271 }
1272
1273 #[tokio::test]
1274 async fn closed() {
1275 let mut producer = Info::new().produce();
1276 let dynamic = producer.dynamic();
1277
1278 let consumer = producer.consume();
1279 consumer.assert_not_closed();
1280
1281 let track1 = producer.assert_create_track("track1", None);
1283 let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1284
1285 let track2_fut = subscribe_pending!(consumer, "track2");
1287
1288 drop(dynamic);
1291
1292 assert!(track2_fut.await.is_err());
1294
1295 assert!(!track1.is_closed());
1297 track1c.assert_not_closed();
1298 }
1299
1300 #[tokio::test]
1302 async fn close_ends_every_clone() {
1303 let producer = Info::new().produce();
1304 let clone = producer.clone();
1305 let consumer = producer.consume();
1306
1307 producer.close();
1308 assert!(matches!(consumer.closed().await, Error::Dropped));
1309 assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
1310 assert!(matches!(clone.consume().track("video"), Err(Error::Unroutable)));
1311
1312 clone.close();
1313 producer.close();
1314 }
1315
1316 #[tokio::test]
1318 async fn drop_ends_like_close() {
1319 let producer = Info::new().produce();
1320 let consumer = producer.consume();
1321 drop(producer);
1322 assert!(matches!(consumer.closed().await, Error::Dropped));
1323 assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
1324 }
1325
1326 #[tokio::test]
1328 #[allow(deprecated)]
1329 async fn deprecated_end_causes() {
1330 let producer = Info::new().produce();
1331 let consumer = producer.consume();
1332 producer.abort(Error::Timeout).unwrap();
1333 assert!(matches!(consumer.closed().await, Error::Timeout));
1334 assert!(!consumer.is_finished());
1335
1336 let producer = Info::new().produce();
1337 let consumer = producer.consume();
1338 producer.finish();
1339 assert!(matches!(consumer.closed().await, Error::Dropped));
1340 assert!(consumer.is_finished());
1341 }
1342
1343 #[tokio::test]
1344 async fn requests() {
1345 let mut producer = Info::new().produce().dynamic();
1346
1347 let consumer = producer.consume();
1348 let consumer2 = consumer.clone();
1349
1350 let track1_fut = subscribe_pending!(consumer, "track1");
1352 let track2_fut = subscribe_pending!(consumer2, "track1");
1353
1354 let request = producer.assert_request();
1356 producer.assert_no_request();
1357 assert_eq!(request.name(), "track1");
1358
1359 let track3 = request.accept(None);
1361 let mut track1 = track1_fut.await.unwrap();
1362 let mut track2 = track2_fut.await.unwrap();
1363
1364 track1.assert_not_closed();
1365 track1.assert_is_clone(&track2);
1366 track3.subscribe(None).assert_is_clone(&track1);
1367
1368 track3.append_group().unwrap();
1370 track1.assert_group();
1371 track2.assert_group();
1372
1373 let track4_fut = subscribe_pending!(consumer, "track2");
1375 drop(producer);
1376 assert!(track4_fut.await.is_err());
1377
1378 let track5 = consumer2.track("track3");
1380 assert!(track5.is_err(), "should have errored");
1381 }
1382
1383 #[tokio::test]
1384 async fn stale_producer() {
1385 let mut broadcast = Info::new().produce().dynamic();
1386 let consumer = broadcast.consume();
1387
1388 let track1_fut = subscribe_pending!(consumer, "track1");
1390 let producer1 = broadcast.assert_request().accept(None);
1391 let mut track1 = track1_fut.await.unwrap();
1392
1393 producer1.append_group().unwrap();
1395 producer1.finish().unwrap();
1396 drop(producer1);
1397
1398 track1.assert_closed();
1400
1401 let track2_fut = subscribe_pending!(consumer, "track1");
1403 let producer2 = broadcast.assert_request().accept(None);
1404 let mut track2 = track2_fut.await.unwrap();
1405 track2.assert_not_closed();
1406 track2.assert_not_clone(&track1);
1407
1408 producer2.append_group().unwrap();
1410 track2.assert_group();
1411 }
1412
1413 #[tokio::test(start_paused = true)]
1414 async fn requested_unused() {
1415 let mut broadcast = Info::new().produce().dynamic();
1416 let bc = broadcast.consume();
1417
1418 let c1_fut = subscribe_pending!(bc, "unknown_track");
1420 let producer1 = broadcast.assert_request().accept(None);
1421 let consumer1 = c1_fut.await.unwrap();
1422
1423 assert!(
1425 producer1.unused().now_or_never().is_none(),
1426 "track producer should be used"
1427 );
1428
1429 let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1431 consumer2.assert_is_clone(&consumer1);
1432
1433 drop(consumer1);
1434 assert!(
1435 producer1.unused().now_or_never().is_none(),
1436 "track producer should be used"
1437 );
1438
1439 drop(consumer2);
1440 assert!(
1441 producer1.unused().now_or_never().is_some(),
1442 "track producer should be unused after all consumers are dropped"
1443 );
1444
1445 let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1449 consumer3.assert_is_clone(&producer1.subscribe(None));
1450 broadcast.assert_no_request();
1451 drop(consumer3);
1452
1453 producer1.abort(Error::Cancel).unwrap();
1456
1457 let c4_fut = subscribe_pending!(bc, "unknown_track");
1458 let producer2 = broadcast.assert_request().accept(None);
1459 let consumer4 = c4_fut.await.unwrap();
1460 drop(consumer4);
1461 assert!(
1462 producer2.unused().now_or_never().is_some(),
1463 "new track producer should be unused after its consumer is dropped"
1464 );
1465 }
1466
1467 #[tokio::test]
1473 async fn create_track_fulfills_queued_request() {
1474 let producer = Info::new().produce();
1475 let mut dynamic = producer.dynamic();
1476 let bc = dynamic.consume();
1477
1478 let subscribing = subscribe_pending!(bc, "video");
1480
1481 let track = producer.create_track("video", None).unwrap();
1483 let mut sub = subscribing.await.expect("fulfilled by create_track");
1484
1485 track.append_group().unwrap();
1487 sub.recv_group().await.expect("recv").expect("group");
1488
1489 dynamic.assert_no_request();
1491 let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
1492 again.assert_is_clone(&track.subscribe(None));
1493 }
1494
1495 #[tokio::test]
1501 async fn dynamic_clone_keeps_alive() {
1502 let broadcast = Info::new().produce().dynamic();
1503 let consumer = broadcast.consume();
1504
1505 let clone = broadcast.clone();
1506 drop(clone);
1507
1508 let _fut = subscribe_pending!(consumer, "track1");
1511 }
1512
1513 #[tokio::test]
1518 async fn close_resolves_a_reserved_name() {
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.close();
1526 assert!(matches!(pending.await, Err(Error::Unroutable)));
1527 }
1528
1529 #[tokio::test]
1532 #[allow(deprecated)]
1533 async fn abort_resolves_a_reserved_name_with_its_reason() {
1534 let producer = Info::new().produce();
1535 let consumer = producer.consume();
1536
1537 let request = producer.reserve_track("track1").unwrap();
1538 let pending = subscribe_pending!(consumer, "track1");
1539
1540 producer.abort(Error::Cancel).unwrap();
1541 assert!(matches!(pending.await, Err(Error::Cancel)));
1542
1543 let track = request.accept(None);
1544 let mut subscriber = track.subscribe(None);
1545 assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
1546 }
1547
1548 #[tokio::test]
1551 async fn close_resolves_a_queued_request() {
1552 let producer = Info::new().produce();
1553 let dynamic = producer.dynamic();
1554 let consumer = dynamic.consume();
1555
1556 let pending = subscribe_pending!(consumer, "track1");
1557
1558 producer.close();
1559 assert!(matches!(pending.await, Err(Error::Unroutable)));
1560 drop(dynamic);
1561 }
1562
1563 #[tokio::test]
1566 async fn dropping_the_last_handle_resolves_a_queued_request() {
1567 let dynamic = Info::new().produce().dynamic();
1568 let consumer = dynamic.consume();
1569
1570 let pending = subscribe_pending!(consumer, "track1");
1571
1572 drop(dynamic);
1573 assert!(matches!(pending.await, Err(Error::Unroutable)));
1574 }
1575
1576 #[tokio::test]
1579 async fn dropping_the_last_handler_resolves_a_queued_request_dropped() {
1580 let producer = Info::new().produce();
1581 let dynamic = producer.dynamic();
1582 let consumer = dynamic.consume();
1583
1584 let pending = subscribe_pending!(consumer, "track1");
1585
1586 drop(dynamic);
1587 assert!(matches!(pending.await, Err(Error::Dropped)));
1588 producer.close();
1589 }
1590
1591 #[tokio::test]
1595 async fn close_leaves_a_claimed_request_to_its_handler() {
1596 let producer = Info::new().produce();
1597 let mut dynamic = producer.dynamic();
1598 let consumer = dynamic.consume();
1599
1600 let accepted = subscribe_pending!(consumer, "track1");
1601 let request = dynamic.requested_track().await.unwrap();
1602 let dropped = subscribe_pending!(consumer, "track2");
1603 let abandoned = dynamic.requested_track().await.unwrap();
1604
1605 producer.close();
1606 assert!(
1607 accepted.poll_ok(&kio::Waiter::noop()).is_pending(),
1608 "close rejected a claimed request"
1609 );
1610 assert!(
1611 dropped.poll_ok(&kio::Waiter::noop()).is_pending(),
1612 "close rejected a claimed request"
1613 );
1614
1615 let _track = request.accept(None);
1616 assert!(accepted.await.is_ok(), "the handler's accept reaches the consumer");
1617 drop(abandoned);
1618 assert!(dropped.await.is_err(), "the handler dropping it rejects the consumer");
1619 drop(dynamic);
1620 }
1621
1622 #[tokio::test]
1627 async fn close_resolves_an_unaccepted_track_with_fetched_info() {
1628 let producer = Info::new().produce();
1629 let consumer = producer.consume();
1630
1631 let request = producer.reserve_track("track1").unwrap();
1632 let dynamic = request.dynamic();
1633 let track = consumer.track("track1").unwrap();
1634 let pending_fetch = track.fetch_group(0, None);
1635 let fetch = dynamic.requested_group().await.unwrap();
1636 let group = fetch.accept(None).unwrap();
1637 group.finish().unwrap();
1638 pending_fetch.await.unwrap();
1639
1640 let mut subscriber = track.subscribe(None).await.unwrap();
1641 producer.close();
1642 assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
1643
1644 let stale = request.accept(None);
1645 assert!(stale.append_group().is_err());
1646 }
1647
1648 #[tokio::test]
1651 async fn close_spares_a_served_track() {
1652 let producer = Info::new().produce();
1653 let consumer = producer.consume();
1654
1655 let track = producer.create_track("track1", None).unwrap();
1656 let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1657
1658 producer.close();
1659 assert!(matches!(consumer.track("track1"), Err(Error::Unroutable)));
1660
1661 track.append_group().unwrap();
1662 subscriber.assert_group();
1663 track.finish().unwrap();
1664 }
1665
1666 #[tokio::test]
1670 async fn close_leaves_a_stale_reservation_inert() {
1671 let producer = Info::new().produce();
1672 let consumer = producer.consume();
1673
1674 let request = producer.reserve_track("track1").unwrap();
1675 let pending = subscribe_pending!(consumer, "track1");
1676
1677 producer.close();
1678 assert!(matches!(pending.await, Err(Error::Unroutable)));
1679
1680 let track = request.accept(None);
1681 assert!(track.append_group().is_err());
1682 let mut subscriber = track.subscribe(None);
1683 assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
1684 assert!(consumer.track("track1").is_err());
1685 }
1686
1687 #[tokio::test]
1691 async fn dropping_a_reserved_request_resolves_dropped() {
1692 let producer = Info::new().produce();
1693 let consumer = producer.consume();
1694
1695 let request = producer.reserve_track("track1").unwrap();
1696 let pending = subscribe_pending!(consumer, "track1");
1697
1698 drop(request);
1699 assert!(matches!(pending.await, Err(Error::Dropped)));
1700 producer.close();
1701 }
1702
1703 #[tokio::test]
1706 async fn rejecting_a_reserved_request_carries_the_reason() {
1707 let producer = Info::new().produce();
1708 let consumer = producer.consume();
1709
1710 let request = producer.reserve_track("track1").unwrap();
1711 let pending = subscribe_pending!(consumer, "track1");
1712
1713 request.reject(Error::NotFound);
1714 assert!(matches!(pending.await, Err(Error::NotFound)));
1715 producer.close();
1716 }
1717
1718 #[tokio::test]
1724 async fn an_idle_teardown_yields_to_a_returning_viewer() {
1725 let producer = Info::new().produce();
1726 let consumer = producer.consume();
1727 let track = producer.create_track("video", None).unwrap();
1728
1729 assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
1731
1732 let viewer = consumer.track("video").unwrap();
1734 let track = track
1735 .abort_unused(Error::Cancel)
1736 .expect_err("viewer keeps the track alive");
1737
1738 assert!(!track.is_closed());
1740 let mut subscriber = viewer.subscribe(None).await.unwrap();
1741 subscriber.assert_no_group();
1742 track.append_group().unwrap();
1743 assert!(subscriber.recv_group().await.unwrap().is_some());
1744
1745 drop(subscriber);
1748 drop(viewer);
1749 assert!(track.abort_unused(Error::Cancel).is_ok());
1750 assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
1751
1752 producer.close();
1753 }
1754
1755 #[test]
1756 fn abort_unused_accepts_an_already_closed_track_with_consumers() {
1757 let producer = Info::new().produce();
1758 let consumer = producer.consume();
1759 let track = producer.create_track("video", None).unwrap();
1760 let _viewer = consumer.track("video").unwrap();
1761 assert!(track.is_used());
1762 track.clone().abort(Error::Cancel).unwrap();
1763 assert!(!track.is_used());
1764 assert!(track.abort_unused(Error::Cancel).is_ok());
1765 producer.close();
1766 }
1767}