1use crate::{broadcast, cache, stats, track};
2use kio::Pollable;
3use std::{
4 cmp::Reverse,
5 collections::{BTreeMap, HashMap, HashSet},
6 fmt,
7 sync::Arc,
8 sync::atomic::{AtomicU64, Ordering},
9 task::{Poll, ready},
10 time::Duration,
11};
12
13use rand::RngExt;
14use web_async::Lock;
15
16use super::{Requests, WeakCache};
17use crate::{
18 AsPath, Error, Path, PathOwned, PathPrefixes,
19 coding::{BoundsExceeded, Decode, DecodeError, Encode, EncodeError},
20};
21
22#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
28pub struct Origin {
29 id: u64,
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35#[non_exhaustive]
36pub struct InvalidOrigin;
37
38impl fmt::Display for InvalidOrigin {
39 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
40 write!(f, "local origin id must be non-zero and below 2^62")
41 }
42}
43
44impl std::error::Error for InvalidOrigin {}
45
46impl Origin {
47 pub(crate) const UNKNOWN: Self = Self { id: 0 };
50
51 pub fn new(id: u64) -> Result<Self, InvalidOrigin> {
57 if id == 0 || id >= 1u64 << 62 {
58 return Err(InvalidOrigin);
59 }
60 Ok(Self { id })
61 }
62
63 pub fn random() -> Self {
72 let mut rng = rand::rng();
73 let id = rng.random_range(1..(1u64 << 53));
74 Self { id }
75 }
76
77 pub fn id(self) -> u64 {
79 self.id
80 }
81
82 pub fn produce(self) -> Producer {
85 Info::new(self).produce()
86 }
87}
88
89#[derive(Clone, Debug)]
99#[non_exhaustive]
100pub struct Info {
101 pub id: Origin,
104
105 pub pool: cache::Pool,
111
112 pub cache_duration: Duration,
119
120 pub linger: Duration,
129}
130
131impl Default for Info {
132 fn default() -> Self {
135 Self {
136 id: Origin::UNKNOWN,
137 pool: cache::Pool::default(),
138 cache_duration: Duration::MAX,
139 linger: Duration::ZERO,
140 }
141 }
142}
143
144impl Info {
145 pub fn new(id: Origin) -> Self {
147 Self { id, ..Self::default() }
148 }
149
150 pub fn with_pool(mut self, pool: cache::Pool) -> Self {
152 self.pool = pool;
153 self
154 }
155
156 pub fn with_cache_duration(mut self, cache_duration: Duration) -> Self {
159 self.cache_duration = cache_duration;
160 self
161 }
162
163 pub fn with_linger(mut self, linger: Duration) -> Self {
166 self.linger = linger;
167 self
168 }
169
170 pub fn produce(self) -> Producer {
172 Producer::new(self)
173 }
174}
175
176impl TryFrom<u64> for Origin {
177 type Error = InvalidOrigin;
178
179 fn try_from(id: u64) -> Result<Self, Self::Error> {
180 Self::new(id)
181 }
182}
183
184impl fmt::Display for Origin {
185 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
186 self.id.fmt(f)
187 }
188}
189
190impl<V: Copy> Encode<V> for Origin
191where
192 u64: Encode<V>,
193{
194 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
195 self.id.encode(w, version)
196 }
197}
198
199impl<V: Copy> Decode<V> for Origin
200where
201 u64: Decode<V>,
202{
203 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
204 let id = u64::decode(r, version)?;
205 if id >= 1u64 << 62 {
206 return Err(DecodeError::InvalidValue);
207 }
208 Ok(Self { id })
209 }
210}
211
212pub(crate) const MAX_HOPS: usize = 32;
218
219#[derive(Debug, Clone, Default, PartialEq, Eq)]
224pub struct OriginList(Vec<Origin>);
225
226#[derive(Debug, Clone, Copy, PartialEq, Eq)]
228#[non_exhaustive]
229pub struct TooManyOrigins;
230
231impl fmt::Display for TooManyOrigins {
232 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
233 write!(f, "too many origins (max {MAX_HOPS})")
234 }
235}
236
237impl std::error::Error for TooManyOrigins {}
238
239impl From<TooManyOrigins> for DecodeError {
240 fn from(_: TooManyOrigins) -> Self {
241 DecodeError::BoundsExceeded
242 }
243}
244
245impl OriginList {
246 pub fn new() -> Self {
248 Self(Vec::new())
249 }
250
251 pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
253 if self.0.len() >= MAX_HOPS {
254 return Err(TooManyOrigins);
255 }
256 self.0.push(origin);
257 Ok(())
258 }
259
260 pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
263 for entry in &mut self.0 {
264 if *entry == target {
265 *entry = replacement;
266 return true;
267 }
268 }
269 false
270 }
271
272 pub fn contains(&self, origin: &Origin) -> bool {
274 self.0.contains(origin)
275 }
276
277 pub fn len(&self) -> usize {
279 self.0.len()
280 }
281
282 pub fn is_empty(&self) -> bool {
284 self.0.is_empty()
285 }
286
287 pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
289 self.0.iter()
290 }
291
292 pub fn as_slice(&self) -> &[Origin] {
294 &self.0
295 }
296}
297
298impl TryFrom<Vec<Origin>> for OriginList {
299 type Error = TooManyOrigins;
300
301 fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
302 if v.len() > MAX_HOPS {
303 return Err(TooManyOrigins);
304 }
305 Ok(Self(v))
306 }
307}
308
309impl<'a> IntoIterator for &'a OriginList {
310 type Item = &'a Origin;
311 type IntoIter = std::slice::Iter<'a, Origin>;
312
313 fn into_iter(self) -> Self::IntoIter {
314 self.iter()
315 }
316}
317
318impl<V: Copy> Encode<V> for OriginList
319where
320 u64: Encode<V>,
321 Origin: Encode<V>,
322{
323 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
324 (self.0.len() as u64).encode(w, version)?;
325 for origin in &self.0 {
326 origin.encode(w, version)?;
327 }
328 Ok(())
329 }
330}
331
332impl<V: Copy> Decode<V> for OriginList
333where
334 u64: Decode<V>,
335 Origin: Decode<V>,
336{
337 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
338 let count = u64::decode(r, version)? as usize;
339 if count > MAX_HOPS {
340 return Err(DecodeError::BoundsExceeded);
341 }
342 let mut list = Vec::with_capacity(count);
343 for _ in 0..count {
344 list.push(Origin::decode(r, version)?);
345 }
346 Ok(Self(list))
347 }
348}
349
350static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
351
352#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
353struct ConsumerId(u64);
354
355impl ConsumerId {
356 fn new() -> Self {
357 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
358 }
359}
360
361struct OriginBroadcast {
367 path: PathOwned,
368 broadcast: broadcast::Producer,
370 state: kio::Producer<FrontState>,
373 announced: bool,
374}
375
376fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
385 (hops.len(), fnv_key(name, hops.iter().copied()))
386}
387
388fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
401 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
403
404 let mut hash = SEED;
405 for &byte in name.as_str().as_bytes() {
406 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
407 }
408 for origin in origins {
409 for &byte in &origin.id().to_le_bytes() {
410 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
411 }
412 }
413
414 hash
415}
416
417fn route_order(name: &Path, route: &FrontRoute) -> (bool, u64, usize, u64, Reverse<u64>) {
436 let (len, hash) = route_key(name, &route.route.hops);
437 (!route.route.announce, route.route.cost, len, hash, Reverse(route.id))
438}
439
440enum PendingUpdate {
449 Announce(broadcast::Consumer),
450 Unannounce,
451 UnannounceAnnounce(broadcast::Consumer),
452}
453
454#[derive(Default)]
459struct OriginConsumerState {
460 pending: BTreeMap<PathOwned, PendingUpdate>,
461}
462
463impl OriginConsumerState {
464 fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
465 let new = match self.pending.remove(&path) {
466 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
468 Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
470 PendingUpdate::UnannounceAnnounce(broadcast)
471 }
472 };
473 self.pending.insert(path, new);
474 }
475
476 fn apply_unannounce(&mut self, path: PathOwned) {
477 match self.pending.remove(&path) {
478 Some(PendingUpdate::Announce(_)) => {}
480 None | Some(PendingUpdate::Unannounce) => {
481 self.pending.insert(path, PendingUpdate::Unannounce);
482 }
483 Some(PendingUpdate::UnannounceAnnounce(_)) => {
486 self.pending.insert(path, PendingUpdate::Unannounce);
487 }
488 }
489 }
490
491 fn take(&mut self) -> Option<OriginAnnounce> {
493 let path = self.pending.keys().next()?.clone();
494 let broadcast = match self.pending.remove(&path).unwrap() {
495 PendingUpdate::Announce(broadcast) => Some(broadcast),
496 PendingUpdate::Unannounce => None,
497 PendingUpdate::UnannounceAnnounce(broadcast) => {
498 self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
501 None
502 }
503 };
504 Some(OriginAnnounce { path, broadcast })
505 }
506}
507
508#[derive(Clone)]
509struct AnnounceConsumerNotify {
510 root: PathOwned,
511 state: kio::Producer<OriginConsumerState>,
512}
513
514impl AnnounceConsumerNotify {
515 fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer) {
516 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
517 self.state
518 .write()
519 .ok()
520 .expect("consumer closed")
521 .apply_announce(path, broadcast);
522 }
523
524 fn unannounce(&self, path: impl AsPath) {
525 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
526 self.state.write().ok().expect("consumer closed").apply_unannounce(path);
527 }
528}
529
530struct NotifyNode {
531 parent: Option<Lock<NotifyNode>>,
532
533 consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
536}
537
538impl NotifyNode {
539 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
540 Self {
541 parent,
542 consumers: HashMap::new(),
543 }
544 }
545
546 fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer) {
547 for consumer in self.consumers.values() {
548 consumer.announce(path.as_path(), broadcast.clone());
549 }
550
551 if let Some(parent) = &self.parent {
552 parent.lock().announce(path, broadcast);
553 }
554 }
555
556 fn unannounce(&mut self, path: impl AsPath) {
557 for consumer in self.consumers.values() {
558 consumer.unannounce(path.as_path());
559 }
560
561 if let Some(parent) = &self.parent {
562 parent.lock().unannounce(path);
563 }
564 }
565}
566
567pub(crate) struct ExclusionGuard {
574 state: kio::Producer<FrontState>,
575 peer: Origin,
576}
577
578impl ExclusionGuard {
579 fn new(state: &kio::Producer<FrontState>, peer: Origin) -> Option<Arc<Self>> {
582 let mut s = state.write().ok()?;
583 if s.closed {
584 return None;
585 }
586 *s.excluded.entry(peer).or_default() += 1;
587 drop(s);
588 Some(Arc::new(Self {
589 state: state.clone(),
590 peer,
591 }))
592 }
593}
594
595impl Drop for ExclusionGuard {
596 fn drop(&mut self) {
597 let Ok(mut state) = self.state.write() else { return };
598 if let std::collections::hash_map::Entry::Occupied(mut entry) = state.excluded.entry(self.peer) {
599 match entry.get() {
600 1 => drop(entry.remove()),
601 n => *entry.get_mut() = n - 1,
602 }
603 }
604 }
605}
606
607enum Resolved {
614 Found(broadcast::Consumer),
617 Excluded,
619 Missing,
621}
622
623struct OriginNode {
624 broadcast: Option<OriginBroadcast>,
627
628 nested: HashMap<String, Lock<OriginNode>>,
630
631 notify: Lock<NotifyNode>,
633}
634
635impl OriginNode {
636 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
637 Self {
638 broadcast: None,
639 nested: HashMap::new(),
640 notify: Lock::new(NotifyNode::new(parent)),
641 }
642 }
643
644 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
645 let (dir, rest) = path.next_part().expect("leaf called with empty path");
646
647 let next = self.entry(dir);
648 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
649 }
650
651 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
652 match self.nested.get(dir) {
653 Some(next) => next.clone(),
654 None => {
655 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
656 self.nested.insert(dir.to_string(), next.clone());
657 next
658 }
659 }
660 }
661
662 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
666 let Some(existing) = &mut self.broadcast else { return };
667 if !existing.state.same_channel(expect) || existing.announced == announce {
668 return;
669 }
670 existing.announced = announce;
671 let path = existing.path.clone();
672 let consumer = existing.broadcast.consume();
673 let mut notify = self.notify.lock();
674 if announce {
675 notify.announce(&path, &consumer);
676 } else {
677 notify.unannounce(&path);
678 }
679 }
680
681 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
682 self.consume_initial(&mut notify);
683 self.notify.lock().consumers.insert(id, notify);
684 }
685
686 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
687 if let Some(broadcast) = &self.broadcast
690 && broadcast.announced
691 {
692 notify.announce(&broadcast.path, broadcast.broadcast.consume());
693 }
694
695 for nested in self.nested.values() {
697 nested.lock().consume_initial(notify);
698 }
699 }
700
701 fn resolve_broadcast(&self, rest: impl AsPath, exclude: Option<Origin>) -> Resolved {
702 let rest = rest.as_path();
703
704 if let Some((dir, rest)) = rest.next_part() {
705 let Some(node) = self.nested.get(dir) else {
706 return Resolved::Missing;
707 };
708 let node = node.lock();
709 return node.resolve_broadcast(&rest, exclude);
710 }
711
712 let Some(broadcast) = self.broadcast.as_ref() else {
713 return Resolved::Missing;
714 };
715 let Some(origin) = exclude else {
716 return Resolved::Found(broadcast.broadcast.consume());
717 };
718
719 let state = broadcast.state.read();
735 if !state.routes.iter().any(|r| r.route.hops.contains(&origin)) {
736 drop(state);
737 let shared = broadcast.broadcast.consume();
738 return match ExclusionGuard::new(&broadcast.state, origin) {
739 Some(guard) => Resolved::Found(shared.with_exclusion(guard)),
740 None => Resolved::Found(shared),
743 };
744 }
745 match state
746 .dispatch(Some(origin))
747 .and_then(|clean| state.routes.iter().find(|r| r.id == clean))
748 {
749 Some(route) => Resolved::Found(route.source.clone()),
750 None => Resolved::Excluded,
751 }
752 }
753
754 fn unconsume(&mut self, id: ConsumerId) {
755 self.notify.lock().consumers.remove(&id).expect("consumer not found");
756 if self.is_empty() {
757 }
760 }
761
762 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
766 let relative = relative.as_path();
767
768 if let Some((dir, relative)) = relative.next_part() {
769 let Some(nested) = self.nested.get(dir) else { return };
770 let nested = nested.clone();
771 let mut locked = nested.lock();
772 locked.remove(expect, &relative);
773
774 if locked.is_empty() {
775 drop(locked);
776 self.nested.remove(dir);
777 }
778 } else if let Some(existing) = &self.broadcast
779 && existing.state.same_channel(expect)
780 {
781 let existing = self.broadcast.take().expect("checked above");
782 if existing.announced {
783 self.notify.lock().unannounce(&existing.path);
784 }
785 }
786 }
787
788 fn is_empty(&self) -> bool {
789 self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
790 }
791}
792
793#[derive(Clone)]
794struct OriginNodes {
795 nodes: Vec<(PathOwned, Lock<OriginNode>)>,
796}
797
798impl OriginNodes {
799 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
802 let mut roots = Vec::new();
803
804 for (root, state) in &self.nodes {
805 for prefix in prefixes {
806 if root.has_prefix(prefix) {
807 roots.push((root.to_owned(), state.clone()));
809 continue;
810 }
811
812 if let Some(suffix) = prefix.strip_prefix(root) {
813 let nested = state.lock().leaf(&suffix);
815 roots.push((prefix.to_owned(), nested));
816 }
817 }
818 }
819
820 if roots.is_empty() {
821 None
822 } else {
823 Some(Self { nodes: roots })
824 }
825 }
826
827 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
828 let new_root = new_root.as_path();
829 let mut roots = Vec::new();
830
831 if new_root.is_empty() {
832 return Some(self.clone());
833 }
834
835 for (root, state) in &self.nodes {
836 if let Some(suffix) = root.strip_prefix(&new_root) {
837 roots.push((suffix.to_owned(), state.clone()));
839 } else if let Some(suffix) = new_root.strip_prefix(root) {
840 let nested = state.lock().leaf(&suffix);
843 roots.push(("".into(), nested));
844 }
845 }
846
847 if roots.is_empty() {
848 None
849 } else {
850 Some(Self { nodes: roots })
851 }
852 }
853
854 pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
856 let path = path.as_path();
857
858 for (root, state) in &self.nodes {
859 if let Some(suffix) = path.strip_prefix(root) {
860 return Some((state.clone(), suffix.to_owned()));
861 }
862 }
863
864 None
865 }
866}
867
868impl Default for OriginNodes {
869 fn default() -> Self {
870 Self {
871 nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
872 }
873 }
874}
875
876#[derive(Clone)]
878pub struct OriginAnnounce {
879 pub path: PathOwned,
881 pub broadcast: Option<broadcast::Consumer>,
887}
888
889#[derive(Clone)]
891pub struct Producer {
892 info: Origin,
896
897 nodes: OriginNodes,
900
901 root: PathOwned,
903
904 dynamic: kio::Shared<OriginDynamicState>,
908
909 pool: cache::Pool,
912
913 cache_duration: Duration,
916
917 linger: Duration,
920
921 stats: stats::Session,
925}
926
927impl std::ops::Deref for Producer {
928 type Target = Origin;
929
930 fn deref(&self) -> &Self::Target {
931 &self.info
932 }
933}
934
935impl Producer {
936 pub fn new(info: Info) -> Self {
940 Self {
941 info: info.id,
942 nodes: OriginNodes::default(),
943 root: PathOwned::default(),
944 dynamic: kio::Shared::default(),
945 pool: info.pool,
946 cache_duration: info.cache_duration,
947 linger: info.linger,
948 stats: stats::Session::default(),
949 }
950 }
951
952 pub fn with_stats(mut self, session: stats::Session) -> Self {
956 self.stats = session;
957 self
958 }
959
960 pub fn with_linger(mut self, linger: Duration) -> Self {
969 self.linger = linger;
970 self
971 }
972
973 pub fn info(&self) -> Info {
976 Info {
977 id: self.info,
978 pool: self.pool.clone(),
979 cache_duration: self.cache_duration,
980 linger: self.linger,
981 }
982 }
983
984 pub(crate) fn empty(info: Origin) -> Self {
989 Self {
990 info,
991 nodes: OriginNodes { nodes: Vec::new() },
992 root: PathOwned::default(),
993 dynamic: kio::Shared::default(),
994 pool: cache::Pool::default(),
995 cache_duration: Duration::MAX,
996 linger: Duration::ZERO,
997 stats: stats::Session::default(),
998 }
999 }
1000
1001 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
1052 let path = path.as_path();
1053
1054 debug_assert!(
1055 !route.hops.contains(&self.info),
1056 "create_broadcast called with a looping hop chain",
1057 );
1058
1059 let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
1060 let full = self.root.join(&path).to_owned();
1061
1062 if full.parts().count() > Path::MAX_PARTS {
1066 return Err(BoundsExceeded.into());
1067 }
1068
1069 let ingress = self.stats.ingress(&full);
1073
1074 let mut source = broadcast::Info { origin: self.info() }
1075 .produce()
1076 .with_stats(ingress.clone());
1077 source.set_route(route).expect("fresh producer");
1078
1079 web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
1080
1081 Ok(source)
1082 }
1083
1084 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
1090 let prefixes = PathPrefixes::new(prefixes);
1091 Some(Producer {
1092 info: self.info,
1093 nodes: self.nodes.select(&prefixes)?,
1094 root: self.root.clone(),
1095 dynamic: self.dynamic.clone(),
1096 pool: self.pool.clone(),
1097 cache_duration: self.cache_duration,
1098 linger: self.linger,
1099 stats: self.stats.clone(),
1100 })
1101 }
1102
1103 pub fn dynamic(&self) -> Dynamic {
1112 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
1113 }
1114
1115 pub fn consume(&self) -> Consumer {
1120 Consumer::new(
1123 self.info,
1124 self.root.clone(),
1125 self.nodes.clone(),
1126 self.dynamic.clone(),
1127 stats::Session::default(),
1128 )
1129 }
1130
1131 pub fn announces(&self) -> AnnounceProducer {
1137 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
1138 }
1139
1140 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1145 let prefix = prefix.as_path();
1146
1147 Some(Self {
1148 info: self.info,
1149 root: self.root.join(&prefix).to_owned(),
1150 nodes: self.nodes.root(&prefix)?,
1151 dynamic: self.dynamic.clone(),
1152 pool: self.pool.clone(),
1153 cache_duration: self.cache_duration,
1154 linger: self.linger,
1155 stats: self.stats.clone(),
1156 })
1157 }
1158
1159 pub fn root(&self) -> &Path<'_> {
1161 &self.root
1162 }
1163
1164 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1167 self.nodes.nodes.iter().map(|(root, _)| root)
1168 }
1169
1170 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1172 self.root.join(path)
1173 }
1174}
1175
1176const MAX_TRACK_RETRIES: u32 = 3;
1180
1181const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1194
1195struct FrontRoute {
1197 id: u64,
1198 route: broadcast::Route,
1201 source: broadcast::Consumer,
1203}
1204
1205struct FrontState {
1207 path: PathOwned,
1209 self_origin: Origin,
1211 publisher: Option<Origin>,
1219 next_route: u64,
1222 routes: Vec<FrontRoute>,
1223 excluded: HashMap<Origin, usize>,
1230 active: Option<u64>,
1232 linger: Duration,
1235 closed: bool,
1240}
1241
1242impl FrontState {
1243 fn best_route(&self) -> Option<u64> {
1249 let candidates: Vec<&FrontRoute> = self.routes.iter().collect();
1250 self.prefer_untainted(&candidates)
1251 .into_iter()
1252 .min_by_key(|r| route_order(&self.path.as_path(), r))
1253 .map(|r| r.id)
1254 }
1255
1256 fn dispatch(&self, exclude: Option<Origin>) -> Option<u64> {
1263 self.routes
1264 .iter()
1265 .filter(|r| exclude.is_none_or(|origin| !r.route.hops.contains(&origin)))
1266 .min_by_key(|r| route_order(&self.path.as_path(), r))
1267 .map(|r| r.id)
1268 }
1269
1270 fn taints_a_reader(&self, route: &broadcast::Route) -> bool {
1272 route.hops.iter().any(|hop| self.excluded.contains_key(hop))
1273 }
1274
1275 fn prefer_untainted<'a>(&self, candidates: &[&'a FrontRoute]) -> Vec<&'a FrontRoute> {
1285 if self.excluded.is_empty() {
1286 return candidates.to_vec();
1287 }
1288 let clean: Vec<&FrontRoute> = candidates
1289 .iter()
1290 .copied()
1291 .filter(|r| !self.taints_a_reader(&r.route))
1292 .collect();
1293 match clean.is_empty() {
1294 true => candidates.to_vec(),
1295 false => clean,
1296 }
1297 }
1298
1299 fn serve_route(&self, skip: impl Fn(u64) -> bool) -> Option<u64> {
1308 if let Some(active) = self.active
1309 && !skip(active)
1310 && let Some(route) = self.routes.iter().find(|r| r.id == active)
1311 && !self.taints_a_reader(&route.route)
1312 {
1313 return Some(active);
1314 }
1315 let candidates: Vec<&FrontRoute> = self.routes.iter().filter(|r| !skip(r.id)).collect();
1316 self.prefer_untainted(&candidates)
1317 .into_iter()
1318 .min_by_key(|r| route_order(&self.path.as_path(), r))
1319 .map(|r| r.id)
1320 }
1321
1322 fn reselect(&mut self, carrying: bool) {
1340 let best = self.best_route();
1341 if carrying
1342 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1343 && best_id != cur_id
1344 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1345 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1346 && incumbent.route.announce
1347 && candidate.route.cost < incumbent.route.cost
1348 && candidate.route.advertised == 0
1349 && candidate.route.hops.len() >= 2
1350 && !self.handover_allowed(&candidate.route)
1351 {
1352 return;
1354 }
1355 self.active = best;
1356 }
1357
1358 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1370 let name = self.path.as_path();
1371 match route.hops.iter().last() {
1372 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1373 None => true,
1374 }
1375 }
1376
1377 fn routes_snapshot(&self) -> Vec<broadcast::Route> {
1383 let mut routes: Vec<&FrontRoute> = self.routes.iter().collect();
1384 routes.sort_by_key(|r| route_order(&self.path.as_path(), r));
1385 routes.sort_by_key(|r| Some(r.id) != self.active);
1386 routes.into_iter().map(|r| r.route.clone()).collect()
1387 }
1388}
1389
1390fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1400 let mut leaf_guard = leaf.lock();
1405 let routes = state.read().routes_snapshot();
1406 if let Some(advert) = routes.first() {
1407 let announce = advert.announce;
1408 broadcast.clone().set_routes(routes);
1409 leaf_guard.set_announced(state, announce);
1410 }
1411}
1412
1413fn detach_source(
1426 state: &kio::Producer<FrontState>,
1427 broadcast: &broadcast::Producer,
1428 leaf: &Lock<OriginNode>,
1429 id: u64,
1430 graceful: bool,
1431) {
1432 let close = {
1433 let carrying = broadcast.demand().is_used();
1434 let Ok(mut s) = state.write() else { return };
1435 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1436 return;
1437 };
1438 s.routes.remove(pos);
1439 s.reselect(carrying);
1440 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1441 s.closed = true;
1444 true
1445 } else {
1446 false
1447 }
1448 };
1449 if close {
1450 broadcast.abort_spliced(Error::Dropped);
1451 }
1452 sync_front(state, broadcast, leaf);
1453}
1454
1455async fn run_source(
1465 origin: Info,
1466 node: Lock<OriginNode>,
1467 full: PathOwned,
1468 rest: PathOwned,
1469 mut source: broadcast::Consumer,
1470 ingress: stats::Scope,
1471) {
1472 let ctx = AttachContext {
1473 origin: &origin,
1474 node: &node,
1475 full: &full,
1476 rest: &rest,
1477 };
1478
1479 let Ok(mut route) = source.route_changed().await else {
1483 return;
1485 };
1486
1487 let mut announce = route.announce.then(|| ingress.announce());
1492
1493 'attach: loop {
1494 let leaf = if rest.is_empty() {
1499 node.clone()
1500 } else {
1501 node.lock().leaf(&rest)
1502 };
1503
1504 let (state, broadcast, id) = match attach_source(&ctx, &leaf, &source, route.clone()) {
1505 Attach::Ready(state, broadcast, id) => (state, broadcast, id),
1506 Attach::Parked(incumbent) => {
1507 tracing::warn!(
1508 broadcast = %full,
1509 "path already live with a different publisher; parking an offline source until it ends",
1510 );
1511 let update = kio::wait(|waiter| {
1515 if let Poll::Ready(update) = source.poll_route_changed(waiter) {
1516 return Poll::Ready(Some(update));
1517 }
1518 match incumbent.poll(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending }) {
1521 Poll::Ready(_) => Poll::Ready(None),
1522 Poll::Pending => Poll::Pending,
1523 }
1524 })
1525 .await;
1526 match update {
1527 Some(Ok(update)) => {
1529 match (update.announce, announce.is_some()) {
1530 (true, false) => announce = Some(ingress.announce()),
1531 (false, true) => announce = None,
1532 _ => {}
1533 }
1534 route = update;
1535 }
1536 Some(Err(_)) => return,
1538 None => {}
1540 }
1541 continue 'attach;
1542 }
1543 };
1544 let publisher = route.hops.iter().next().copied();
1545
1546 loop {
1547 match source.route_changed().await {
1548 Ok(update) => {
1549 let announced = update.announce;
1550 if update.hops.iter().next().copied() != publisher {
1554 detach_source(&state, &broadcast, &leaf, id, true);
1555 match (announced, announce.is_some()) {
1556 (true, false) => announce = Some(ingress.announce()),
1557 (false, true) => announce = None,
1558 _ => {}
1559 }
1560 route = update;
1561 continue 'attach;
1562 }
1563 {
1564 let carrying = broadcast.demand().is_used();
1565 let Ok(mut s) = state.write() else { return };
1566 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1567 return;
1568 };
1569 if entry.route == update {
1570 continue;
1571 }
1572 entry.route = update;
1573 s.reselect(carrying);
1574 }
1575 match (announced, announce.is_some()) {
1577 (true, false) => announce = Some(ingress.announce()),
1578 (false, true) => announce = None,
1579 _ => {}
1580 }
1581 sync_front(&state, &broadcast, &leaf);
1582 }
1583 Err(_) => {
1584 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1587 return;
1588 }
1589 }
1590 }
1591 }
1592}
1593
1594enum Attach {
1596 Ready(kio::Producer<FrontState>, broadcast::Producer, u64),
1599 Parked(kio::Producer<FrontState>),
1604}
1605
1606struct AttachContext<'a> {
1608 origin: &'a Info,
1609 node: &'a Lock<OriginNode>,
1610 full: &'a PathOwned,
1612 rest: &'a PathOwned,
1614}
1615
1616fn attach_source(
1633 ctx: &AttachContext,
1634 leaf: &Lock<OriginNode>,
1635 source: &broadcast::Consumer,
1636 route: broadcast::Route,
1637) -> Attach {
1638 let publisher = route.hops.iter().next().copied();
1639 let mut leaf_guard = leaf.lock();
1640
1641 if let Some(existing) = &leaf_guard.broadcast {
1644 let mut joined = None;
1645 let carrying = existing.broadcast.demand().is_used();
1646 if let Ok(mut s) = existing.state.write()
1647 && !s.closed
1648 {
1649 if s.publisher == publisher {
1650 let id = s.next_route;
1651 s.next_route += 1;
1652 s.routes.push(FrontRoute {
1653 id,
1654 route: route.clone(),
1655 source: source.clone(),
1656 });
1657 s.reselect(carrying);
1658 joined = Some(id);
1659 } else if !route.announce {
1660 return Attach::Parked(existing.state.clone());
1661 } else {
1662 s.closed = true;
1670 tracing::warn!(broadcast = %ctx.full, "replacing a live broadcast from a different publisher");
1671 }
1672 }
1673 if let Some(id) = joined {
1674 let state = existing.state.clone();
1675 let broadcast = existing.broadcast.clone();
1676 drop(leaf_guard);
1677 sync_front(&state, &broadcast, leaf);
1678 return Attach::Ready(state, broadcast, id);
1679 }
1680 }
1681
1682 let announce = route.announce;
1684 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
1685 origin: ctx.origin.clone(),
1686 });
1687 let _ = broadcast.clone().set_route(route.clone());
1688 let state = kio::Producer::new(FrontState {
1689 path: ctx.full.clone(),
1690 self_origin: ctx.origin.id,
1691 publisher,
1692 next_route: 1,
1693 excluded: HashMap::new(),
1694 routes: vec![FrontRoute {
1695 id: 0,
1696 route,
1697 source: source.clone(),
1698 }],
1699 active: Some(0),
1700 linger: ctx.origin.linger,
1701 closed: false,
1702 });
1703
1704 if let Some(stale) = leaf_guard.broadcast.take()
1708 && stale.announced
1709 {
1710 leaf_guard.notify.lock().unannounce(&stale.path);
1711 }
1712 let entry = OriginBroadcast {
1713 path: ctx.full.clone(),
1714 broadcast: broadcast.clone(),
1715 state: state.clone(),
1716 announced: announce,
1717 };
1718 if entry.announced {
1719 leaf_guard.notify.lock().announce(ctx.full, &broadcast.consume());
1720 }
1721 leaf_guard.broadcast = Some(entry);
1722 drop(leaf_guard);
1723
1724 web_async::spawn(run_front(
1725 state.clone(),
1726 broadcast.clone(),
1727 ctx.node.clone(),
1728 ctx.rest.clone(),
1729 ));
1730
1731 Attach::Ready(state, broadcast, 0)
1732}
1733
1734async fn run_front(
1737 state: kio::Producer<FrontState>,
1738 mut broadcast: broadcast::Producer,
1739 node: Lock<OriginNode>,
1740 rest: PathOwned,
1741) {
1742 enum Step {
1743 Serve(Arc<str>, super::resume::Producer),
1744 Changed,
1746 Expired,
1748 Closed,
1749 }
1750
1751 let linger = state.read().linger;
1752 let mut deadline = kio::time::Deadline::new();
1757
1758 loop {
1759 let empty = {
1760 let s = state.read();
1761 !s.closed && s.routes.is_empty()
1762 };
1763 deadline.set(match (empty, deadline.deadline()) {
1764 (true, None) => web_async::time::Instant::now().checked_add(linger),
1767 (true, at) => at,
1768 (false, _) => None,
1769 });
1770
1771 let step = {
1772 kio::wait(|waiter| {
1773 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1774 return Poll::Ready(Step::Serve(name, resume));
1775 }
1776 match state.poll(waiter, |s| {
1779 if s.closed || s.routes.is_empty() != empty {
1780 Poll::Ready(())
1781 } else {
1782 Poll::Pending
1783 }
1784 }) {
1785 Poll::Ready(Ok(guard)) => {
1786 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1787 }
1788 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1789 Poll::Pending => {}
1790 }
1791 deadline.poll(waiter).map(|_| Step::Expired)
1792 })
1793 .await
1794 };
1795
1796 match step {
1797 Step::Serve(name, resume) => {
1798 web_async::spawn(serve_track(state.clone(), name, resume));
1801 }
1802 Step::Changed => {}
1803 Step::Expired => {
1804 let close = {
1808 let Ok(mut s) = state.write() else { break };
1809 if !s.closed && s.routes.is_empty() {
1810 s.closed = true;
1811 true
1812 } else {
1813 false
1814 }
1815 };
1816 if close {
1817 break;
1818 }
1819 }
1820 Step::Closed => break,
1821 }
1822 }
1823
1824 broadcast.abort_spliced(Error::Dropped);
1826
1827 broadcast.finish();
1829
1830 node.lock().remove(&state, &rest);
1833}
1834
1835async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1844 enum Step {
1845 Closed,
1846 Splice(u64, broadcast::Consumer),
1847 Complete,
1848 Failed,
1849 Resweep,
1852 Idle,
1854 Demand,
1856 }
1857
1858 let mut fails = 0u32;
1859 let mut serving: Option<(u64, track::Consumer)> = None;
1861 let mut refused: HashSet<u64> = HashSet::new();
1867 let mut dead: HashSet<u64> = HashSet::new();
1872 let mut idle_since: Option<web_async::time::Instant> = None;
1874 let mut deadline = kio::time::Deadline::new();
1875
1876 loop {
1877 let serving_id = serving.as_ref().map(|(id, _)| *id);
1878
1879 {
1887 let s = state.read();
1888 refused.retain(|id| s.routes.iter().any(|r| r.id == *id));
1889 dead.retain(|id| s.routes.iter().any(|r| r.id == *id));
1890 let exhausted = !s.routes.is_empty()
1891 && s.serve_route(|id| refused.contains(&id) || dead.contains(&id))
1892 .is_none();
1893 if exhausted && dead.is_empty() {
1896 drop(s);
1897 fails += 1;
1898 if fails >= MAX_TRACK_RETRIES {
1899 tracing::debug!(name = %name, "aborting unservable track");
1900 let _ = resume.abort(Error::Unroutable);
1901 return;
1902 }
1903 refused.clear();
1904 }
1905 }
1906
1907 let used = resume.is_used();
1911 idle_since = match (serving.is_some(), used) {
1912 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
1913 _ => None,
1914 };
1915 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
1916
1917 let step = {
1918 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
1919 kio::wait(|waiter| {
1920 match state.poll(waiter, |s| {
1926 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
1927 if s.closed
1928 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
1929 {
1930 Poll::Ready(())
1931 } else {
1932 Poll::Pending
1933 }
1934 }) {
1935 Poll::Ready(Ok(guard)) => {
1936 if guard.closed {
1937 return Poll::Ready(Step::Closed);
1938 }
1939 let Some(next) = guard.serve_route(skip) else {
1940 return Poll::Ready(Step::Resweep);
1941 };
1942 let source = guard
1943 .routes
1944 .iter()
1945 .find(|r| r.id == next)
1946 .expect("servable source in table")
1947 .source
1948 .clone();
1949 return Poll::Ready(Step::Splice(next, source));
1950 }
1951 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1952 Poll::Pending => {}
1953 }
1954
1955 let edge = match used {
1960 true => resume.poll_unused(waiter),
1961 false => resume.poll_used(waiter),
1962 };
1963 if edge.is_ready() {
1964 return Poll::Ready(Step::Demand);
1965 }
1966
1967 if let Some((_, track)) = &serving
1970 && let Poll::Ready(result) = track.poll_complete(waiter)
1971 {
1972 return Poll::Ready(match result {
1973 Ok(()) => Step::Complete,
1974 Err(_) => Step::Failed,
1975 });
1976 }
1977
1978 deadline.poll(waiter).map(|_| Step::Idle)
1979 })
1980 .await
1981 };
1982
1983 match step {
1984 Step::Closed => return,
1986 Step::Complete => {
1987 let _ = resume.finish();
1988 return;
1989 }
1990 Step::Failed => {
1991 serving = None;
1994 }
1995 Step::Demand => {}
1997 Step::Resweep => {}
2001 Step::Idle => {
2002 if resume.release().is_err() {
2007 return;
2009 }
2010 serving = None;
2011 }
2012 Step::Splice(id, source) => {
2013 let attempt = match source.track(&name) {
2017 Ok(track) => {
2018 let query = track.info().into_inner();
2021 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2022 let info = kio::wait(|waiter| {
2023 if let Poll::Ready(result) = query.poll(waiter) {
2024 return Poll::Ready(Some(result));
2025 }
2026 match state.poll(waiter, |s| {
2027 if s.closed || s.serve_route(skip) != Some(id) {
2028 Poll::Ready(())
2029 } else {
2030 Poll::Pending
2031 }
2032 }) {
2033 Poll::Ready(_) => Poll::Ready(None),
2034 Poll::Pending => Poll::Pending,
2035 }
2036 })
2037 .await;
2038 match info {
2039 None => continue,
2042 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2046 Poll::Ready(Err(err)) => Err(err),
2047 _ => Ok(track),
2048 },
2049 Some(Err(err)) => Err(err),
2050 }
2051 }
2052 Err(err) => Err(err),
2053 };
2054
2055 match attempt {
2056 Ok(track) => {
2057 if resume.takeover(&track).is_err() {
2058 return;
2061 }
2062 fails = 0;
2067 dead.clear();
2068 serving = Some((id, track));
2069 }
2070 Err(_) if source.is_closing() => {
2074 dead.insert(id);
2075 serving = None;
2076 }
2077 Err(err) => {
2081 tracing::debug!(name = %name, source = id, %err, "source refused track");
2082 refused.insert(id);
2083 serving = None;
2084 }
2085 }
2086 }
2087 }
2088 }
2089}
2090
2091#[derive(Default)]
2097struct OriginDynamicState {
2098 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2101
2102 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2108}
2109
2110#[derive(Default)]
2117struct PendingBroadcast {
2118 resolved: Option<Result<broadcast::Consumer, Error>>,
2119}
2120
2121pub struct Dynamic {
2132 info: Origin,
2133 root: PathOwned,
2134 state: kio::Shared<OriginDynamicState>,
2135}
2136
2137impl Clone for Dynamic {
2138 fn clone(&self) -> Self {
2139 self.state.lock().requests.add_handler();
2143
2144 Self {
2145 info: self.info,
2146 root: self.root.clone(),
2147 state: self.state.clone(),
2148 }
2149 }
2150}
2151
2152impl Dynamic {
2153 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2154 state.lock().requests.add_handler();
2155
2156 Self { info, root, state }
2157 }
2158
2159 pub fn info(&self) -> &Origin {
2161 &self.info
2162 }
2163
2164 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2166 let mut state = ready!(self.state.poll(waiter, |state| {
2167 if state.requests.has_queued() {
2168 Poll::Ready(())
2169 } else {
2170 Poll::Pending
2171 }
2172 }));
2173
2174 let path = state.requests.pop().expect("predicate guaranteed a request");
2175 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2181 Poll::Ready(Ok(Request {
2182 path,
2183 producer,
2184 state: self.state.clone(),
2185 }))
2186 }
2187
2188 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2191 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2192 }
2193
2194 pub fn root(&self) -> &Path<'_> {
2196 &self.root
2197 }
2198}
2199
2200impl Drop for Dynamic {
2201 fn drop(&mut self) {
2202 let mut state = self.state.lock();
2205 if state.requests.remove_handler() {
2206 state.requests.drain_queued();
2210 }
2211 }
2212}
2213
2214pub struct Request {
2221 path: PathOwned,
2223
2224 producer: kio::Producer<PendingBroadcast>,
2227
2228 state: kio::Shared<OriginDynamicState>,
2230}
2231
2232impl Request {
2233 pub fn path(&self) -> &Path<'_> {
2235 &self.path
2236 }
2237
2238 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2244 let broadcast = broadcast.consume();
2245
2246 let resolved = {
2252 let mut state = self.state.lock();
2253 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2254 state
2255 .requests
2256 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2257 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2258 };
2259
2260 if let Ok(mut pending) = self.producer.write() {
2261 pending.resolved = Some(Ok(resolved));
2262 }
2263 }
2265
2266 pub fn reject(self, err: Error) {
2268 self.state
2269 .lock()
2270 .requests
2271 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2272 if let Ok(mut state) = self.producer.write() {
2273 state.resolved = Some(Err(err));
2274 }
2275 }
2276}
2277
2278impl Drop for Request {
2279 fn drop(&mut self) {
2280 self.state
2288 .lock()
2289 .requests
2290 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2291 }
2292}
2293
2294pub struct Requesting {
2301 inner: RequestState,
2302 stats: stats::Scope,
2305}
2306
2307enum RequestState {
2308 Ready(broadcast::Consumer),
2310 Failed(Error),
2313 Pending(kio::Consumer<PendingBroadcast>),
2315}
2316
2317impl Requesting {
2318 fn ready(broadcast: broadcast::Consumer) -> Self {
2319 Self {
2320 inner: RequestState::Ready(broadcast),
2321 stats: stats::Scope::default(),
2322 }
2323 }
2324
2325 fn failed(error: Error) -> Self {
2326 Self {
2327 inner: RequestState::Failed(error),
2328 stats: stats::Scope::default(),
2329 }
2330 }
2331
2332 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2333 Self {
2334 inner: RequestState::Pending(consumer),
2335 stats: stats::Scope::default(),
2336 }
2337 }
2338
2339 fn with_stats(mut self, scope: stats::Scope) -> Self {
2340 self.stats = scope;
2341 self
2342 }
2343
2344 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2346 match &self.inner {
2347 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2348 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2349 RequestState::Pending(consumer) => Poll::Ready(
2350 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2351 Some(result) => Poll::Ready(result.clone()),
2352 None => Poll::Pending,
2353 })) {
2354 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2355 Err(_closed) => Err(Error::Unroutable),
2357 },
2358 ),
2359 }
2360 }
2361}
2362
2363impl kio::Pollable for Requesting {
2364 type Output = Result<broadcast::Consumer, Error>;
2365
2366 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2367 self.poll_ok(waiter)
2368 }
2369}
2370
2371pub trait Consume<T> {
2379 fn consume(&self) -> T;
2381}
2382
2383impl<T, U: Consume<T>> Consume<T> for &U {
2384 fn consume(&self) -> T {
2385 (**self).consume()
2386 }
2387}
2388
2389impl Consume<Consumer> for Producer {
2390 fn consume(&self) -> Consumer {
2391 Consumer::new(
2395 self.info,
2396 self.root.clone(),
2397 self.nodes.clone(),
2398 self.dynamic.clone(),
2399 stats::Session::default(),
2400 )
2401 }
2402}
2403
2404impl Consume<Consumer> for Consumer {
2405 fn consume(&self) -> Consumer {
2406 self.clone()
2407 }
2408}
2409
2410impl Consume<broadcast::Consumer> for broadcast::Producer {
2411 fn consume(&self) -> broadcast::Consumer {
2412 self.consume()
2414 }
2415}
2416
2417impl Consume<broadcast::Consumer> for broadcast::Consumer {
2418 fn consume(&self) -> broadcast::Consumer {
2419 self.clone()
2420 }
2421}
2422
2423impl Consume<track::Consumer> for track::Producer {
2424 fn consume(&self) -> track::Consumer {
2425 self.consume()
2426 }
2427}
2428
2429impl Consume<track::Consumer> for track::Consumer {
2430 fn consume(&self) -> track::Consumer {
2431 self.clone()
2432 }
2433}
2434
2435#[derive(Clone)]
2441pub struct Consumer {
2442 info: Origin,
2444 nodes: OriginNodes,
2445
2446 root: PathOwned,
2448
2449 dynamic: kio::Shared<OriginDynamicState>,
2452
2453 stats: stats::Session,
2457
2458 exclude: Option<Origin>,
2462}
2463
2464impl std::ops::Deref for Consumer {
2465 type Target = Origin;
2466
2467 fn deref(&self) -> &Self::Target {
2468 &self.info
2469 }
2470}
2471
2472impl Consumer {
2473 fn new(
2474 info: Origin,
2475 root: PathOwned,
2476 nodes: OriginNodes,
2477 dynamic: kio::Shared<OriginDynamicState>,
2478 stats: stats::Session,
2479 ) -> Self {
2480 Self {
2481 info,
2482 nodes,
2483 root,
2484 dynamic,
2485 stats,
2486 exclude: None,
2487 }
2488 }
2489
2490 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2495 self.exclude = Some(peer);
2496 self
2497 }
2498
2499 pub fn with_stats(mut self, session: stats::Session) -> Self {
2503 self.stats = session;
2504 self
2505 }
2506
2507 fn untagged(&self) -> Self {
2511 Self {
2512 stats: stats::Session::default(),
2513 ..self.clone()
2514 }
2515 }
2516
2517 pub(crate) fn empty(&self) -> Self {
2522 Self {
2523 info: self.info,
2524 nodes: OriginNodes { nodes: Vec::new() },
2525 root: self.root.clone(),
2526 dynamic: self.dynamic.clone(),
2527 stats: self.stats.clone(),
2528 exclude: self.exclude,
2529 }
2530 }
2531
2532 pub fn announced(&self) -> AnnounceConsumer {
2539 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2540 }
2541
2542 pub fn consume(&self) -> Self {
2544 self.clone()
2545 }
2546
2547 fn resolve(&self, path: impl AsPath) -> Resolved {
2556 let path = path.as_path();
2557 let Some((root, rest)) = self.nodes.get(&path) else {
2558 return Resolved::Missing;
2559 };
2560 let state = root.lock();
2561 state.resolve_broadcast(&rest, self.exclude)
2562 }
2563
2564 #[cfg(test)]
2566 fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2567 match self.resolve(path) {
2568 Resolved::Found(broadcast) => Some(broadcast),
2569 Resolved::Excluded | Resolved::Missing => None,
2570 }
2571 }
2572
2573 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2585 let path = path.as_path();
2586
2587 let consumer = self.scope(std::slice::from_ref(&path))?;
2589
2590 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2594 return None;
2595 }
2596
2597 let mut announced = consumer.untagged().announced();
2601 let scope = self.stats.egress(self.root.join(&path).to_owned());
2602 loop {
2603 let OriginAnnounce {
2604 path: announced_path,
2605 broadcast,
2606 } = announced.next().await?;
2607 if announced_path.as_path() == path
2609 && let Some(broadcast) = broadcast
2610 {
2611 return Some(broadcast.with_stats(scope));
2612 }
2613 }
2614 }
2615
2616 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2622 let prefixes = PathPrefixes::new(prefixes);
2623 Some(Consumer {
2624 info: self.info,
2625 root: self.root.clone(),
2626 nodes: self.nodes.select(&prefixes)?,
2627 dynamic: self.dynamic.clone(),
2628 stats: self.stats.clone(),
2629 exclude: self.exclude,
2630 })
2631 }
2632
2633 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2652 let path = path.as_path();
2653
2654 let absolute = self.root.join(&path).to_owned();
2658 let scope = self.stats.egress(&absolute);
2659
2660 match self.resolve(&path) {
2666 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
2667 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
2668 Resolved::Missing => {}
2669 }
2670
2671 let mut state = self.dynamic.lock();
2672
2673 if let Some(weak) = state.served.get(&absolute) {
2677 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2678 }
2679
2680 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2683 producer.consume()
2684 } else {
2685 let producer = kio::Producer::<PendingBroadcast>::default();
2686 let consumer = producer.consume();
2687 if state.requests.insert(absolute, producer).is_err() {
2688 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2689 }
2690 consumer
2691 };
2692
2693 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2694 }
2695
2696 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2701 let prefix = prefix.as_path();
2702
2703 Some(Self {
2704 info: self.info,
2705 root: self.root.join(&prefix).to_owned(),
2706 nodes: self.nodes.root(&prefix)?,
2707 dynamic: self.dynamic.clone(),
2708 stats: self.stats.clone(),
2709 exclude: self.exclude,
2710 })
2711 }
2712
2713 pub fn root(&self) -> &Path<'_> {
2715 &self.root
2716 }
2717
2718 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2721 self.nodes.nodes.iter().map(|(root, _)| root)
2722 }
2723
2724 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2726 self.root.join(path)
2727 }
2728}
2729
2730#[derive(Clone)]
2735pub struct AnnounceProducer {
2736 nodes: OriginNodes,
2737 root: PathOwned,
2738}
2739
2740impl AnnounceProducer {
2741 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2742 Self { nodes, root }
2743 }
2744
2745 pub fn consume(&self) -> AnnounceConsumer {
2751 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2754 }
2755
2756 pub fn root(&self) -> &Path<'_> {
2758 &self.root
2759 }
2760}
2761
2762pub struct AnnounceConsumer {
2767 id: ConsumerId,
2768 nodes: OriginNodes,
2769 root: PathOwned,
2770
2771 state: kio::Producer<OriginConsumerState>,
2774
2775 stats: stats::Session,
2778
2779 guards: HashMap<PathOwned, stats::Announce>,
2783}
2784
2785impl AnnounceConsumer {
2786 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2787 let state = kio::Producer::<OriginConsumerState>::default();
2788 let id = ConsumerId::new();
2789
2790 for (_, node) in &nodes.nodes {
2791 let notify = AnnounceConsumerNotify {
2792 root: root.clone(),
2793 state: state.clone(),
2794 };
2795 node.lock().consume(id, notify);
2796 }
2797
2798 Self {
2799 id,
2800 nodes,
2801 root,
2802 state,
2803 stats,
2804 guards: HashMap::new(),
2805 }
2806 }
2807
2808 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2814 let OriginAnnounce { path, broadcast } = update;
2815 let absolute = self.root.join(&path).to_owned();
2816 match broadcast {
2817 Some(broadcast) => {
2818 let scope = self.stats.egress(&absolute);
2819 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2820 OriginAnnounce {
2821 path,
2822 broadcast: Some(broadcast.with_stats(scope)),
2823 }
2824 }
2825 None => {
2826 self.guards.remove(&absolute);
2827 OriginAnnounce { path, broadcast: None }
2828 }
2829 }
2830 }
2831
2832 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2839 kio::wait(|waiter| self.poll_next(waiter)).await
2840 }
2841
2842 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2848 let update = {
2849 let mut state = match ready!(self.state.poll(waiter, |state| {
2850 if state.pending.is_empty() {
2851 Poll::Pending
2852 } else {
2853 Poll::Ready(())
2854 }
2855 })) {
2856 Ok(state) => state,
2857 Err(_) => return Poll::Ready(None),
2859 };
2860 state.take().expect("predicate guaranteed an update")
2861 };
2862 Poll::Ready(Some(self.attribute(update)))
2863 }
2864
2865 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2870 let update = self.state.write().ok()?.take()?;
2871 Some(self.attribute(update))
2872 }
2873
2874 pub fn is_closed(&self) -> bool {
2876 self.state.write().is_err()
2877 }
2878
2879 pub fn root(&self) -> &Path<'_> {
2881 &self.root
2882 }
2883
2884 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2886 self.root.join(path)
2887 }
2888}
2889
2890impl Drop for AnnounceConsumer {
2891 fn drop(&mut self) {
2892 for (_, root) in &self.nodes.nodes {
2893 root.lock().unconsume(self.id);
2894 }
2895 }
2896}
2897
2898#[cfg(test)]
2899use futures::FutureExt;
2900
2901#[cfg(test)]
2902#[allow(missing_docs)] impl AnnounceConsumer {
2904 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2905 let expected = expected.as_path();
2906 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2907 assert_eq!(announce.path, expected, "wrong path");
2908 let announced = announce.broadcast.expect("should be an active announce");
2909 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2910 }
2911
2912 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2916 let expected = expected.as_path();
2917 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2918 assert_eq!(announce.path, expected, "wrong path");
2919 announce.broadcast.expect("should be an active announce")
2920 }
2921
2922 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2923 let expected = expected.as_path();
2924 let announce = self.try_next().expect("no next");
2925 assert_eq!(announce.path, expected, "wrong path");
2926 let announced = announce.broadcast.expect("should be an active announce");
2927 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2928 }
2929
2930 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2932 let expected = expected.as_path();
2933 let announce = self.try_next().expect("no next");
2934 assert_eq!(announce.path, expected, "wrong path");
2935 announce.broadcast.expect("should be an active announce")
2936 }
2937
2938 pub fn assert_next_none(&mut self, expected: impl AsPath) {
2939 let expected = expected.as_path();
2940 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2941 assert_eq!(announce.path, expected, "wrong path");
2942 assert!(announce.broadcast.is_none(), "should be unannounced");
2943 }
2944
2945 pub fn assert_next_wait(&mut self) {
2946 if let Some(res) = self.next().now_or_never() {
2947 panic!("next should block: got {:?}", res.map(|a| a.path));
2948 }
2949 }
2950
2951 }
2960
2961#[cfg(test)]
2962mod tests {
2963 use crate::coding::Decode;
2964 use crate::group;
2965
2966 use super::*;
2967
2968 fn announce() -> broadcast::Route {
2970 broadcast::Route::new().with_announce(true)
2971 }
2972
2973 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
2979 let name = Path::new(name);
2980 let peer_key = fnv_key(&name, [peer]);
2981 (100u64..)
2982 .map(|id| Origin::new(id).unwrap())
2983 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
2984 .unwrap()
2985 }
2986
2987 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
2990 let source = broadcast::Info::new().produce().consume();
2991 FrontState {
2992 path: Path::new("test").to_owned(),
2993 self_origin,
2994 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
2995 next_route: routes.len() as u64,
2996 excluded: HashMap::new(),
2997 routes: routes
2998 .into_iter()
2999 .enumerate()
3000 .map(|(id, route)| FrontRoute {
3001 id: id as u64,
3002 route,
3003 source: source.clone(),
3004 })
3005 .collect(),
3006 active: Some(0),
3007 linger: Duration::ZERO,
3008 closed: false,
3009 }
3010 }
3011
3012 fn sibling_route(peer: Origin) -> broadcast::Route {
3015 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3016 announce().with_hops(hops)
3017 }
3018
3019 fn upstream_route(cost: u64) -> broadcast::Route {
3021 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3022 announce().with_hops(hops).with_cost(cost)
3023 }
3024
3025 #[test]
3029 fn test_carrying_gate_keys() {
3030 let peer = Origin::new(3).unwrap();
3031
3032 let mut lost = front_state(
3034 origin_keyed("test", peer, false),
3035 vec![upstream_route(10), sibling_route(peer)],
3036 );
3037 lost.reselect(true);
3038 assert_eq!(
3039 lost.active,
3040 Some(0),
3041 "carrying front re-parented onto a higher-keyed peer"
3042 );
3043 lost.reselect(false);
3044 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3045
3046 let mut won = front_state(
3048 origin_keyed("test", peer, true),
3049 vec![upstream_route(10), sibling_route(peer)],
3050 );
3051 won.reselect(true);
3052 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3053 }
3054
3055 #[test]
3060 fn test_carrying_gate_symmetric_race() {
3061 let a = Origin::new(1).unwrap();
3062 let b = Origin::new(2).unwrap();
3063
3064 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3065 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3066 a_view.reselect(true);
3067 b_view.reselect(true);
3068
3069 let a_moved = a_view.active == Some(1);
3070 let b_moved = b_view.active == Some(1);
3071 assert!(
3072 a_moved != b_moved,
3073 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3074 );
3075 }
3076
3077 #[test]
3082 fn test_carrying_switches_to_benign_routes() {
3083 let peer = Origin::new(3).unwrap();
3084 let lost = origin_keyed("test", peer, false);
3085
3086 let mut forwarder = sibling_route(peer).with_cost(4);
3088 forwarder.advertised = 4;
3089 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3090 state.reselect(true);
3091 assert_eq!(
3092 state.active,
3093 Some(1),
3094 "a cheaper forwarder path must win while carrying"
3095 );
3096
3097 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3099 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3100 state.reselect(true);
3101 assert_eq!(
3102 state.active,
3103 Some(1),
3104 "a direct publisher route must win while carrying"
3105 );
3106
3107 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3112 state.reselect(true);
3113 assert_eq!(
3114 state.active,
3115 Some(1),
3116 "a reconnect on an identical chain must win while carrying"
3117 );
3118 }
3119
3120 #[test]
3123 fn test_carrying_gate_ignores_unannounced_incumbent() {
3124 let peer = Origin::new(3).unwrap();
3125 let unannounced = upstream_route(10).with_announce(false);
3126 let mut state = front_state(
3127 origin_keyed("test", peer, false),
3128 vec![unannounced, sibling_route(peer)],
3129 );
3130 state.reselect(true);
3131 assert_eq!(
3132 state.active,
3133 Some(1),
3134 "an unannounced incumbent must always be displaced"
3135 );
3136 }
3137
3138 async fn settle() {
3141 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3142 }
3143
3144 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3147 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3148 .await
3149 .expect("timed out waiting for a track request")
3150 .expect("source closed");
3151 assert_eq!(request.name(), name, "unexpected track dispatched");
3152 request.accept(None)
3153 }
3154
3155 #[tokio::test]
3159 async fn test_stats_tagged_end_to_end() {
3160 use crate::Timestamp;
3161 use crate::stats::{Config, Registry, Tier};
3162 use bytes::Bytes;
3163
3164 tokio::time::pause();
3165
3166 let registry = Registry::new(Config::new());
3167 let ctx = registry.tier(Tier::default()).session("acme");
3168
3169 let origin = Origin::random().produce();
3170 let ingress = origin.clone().with_stats(ctx.clone());
3171 let egress = origin.consume().with_stats(ctx.clone());
3172
3173 let mut announced = egress.announced();
3176
3177 let source = ingress.create_broadcast("demo", announce()).unwrap();
3179 let mut dynamic = source.dynamic();
3180 settle().await;
3181 settle().await;
3182
3183 let update = announced.next().await.unwrap();
3185 assert_eq!(update.path.as_str(), "demo");
3186 let broadcast = update.broadcast.unwrap();
3187
3188 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3190 let mut producer = accept_track(&mut dynamic, "video").await;
3191 settle().await;
3192 let mut sub = subscribing.await.unwrap();
3193
3194 let mut group = producer.append_group().unwrap();
3196 group
3197 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3198 .unwrap();
3199 group
3200 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3201 .unwrap();
3202 group.finish().unwrap();
3203
3204 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3206 let mut frames = 0;
3207 while let Some(frame) = group_c.read_frame().await.unwrap() {
3208 assert_eq!(frame.payload.len(), 5);
3209 frames += 1;
3210 }
3211 assert_eq!(frames, 2);
3212 settle().await;
3213
3214 let report = registry.report();
3215 let entry = report
3216 .traffic
3217 .iter()
3218 .find(|e| e.path.as_str() == "demo")
3219 .expect("demo tracked");
3220 let path_len = "demo".len() as u64;
3221
3222 let egress = &entry.publisher;
3224 assert_eq!(egress.announced, 1, "one egress announce");
3225 assert_eq!(egress.announced_bytes, path_len);
3226 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3227 assert_eq!(egress.broadcasts, 1, "one viewer");
3228 assert_eq!(egress.groups, 1);
3229 assert_eq!(egress.frames, 2);
3230 assert_eq!(egress.bytes, 10);
3231 assert_eq!(egress.fetches, 0);
3232
3233 let ingress = &entry.subscriber;
3235 assert_eq!(ingress.announced, 1, "one ingress announce");
3236 assert_eq!(ingress.announced_bytes, path_len);
3237 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3238 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3239 assert_eq!(ingress.groups, 1);
3240 assert_eq!(ingress.frames, 2);
3241 assert_eq!(ingress.bytes, 10);
3242
3243 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3245 let _ = fetched;
3246 settle().await;
3247 let report = registry.report();
3248 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3249 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3250 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3251 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3252 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3256 }
3257
3258 #[tokio::test]
3263 async fn test_stats_read_frame_counts_once() {
3264 use crate::Timestamp;
3265 use crate::stats::{Config, Registry, Tier};
3266 use bytes::Bytes;
3267
3268 tokio::time::pause();
3269
3270 let registry = Registry::new(Config::new());
3271 let ctx = registry.tier(Tier::default()).session("acme");
3272
3273 let origin = Origin::random().produce();
3274 let ingress = origin.clone().with_stats(ctx.clone());
3275 let egress = origin.consume().with_stats(ctx.clone());
3276
3277 let mut announced = egress.announced();
3278 let source = ingress.create_broadcast("demo", announce()).unwrap();
3279 let mut dynamic = source.dynamic();
3280 settle().await;
3281 settle().await;
3282
3283 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3284 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3285 let mut producer = accept_track(&mut dynamic, "video").await;
3286 settle().await;
3287 let mut sub = subscribing.await.unwrap();
3288
3289 producer
3291 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3292 .unwrap();
3293
3294 let frame = sub.read_frame().await.unwrap().expect("frame");
3295 assert_eq!(frame.payload.len(), 5);
3296 settle().await;
3297
3298 let report = registry.report();
3299 let entry = report
3300 .traffic
3301 .iter()
3302 .find(|e| e.path.as_str() == "demo")
3303 .expect("demo tracked");
3304 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3305 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3306 assert_eq!(
3307 entry.publisher.bytes, 5,
3308 "payload counted once, not zero and not doubled"
3309 );
3310 }
3311
3312 #[tokio::test]
3316 async fn test_stats_datagrams_counted_both_sides() {
3317 use crate::Timestamp;
3318 use crate::stats::{Config, Registry, Tier};
3319
3320 tokio::time::pause();
3321
3322 let registry = Registry::new(Config::new());
3323 let ctx = registry.tier(Tier::default()).session("acme");
3324
3325 let origin = Origin::random().produce();
3326 let ingress = origin.clone().with_stats(ctx.clone());
3327 let egress = origin.consume().with_stats(ctx.clone());
3328
3329 let mut announced = egress.announced();
3330 let source = ingress.create_broadcast("demo", announce()).unwrap();
3331 let mut dynamic = source.dynamic();
3332 settle().await;
3333 settle().await;
3334
3335 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3336 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3337 let mut producer = accept_track(&mut dynamic, "video").await;
3338 settle().await;
3339 let mut sub = subscribing.await.unwrap();
3340
3341 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3342 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3343 assert_eq!(&datagram.payload[..], b"hello");
3344 settle().await;
3345
3346 let report = registry.report();
3347 let entry = report
3348 .traffic
3349 .iter()
3350 .find(|e| e.path.as_str() == "demo")
3351 .expect("demo tracked");
3352
3353 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3354 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3355 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3356 assert_eq!(traffic.frames, 1, "{side}: one frame");
3357 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3358 }
3359 }
3360
3361 #[test]
3362 fn origin_rejects_reserved_ids() {
3363 assert!(Origin::new(0).is_err());
3364 assert!(Origin::new(1u64 << 62).is_err());
3365 assert_eq!(Origin::new(1).unwrap().id(), 1);
3366
3367 let mut zero = [0u8].as_slice();
3368 assert_eq!(
3369 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3370 Origin::UNKNOWN
3371 );
3372 }
3373
3374 #[test]
3375 fn origin_list_push_fails_at_limit() {
3376 let mut list = OriginList::new();
3377 for _ in 0..MAX_HOPS {
3378 list.push(Origin::random()).unwrap();
3379 }
3380 assert_eq!(list.len(), MAX_HOPS);
3381 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3382 }
3383
3384 #[test]
3385 fn origin_list_replace_first() {
3386 let mut list = OriginList::new();
3387 for _ in 0..3 {
3388 list.push(Origin::UNKNOWN).unwrap();
3389 }
3390
3391 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3393 assert_eq!(
3394 list.as_slice(),
3395 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3396 );
3397
3398 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3400 assert_eq!(list.len(), 3);
3401 }
3402
3403 #[test]
3404 fn origin_list_try_from_vec_enforces_limit() {
3405 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3406 assert!(OriginList::try_from(under).is_ok());
3407
3408 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3409 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3410 }
3411
3412 #[tokio::test]
3413 async fn test_announce() {
3414 tokio::time::pause();
3415
3416 let origin = Origin::random().produce();
3417
3418 let mut consumer1 = origin.consume().announced();
3419 consumer1.assert_next_wait();
3420
3421 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3423 settle().await;
3424
3425 consumer1.assert_next_some("test1");
3426 consumer1.assert_next_wait();
3427
3428 let mut consumer2 = origin.consume().announced();
3431
3432 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3434 settle().await;
3435
3436 consumer1.assert_next_some("test2");
3437 consumer1.assert_next_wait();
3438
3439 consumer2.assert_next_some("test1");
3440 consumer2.assert_next_some("test2");
3441 consumer2.assert_next_wait();
3442
3443 broadcast1.finish();
3445 settle().await;
3446
3447 consumer1.assert_next_none("test1");
3449 consumer2.assert_next_none("test1");
3450 consumer1.assert_next_wait();
3451 consumer2.assert_next_wait();
3452
3453 let mut consumer3 = origin.consume().announced();
3455 consumer3.assert_next_some("test2");
3456 consumer3.assert_next_wait();
3457
3458 broadcast2.finish();
3459 settle().await;
3460
3461 consumer1.assert_next_none("test2");
3462 consumer2.assert_next_none("test2");
3463 consumer3.assert_next_none("test2");
3464 }
3465
3466 #[tokio::test]
3470 async fn test_duplicate() {
3471 tokio::time::pause();
3472
3473 let origin = Origin::random().produce();
3474 let consumer = origin.consume();
3475 let mut announced = consumer.announced();
3476
3477 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3478 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3479 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3480 settle().await;
3481 assert!(consumer.get_broadcast("test").is_some());
3482
3483 announced.assert_next_some("test");
3484 announced.assert_next_wait();
3485
3486 broadcast2.finish();
3488 settle().await;
3489 assert!(consumer.get_broadcast("test").is_some());
3490 announced.assert_next_wait();
3491
3492 broadcast1.finish();
3494 settle().await;
3495 assert!(consumer.get_broadcast("test").is_some());
3496 announced.assert_next_wait();
3497
3498 broadcast3.finish();
3500 settle().await;
3501 assert!(consumer.get_broadcast("test").is_none());
3502
3503 announced.assert_next_none("test");
3504 announced.assert_next_wait();
3505 }
3506
3507 #[tokio::test]
3510 async fn test_route_failover() {
3511 tokio::time::pause();
3512
3513 let origin = Origin::random().produce();
3514 let consumer = origin.consume();
3515 let mut announced = consumer.announced();
3516
3517 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3520 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3521
3522 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3524 let mut dynamic_a = source_a.dynamic();
3525 settle().await;
3526 settle().await;
3527 let broadcast = consumer.request_broadcast("test").await.unwrap();
3528 announced.assert_next_some("test");
3529
3530 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3532 let mut dynamic_b = source_b.dynamic();
3533 settle().await;
3534 settle().await;
3535 announced.assert_next_wait();
3536
3537 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3539 let mut producer = accept_track(&mut dynamic_a, "video").await;
3540 settle().await;
3541 dynamic_b.assert_no_request();
3542
3543 let mut sub = subscribing.await.unwrap();
3544 sub.assert_no_group();
3547 assert_eq!(producer.subscription().unwrap().group_start, None);
3548
3549 producer.append_group().unwrap();
3550 producer.append_group().unwrap();
3551 assert_eq!(sub.assert_group().sequence, 0);
3552 assert_eq!(sub.assert_group().sequence, 1);
3553
3554 producer.abort(Error::Dropped).unwrap();
3558 source_a.abort(Error::Dropped).unwrap();
3559 drop(dynamic_a);
3560 settle().await;
3561 announced.assert_next_wait();
3562
3563 let mut producer = accept_track(&mut dynamic_b, "video").await;
3566 settle().await;
3567 sub.assert_no_group();
3568 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3569 producer.create_group(group::Info { sequence: 1 }).unwrap();
3570 producer.create_group(group::Info { sequence: 2 }).unwrap();
3571 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3572 sub.assert_not_closed();
3573 }
3574
3575 #[tokio::test]
3578 async fn test_broadcast_route_watch() {
3579 let mut producer = broadcast::Info::new().produce();
3580 let mut consumer = producer.consume();
3581
3582 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3584
3585 producer.set_route(broadcast::Route::default()).unwrap();
3587 assert!(consumer.route_changed().now_or_never().is_none());
3588
3589 let mut hops = OriginList::new();
3590 hops.push(Origin::new(7).unwrap()).unwrap();
3591 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3592 producer.set_route(route.clone()).unwrap();
3593 assert_eq!(consumer.route_changed().await.unwrap(), route);
3594
3595 let mut fresh = producer.consume();
3597 assert_eq!(fresh.route_changed().await.unwrap(), route);
3598
3599 drop(producer);
3600 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3601 }
3602
3603 #[tokio::test]
3607 async fn test_route_cost_update() {
3608 tokio::time::pause();
3609
3610 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3614 let consumer = origin.consume();
3615 let mut announced = consumer.announced();
3616
3617 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3620 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3621
3622 let mut source_a = origin
3624 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3625 .unwrap();
3626 let mut dynamic_a = source_a.dynamic();
3627 settle().await;
3628 let broadcast = consumer.request_broadcast("test").await.unwrap();
3629 announced.assert_next_some("test");
3630
3631 let mut watch = broadcast.clone();
3632 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3633
3634 let mut source_b = origin
3635 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3636 .unwrap();
3637 let mut dynamic_b = source_b.dynamic();
3638 settle().await;
3639 assert!(
3640 watch.route_changed().now_or_never().is_none(),
3641 "a losing standby must not change the advertised route"
3642 );
3643
3644 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3646 let mut producer = accept_track(&mut dynamic_a, "video").await;
3647 settle().await;
3648 let mut sub = subscribing.await.unwrap();
3649 producer.append_group().unwrap();
3650 assert_eq!(sub.assert_group().sequence, 0);
3651
3652 source_a
3655 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3656 .unwrap();
3657 settle().await;
3658 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3659 announced.assert_next_wait();
3660
3661 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3662 settle().await;
3663 sub.assert_no_group();
3666 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3667 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3668 assert_eq!(sub.assert_group().sequence, 1);
3669 sub.assert_not_closed();
3670
3671 source_b
3673 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3674 .unwrap();
3675 settle().await;
3676 let advertised = watch.route_changed().await.unwrap();
3677 assert_eq!(advertised.hops, hops_b);
3678 assert_eq!(advertised.cost, 5);
3679 announced.assert_next_wait();
3680 }
3681
3682 #[tokio::test]
3685 async fn test_completed_track_survives_route_churn() {
3686 tokio::time::pause();
3687
3688 let origin = Origin::random().produce();
3689 let consumer = origin.consume();
3690
3691 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3693 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3694
3695 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3696 let mut dynamic_a = source_a.dynamic();
3697 settle().await;
3698 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3699 let mut dynamic_b = source_b.dynamic();
3700 settle().await;
3701 settle().await;
3702 let broadcast = consumer.request_broadcast("test").await.unwrap();
3703
3704 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3706 let mut producer = accept_track(&mut dynamic_a, "video").await;
3707 settle().await;
3708 let mut sub = subscribing.await.unwrap();
3709 producer.append_group().unwrap();
3710 assert_eq!(sub.assert_group().sequence, 0);
3711 producer.finish().unwrap();
3712 drop(producer);
3713 settle().await;
3714 sub.assert_closed();
3715
3716 source_a.abort(Error::Dropped).unwrap();
3718 drop(dynamic_a);
3719 settle().await;
3720 dynamic_b.assert_no_request();
3721
3722 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3724 late.assert_closed();
3725 }
3726
3727 #[tokio::test]
3730 async fn test_serve_resets_retry_budget() {
3731 tokio::time::pause();
3732
3733 let origin = Origin::random().produce();
3734 let consumer = origin.consume();
3735
3736 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3737 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3738 let mut dynamic = source.dynamic();
3739 settle().await;
3740 settle().await;
3741 let broadcast = consumer.request_broadcast("test").await.unwrap();
3742
3743 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3745
3746 for _ in 0..2 * MAX_TRACK_RETRIES {
3749 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3750 .await
3751 .expect("timed out waiting for a retry")
3752 .unwrap();
3753 request.reject(Error::NotFound);
3754 let producer = accept_track(&mut dynamic, "video").await;
3755 settle().await;
3756 drop(producer);
3757 }
3758
3759 let _producer = accept_track(&mut dynamic, "video").await;
3760 settle().await;
3761 let mut sub = subscribing.await.unwrap();
3762 sub.assert_not_closed();
3763 }
3764
3765 #[tokio::test]
3769 async fn test_route_handover() {
3770 tokio::time::pause();
3771
3772 let origin = Origin::random().produce();
3773 let consumer = origin.consume();
3774 let mut announced = consumer.announced();
3775
3776 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3778 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3779
3780 let source_a = origin
3781 .create_broadcast("test", announce().with_hops(hops_long))
3782 .unwrap();
3783 let mut dynamic_a = source_a.dynamic();
3784 settle().await;
3785 settle().await;
3786 let broadcast = consumer.request_broadcast("test").await.unwrap();
3787 announced.assert_next_some("test");
3788
3789 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3790 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3791 settle().await;
3792 let mut sub = subscribing.await.unwrap();
3793 producer_a.append_group().unwrap();
3794 producer_a.append_group().unwrap();
3795 assert_eq!(sub.assert_group().sequence, 0);
3796 assert_eq!(sub.assert_group().sequence, 1);
3797
3798 let source_b = origin
3801 .create_broadcast("test", announce().with_hops(hops_short))
3802 .unwrap();
3803 let mut dynamic_b = source_b.dynamic();
3804 settle().await;
3805 settle().await;
3806 announced.assert_next_wait();
3807
3808 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3809 settle().await;
3810
3811 sub.assert_no_group();
3814 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3815 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3816
3817 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3819 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3820 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3821 assert_eq!(sub.assert_group().sequence, 2);
3822 assert_eq!(sub.assert_group().sequence, 3);
3823 sub.assert_no_group();
3824 sub.assert_not_closed();
3825 }
3826
3827 #[tokio::test(start_paused = true)]
3830 async fn test_route_unannounce_immediate() {
3831 let origin = Origin::random().produce();
3832 let consumer = origin.consume();
3833 let mut announced = consumer.announced();
3834
3835 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3836 let mut source = origin
3837 .create_broadcast("test", announce().with_hops(hops.clone()))
3838 .unwrap();
3839 settle().await;
3840 let broadcast = consumer.request_broadcast("test").await.unwrap();
3841 announced.assert_next_some("test");
3842
3843 source.finish();
3846 settle().await;
3847 announced.assert_next_none("test");
3848
3849 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3851 settle().await;
3852 let fresh = consumer.request_broadcast("test").await.unwrap();
3853 announced.assert_next_some("test");
3854 assert!(
3855 !fresh.is_clone(&broadcast),
3856 "re-create must not splice the old broadcast"
3857 );
3858 }
3859
3860 #[tokio::test(start_paused = true)]
3865 async fn test_route_detach_immediate() {
3866 let origin = Origin::random().produce();
3867 let consumer = origin.consume();
3868 let mut announced = consumer.announced();
3869
3870 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3871 let source = origin
3872 .create_broadcast("test", announce().with_hops(hops.clone()))
3873 .unwrap();
3874 let mut dynamic = source.dynamic();
3875 settle().await;
3876 settle().await;
3877 let broadcast = consumer.request_broadcast("test").await.unwrap();
3878 announced.assert_next_some("test");
3879
3880 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3881 let producer = accept_track(&mut dynamic, "video").await;
3882 settle().await;
3883 let mut sub = subscribing.await.unwrap();
3884
3885 drop(producer);
3887 source.abort(Error::Dropped).unwrap();
3888 drop(dynamic);
3889
3890 settle().await;
3891 announced.assert_next_none("test");
3892 sub.assert_error();
3893
3894 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3897 settle().await;
3898 settle().await;
3899 let fresh = consumer.request_broadcast("test").await.unwrap();
3900 announced.assert_next_some("test");
3901 assert!(
3902 !fresh.is_clone(&broadcast),
3903 "re-create must not splice the old broadcast"
3904 );
3905 }
3906
3907 #[tokio::test(start_paused = true)]
3912 async fn test_idle_track_releases_without_respinning() {
3913 let origin = Info::new(Origin::random()).produce();
3914 let consumer = origin.consume();
3915
3916 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3917 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3918 let mut dynamic = source.dynamic();
3919 settle().await;
3920 let broadcast = consumer.request_broadcast("test").await.unwrap();
3921
3922 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3923 let producer = accept_track(&mut dynamic, "video").await;
3924 settle().await;
3925 let sub = subscribing.await.unwrap();
3926
3927 drop(sub);
3930 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
3931 settle().await;
3932 assert!(
3933 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
3934 "the copy must stay spliced inside the linger",
3935 );
3936
3937 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3940 settle().await;
3941 assert!(
3942 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3943 "an idle copy must be released after the linger",
3944 );
3945
3946 for _ in 0..3 {
3950 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3951 settle().await;
3952 assert!(
3953 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3954 "an unread copy must stay released, not be re-spliced",
3955 );
3956 }
3957 assert!(
3958 dynamic.requested_track().now_or_never().is_none(),
3959 "an unread track must not be re-requested",
3960 );
3961 drop(producer);
3962
3963 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3965 let mut producer = accept_track(&mut dynamic, "video").await;
3966 settle().await;
3967 let mut sub = subscribing.await.unwrap();
3968 producer.append_group().unwrap();
3969 assert_eq!(sub.assert_group().sequence, 0);
3970 }
3971
3972 #[tokio::test(start_paused = true)]
3976 async fn test_back_to_back_fetches_reuse_the_track() {
3977 let origin = Info::new(Origin::random()).produce();
3978 let consumer = origin.consume();
3979
3980 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3981 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3982 let mut dynamic = source.dynamic();
3983 settle().await;
3984 let broadcast = consumer.request_broadcast("test").await.unwrap();
3985
3986 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3988 let mut producer = accept_track(&mut dynamic, "video").await;
3989 producer.append_group().unwrap().finish().unwrap();
3990 settle().await;
3991 let first = fetching.await.expect("first fetch");
3992 drop(first);
3993
3994 settle().await;
3996 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3997 settle().await;
3998 assert!(
3999 dynamic.requested_track().now_or_never().is_none(),
4000 "a fetch inside the linger must reuse the track, not re-request it",
4001 );
4002 drop(fetching.await.expect("second fetch"));
4003
4004 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4006 settle().await;
4007 assert!(
4008 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4009 "the copy must be released once the fetches stop",
4010 );
4011 drop(producer);
4012
4013 settle().await;
4015 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4016 let mut producer = accept_track(&mut dynamic, "video").await;
4017 producer.append_group().unwrap().finish().unwrap();
4018 settle().await;
4019 fetching.await.expect("fetch after the linger");
4020 }
4021
4022 #[tokio::test(start_paused = true)]
4026 async fn test_linger_reconnect_splices() {
4027 let origin = Info::new(Origin::random())
4028 .with_linger(Duration::from_secs(5))
4029 .produce();
4030 let consumer = origin.consume();
4031 let mut announced = consumer.announced();
4032
4033 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4034 let source = origin
4035 .create_broadcast("test", announce().with_hops(hops.clone()))
4036 .unwrap();
4037 let mut dynamic = source.dynamic();
4038 settle().await;
4039 settle().await;
4040 let broadcast = consumer.request_broadcast("test").await.unwrap();
4041 announced.assert_next_some("test");
4042
4043 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4044 let mut producer = accept_track(&mut dynamic, "video").await;
4045 settle().await;
4046 let mut sub = subscribing.await.unwrap();
4047
4048 producer.append_group().unwrap();
4049 producer.append_group().unwrap();
4050 assert_eq!(sub.assert_group().sequence, 0);
4051 assert_eq!(sub.assert_group().sequence, 1);
4052
4053 drop(producer);
4056 source.abort(Error::Dropped).unwrap();
4057 drop(dynamic);
4058 settle().await;
4059
4060 announced.assert_next_wait();
4062 sub.assert_no_group();
4063 sub.assert_not_closed();
4064
4065 let during = consumer.request_broadcast("test").await.unwrap();
4067 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4068
4069 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4072 let mut dynamic = source.dynamic();
4073 settle().await;
4074 settle().await;
4075 announced.assert_next_wait();
4076 let again = consumer.request_broadcast("test").await.unwrap();
4077 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4078
4079 let mut producer = accept_track(&mut dynamic, "video").await;
4083 settle().await;
4084 sub.assert_no_group();
4085 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
4086 producer.create_group(group::Info { sequence: 2 }).unwrap();
4087 assert_eq!(sub.assert_group().sequence, 2);
4088 sub.assert_not_closed();
4089 }
4090
4091 #[tokio::test(start_paused = true)]
4094 async fn test_linger_expiry_closes() {
4095 let origin = Info::new(Origin::random())
4096 .with_linger(Duration::from_secs(5))
4097 .produce();
4098 let consumer = origin.consume();
4099 let mut announced = consumer.announced();
4100
4101 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4102 let source = origin
4103 .create_broadcast("test", announce().with_hops(hops.clone()))
4104 .unwrap();
4105 let mut dynamic = source.dynamic();
4106 settle().await;
4107 settle().await;
4108 let broadcast = consumer.request_broadcast("test").await.unwrap();
4109 announced.assert_next_some("test");
4110
4111 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4112 let producer = accept_track(&mut dynamic, "video").await;
4113 settle().await;
4114 let mut sub = subscribing.await.unwrap();
4115
4116 drop(producer);
4117 source.abort(Error::Dropped).unwrap();
4118 drop(dynamic);
4119 settle().await;
4120 announced.assert_next_wait();
4121
4122 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4124 settle().await;
4125 announced.assert_next_none("test");
4126 sub.assert_error();
4127
4128 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4130 settle().await;
4131 settle().await;
4132 let fresh = consumer.request_broadcast("test").await.unwrap();
4133 announced.assert_next_some("test");
4134 assert!(
4135 !fresh.is_clone(&broadcast),
4136 "a late re-create must not splice the expired broadcast"
4137 );
4138 }
4139
4140 #[tokio::test(start_paused = true)]
4144 async fn test_linger_forever() {
4145 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4146 let consumer = origin.consume();
4147 let mut announced = consumer.announced();
4148
4149 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4150 let source = origin
4151 .create_broadcast("test", announce().with_hops(hops.clone()))
4152 .unwrap();
4153 settle().await;
4154 let broadcast = consumer.request_broadcast("test").await.unwrap();
4155 announced.assert_next_some("test");
4156
4157 source.abort(Error::Dropped).unwrap();
4158 settle().await;
4159
4160 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4162 announced.assert_next_wait();
4163 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4164 settle().await;
4165 settle().await;
4166 let again = consumer.request_broadcast("test").await.unwrap();
4167 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4168 drop(source);
4169 }
4170
4171 #[tokio::test(start_paused = true)]
4174 async fn test_linger_skipped_on_finish() {
4175 let origin = Info::new(Origin::random())
4176 .with_linger(Duration::from_secs(5))
4177 .produce();
4178 let consumer = origin.consume();
4179 let mut announced = consumer.announced();
4180
4181 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4182 let mut source = origin
4183 .create_broadcast("test", announce().with_hops(hops.clone()))
4184 .unwrap();
4185 settle().await;
4186 let broadcast = consumer.request_broadcast("test").await.unwrap();
4187 announced.assert_next_some("test");
4188
4189 source.finish();
4192 settle().await;
4193 announced.assert_next_none("test");
4194
4195 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4197 settle().await;
4198 let fresh = consumer.request_broadcast("test").await.unwrap();
4199 announced.assert_next_some("test");
4200 assert!(
4201 !fresh.is_clone(&broadcast),
4202 "a finish must not leave a lingering broadcast to splice into"
4203 );
4204 }
4205
4206 #[tokio::test]
4209 async fn test_announce_toggle() {
4210 tokio::time::pause();
4211
4212 let origin = Origin::random().produce();
4213 let consumer = origin.consume();
4214 let mut announced = consumer.announced();
4215
4216 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4217 settle().await;
4218
4219 announced.assert_next_wait();
4221 let broadcast = consumer
4222 .get_broadcast("test")
4223 .expect("offline broadcast is still routable");
4224 assert!(!broadcast.route().announce);
4225
4226 let requested = consumer.request_broadcast("test").await.unwrap();
4228 assert!(requested.is_clone(&broadcast));
4229
4230 source.set_route(announce()).unwrap();
4232 settle().await;
4233 let face = announced.assert_next_some("test");
4234 assert!(face.is_clone(&broadcast));
4235
4236 let mut fresh = origin.consume().announced();
4238 fresh.assert_next_some("test");
4239 fresh.assert_next_wait();
4240
4241 source.set_route(broadcast::Route::new()).unwrap();
4243 settle().await;
4244 announced.assert_next_none("test");
4245 assert!(consumer.get_broadcast("test").is_some());
4246 let mut fresh = origin.consume().announced();
4247 fresh.assert_next_wait();
4248
4249 source.finish();
4250 settle().await;
4251 assert!(consumer.get_broadcast("test").is_none());
4252 }
4253
4254 #[tokio::test]
4257 async fn test_announce_beats_offline() {
4258 tokio::time::pause();
4259
4260 let origin = Origin::random().produce();
4261 let consumer = origin.consume();
4262 let mut announced = consumer.announced();
4263
4264 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4266 settle().await;
4267 announced.assert_next_wait();
4268
4269 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
4272 settle().await;
4273 announced.assert_next_some("test");
4274 let face = consumer.get_broadcast("test").unwrap();
4275 assert!(face.route().announce);
4276 assert_eq!(face.route().cost, 10);
4277
4278 announced_source.finish();
4281 settle().await;
4282 announced.assert_next_none("test");
4283 assert!(consumer.get_broadcast("test").is_some());
4284 }
4285
4286 #[tokio::test]
4289 async fn test_better_source_no_churn() {
4290 tokio::time::pause();
4291
4292 let origin = Origin::random().produce();
4293 let mut announced = origin.consume().announced();
4294
4295 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4298 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4299 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4300 settle().await;
4301 let face = announced.assert_next_some("test");
4302
4303 let _b = origin
4304 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4305 .unwrap();
4306 settle().await;
4307 announced.assert_next_wait();
4308 let current = origin.consume().get_broadcast("test").unwrap();
4309 assert!(current.is_clone(&face), "the broadcast identity must not change");
4310 assert_eq!(current.route().hops, hops_b);
4312 }
4313
4314 #[tokio::test]
4320 async fn test_publisher_mismatch_replaces() {
4321 tokio::time::pause();
4322
4323 let origin = Origin::random().produce();
4324 let consumer = origin.consume();
4325 let mut announced = consumer.announced();
4326
4327 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4328 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4329
4330 let mut source_a = origin
4331 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4332 .unwrap();
4333 settle().await;
4334 let face_a = announced.assert_next_some("test");
4335
4336 let _source_b = origin
4339 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4340 .unwrap();
4341 settle().await;
4342 settle().await;
4343 announced.assert_next_none("test");
4344 let face_b = announced.assert_next_some("test");
4345 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
4346 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4347 assert!(face_a.is_closed(), "the displaced front must close");
4351
4352 source_a.finish();
4354 settle().await;
4355 settle().await;
4356 announced.assert_next_wait();
4357 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4358 }
4359
4360 #[tokio::test]
4365 async fn test_reconnect_wins_over_stale_route() {
4366 tokio::time::pause();
4367
4368 let origin = Origin::random().produce();
4369 let consumer = origin.consume();
4370
4371 let publisher = Origin::new(1).unwrap();
4372 let hops = OriginList::try_from(vec![publisher]).unwrap();
4373
4374 let stale = origin
4377 .create_broadcast("test", announce().with_hops(hops.clone()))
4378 .unwrap();
4379 let mut stale_dynamic = stale.dynamic();
4380 settle().await;
4381
4382 let fresh = origin
4384 .create_broadcast("test", announce().with_hops(hops.clone()))
4385 .unwrap();
4386 let mut fresh_dynamic = fresh.dynamic();
4387 settle().await;
4388 settle().await;
4389
4390 let broadcast = consumer.request_broadcast("test").await.unwrap();
4392 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4393 settle().await;
4394 let _producer = accept_track(&mut fresh_dynamic, "video").await;
4395 settle().await;
4396 subscribing.await.unwrap();
4397 stale_dynamic.assert_no_request();
4398 }
4399
4400 #[tokio::test]
4406 async fn test_carrying_reconnect_switches_immediately() {
4407 tokio::time::pause();
4408
4409 let origin = Origin::random().produce();
4410 let consumer = origin.consume();
4411
4412 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4413
4414 let stale = origin
4415 .create_broadcast("test", announce().with_hops(hops.clone()))
4416 .unwrap();
4417 let mut stale_dynamic = stale.dynamic();
4418 settle().await;
4419
4420 let broadcast = consumer.request_broadcast("test").await.unwrap();
4422 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4423 settle().await;
4424 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
4425 settle().await;
4426 let _subscription = subscribing.await.unwrap();
4428
4429 let fresh = origin
4431 .create_broadcast("test", announce().with_hops(hops.clone()))
4432 .unwrap();
4433 let mut fresh_dynamic = fresh.dynamic();
4434 settle().await;
4435 settle().await;
4436
4437 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
4440 }
4441
4442 #[tokio::test]
4452 async fn test_offline_mismatch_never_evicts_a_live_front() {
4453 tokio::time::pause();
4454
4455 let origin = Origin::random().produce();
4456 let consumer = origin.consume();
4457 let mut announced = consumer.announced();
4458
4459 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4460 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4461
4462 let mut live = origin
4463 .create_broadcast("test", announce().with_hops(hops_live.clone()))
4464 .unwrap();
4465 let mut live_dynamic = live.dynamic();
4466 settle().await;
4467 let face = announced.assert_next_some("test");
4468
4469 let broadcast = consumer.request_broadcast("test").await.unwrap();
4471 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4472 settle().await;
4473 let _producer = accept_track(&mut live_dynamic, "video").await;
4474 settle().await;
4475 subscribing.await.unwrap();
4476
4477 let cache = origin
4480 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
4481 .unwrap();
4482 settle().await;
4483 settle().await;
4484 announced.assert_next_wait();
4485 assert!(!face.is_closed(), "the live front must survive");
4486 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
4487
4488 live.finish();
4491 settle().await;
4492 settle().await;
4493 announced.assert_next_none("test");
4494 let taken = consumer
4495 .get_broadcast("test")
4496 .expect("the parked source must take over");
4497 assert_eq!(taken.route().hops, hops_cache);
4498 announced.assert_next_wait();
4500 drop(cache);
4501 }
4502
4503 #[tokio::test]
4509 async fn test_dispatch_excludes_requester() {
4510 tokio::time::pause();
4511
4512 let origin = Origin::random().produce();
4513 let consumer = origin.consume();
4514
4515 let peer = Origin::new(5).unwrap();
4516 let publisher = Origin::new(1).unwrap();
4517 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
4519 let clean = OriginList::try_from(vec![publisher]).unwrap();
4520
4521 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4522 let mut dynamic_a = source_a.dynamic();
4523 settle().await;
4524 let source_b = origin
4525 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
4526 .unwrap();
4527 let mut dynamic_b = source_b.dynamic();
4528 settle().await;
4529 settle().await;
4530
4531 let shared = consumer.request_broadcast("test").await.unwrap();
4534 let subscribing = shared.track("video").unwrap().subscribe(None);
4535 let _producer_a = accept_track(&mut dynamic_a, "video").await;
4536 settle().await;
4537 subscribing.await.unwrap();
4538
4539 let scoped = consumer.clone().excluding(peer);
4544 let pinned = scoped.request_broadcast("test").await.unwrap();
4545 let subscribing = pinned.track("video").unwrap().subscribe(None);
4546 let _producer_b = accept_track(&mut dynamic_b, "video").await;
4547 settle().await;
4548 subscribing.await.unwrap();
4549 dynamic_a.assert_no_request();
4550 }
4551
4552 #[tokio::test]
4558 async fn test_standby_join_splices_live_subscriber() {
4559 tokio::time::pause();
4560
4561 let origin = Origin::random().produce();
4562 let consumer = origin.consume();
4563
4564 let publisher = Origin::new(1).unwrap();
4565 let peer = Origin::new(5).unwrap();
4566 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4567 let local = OriginList::try_from(vec![publisher]).unwrap();
4568
4569 let source_remote = origin
4571 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4572 .unwrap();
4573 let mut dynamic_remote = source_remote.dynamic();
4574 settle().await;
4575 settle().await;
4576 let broadcast = consumer.request_broadcast("test").await.unwrap();
4577 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4578 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
4579 settle().await;
4580 let mut sub = subscribing.await.unwrap();
4581 producer_remote.append_group().unwrap();
4582 assert_eq!(sub.assert_group().sequence, 0);
4583
4584 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4587 let mut dynamic_local = source_local.dynamic();
4588 settle().await;
4589 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4590 settle().await;
4591 sub.assert_no_group();
4592 assert_eq!(producer_local.subscription().unwrap().group_start, Some(1));
4593 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
4594 assert_eq!(sub.assert_group().sequence, 1);
4595 sub.assert_not_closed();
4596 }
4597
4598 #[tokio::test]
4605 async fn test_standby_missing_track_keeps_incumbent() {
4606 tokio::time::pause();
4607
4608 let origin = Origin::random().produce();
4609 let consumer = origin.consume();
4610
4611 let publisher = Origin::new(1).unwrap();
4612 let peer = Origin::new(5).unwrap();
4613 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4614 let local = OriginList::try_from(vec![publisher]).unwrap();
4615
4616 let source_remote = origin
4618 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4619 .unwrap();
4620 let mut dynamic_remote = source_remote.dynamic();
4621 settle().await;
4622 settle().await;
4623 let broadcast = consumer.request_broadcast("test").await.unwrap();
4624 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4625 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
4626 settle().await;
4627 let mut sub = subscribing.await.unwrap();
4628 producer_remote.append_group().unwrap();
4629 assert_eq!(sub.assert_group().sequence, 0);
4630
4631 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4636 let mut dynamic_local = source_local.dynamic();
4637 settle().await;
4638 for _ in 0..2 * MAX_TRACK_RETRIES {
4639 if let Some(Ok(request)) = dynamic_local.requested_track().now_or_never() {
4640 assert_eq!(request.name(), "audio");
4641 request.reject(Error::NotFound);
4642 }
4643 settle().await;
4644 }
4645
4646 producer_remote.append_group().unwrap();
4648 assert_eq!(sub.assert_group().sequence, 1);
4649 sub.assert_not_closed();
4650
4651 source_remote.abort(Error::Dropped).unwrap();
4654 settle().await;
4655 settle().await;
4656 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
4657 settle().await;
4658 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
4659 assert_eq!(sub.assert_group().sequence, 2);
4660 sub.assert_not_closed();
4661 }
4662
4663 #[tokio::test]
4668 async fn test_unservable_track_retried_by_a_later_request() {
4669 tokio::time::pause();
4670
4671 let origin = Origin::random().produce();
4672 let consumer = origin.consume();
4673
4674 let source = origin.create_broadcast("test", announce()).unwrap();
4675 let mut dynamic = source.dynamic();
4676 settle().await;
4677 settle().await;
4678 let broadcast = consumer.request_broadcast("test").await.unwrap();
4679
4680 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4682 for _ in 0..MAX_TRACK_RETRIES {
4683 let request = dynamic.requested_track().await.unwrap();
4684 request.reject(Error::NotFound);
4685 settle().await;
4686 }
4687 assert!(matches!(subscribing.await, Err(Error::Unroutable)));
4688
4689 let retry = broadcast.track("audio").unwrap().subscribe(None);
4691 let mut producer = accept_track(&mut dynamic, "audio").await;
4692 settle().await;
4693 let mut sub = retry.await.expect("a fresh request must reach the source");
4694 producer.append_group().unwrap();
4695 assert_eq!(sub.assert_group().sequence, 0);
4696 }
4697
4698 #[tokio::test]
4704 async fn test_per_track_fallback_respects_exclusion() {
4705 tokio::time::pause();
4706
4707 let origin = Origin::random().produce();
4708 let consumer = origin.consume();
4709
4710 let publisher = Origin::new(1).unwrap();
4711 let peer = Origin::new(5).unwrap();
4712 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4713 let local = OriginList::try_from(vec![publisher]).unwrap();
4714
4715 let source_tainted = origin
4717 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4718 .unwrap();
4719 let mut dynamic_tainted = source_tainted.dynamic();
4720 settle().await;
4721 settle().await;
4722
4723 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4726 let mut dynamic_clean = source_clean.dynamic();
4727 settle().await;
4728
4729 let scoped = consumer.clone().excluding(peer);
4730 let broadcast = scoped.request_broadcast("test").await.unwrap();
4731 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
4732 settle().await;
4733
4734 for _ in 0..2 * MAX_TRACK_RETRIES {
4737 if let Some(Ok(request)) = dynamic_clean.requested_track().now_or_never() {
4738 request.reject(Error::NotFound);
4739 }
4740 settle().await;
4741 }
4742 dynamic_tainted.assert_no_request();
4743 }
4744
4745 #[tokio::test]
4749 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
4750 tokio::time::pause();
4751
4752 let origin = Origin::random().produce();
4753 let consumer = origin.consume();
4754
4755 let publisher = Origin::new(1).unwrap();
4756 let peer = Origin::new(5).unwrap();
4757 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4758 let local = OriginList::try_from(vec![publisher]).unwrap();
4759
4760 let source_tainted = origin
4761 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4762 .unwrap();
4763 let mut dynamic_tainted = source_tainted.dynamic();
4764 settle().await;
4765 settle().await;
4766 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4767 let mut dynamic_clean = source_clean.dynamic();
4768 settle().await;
4769
4770 let scoped = consumer.clone().excluding(peer);
4771 let broadcast = scoped.request_broadcast("test").await.unwrap();
4772 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
4773 let _clean = accept_track(&mut dynamic_clean, "video").await;
4774 settle().await;
4775
4776 source_clean.abort(Error::Dropped).unwrap();
4778 settle().await;
4779 settle().await;
4780 dynamic_tainted.assert_no_request();
4781
4782 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
4785 }
4786
4787 #[tokio::test]
4792 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
4793 tokio::time::pause();
4794
4795 let origin = Origin::random().produce();
4796 let consumer = origin.consume();
4797
4798 let publisher = Origin::new(1).unwrap();
4799 let peer = Origin::new(5).unwrap();
4800 let local = OriginList::try_from(vec![publisher]).unwrap();
4801 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4802
4803 let source_clean = origin
4807 .create_broadcast("test", announce().with_hops(local).with_cost(5))
4808 .unwrap();
4809 let mut dynamic_clean = source_clean.dynamic();
4810 settle().await;
4811 settle().await;
4812 let scoped = consumer.clone().excluding(peer);
4813 let broadcast = scoped.request_broadcast("test").await.unwrap();
4814 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4815 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
4816 settle().await;
4817 let mut sub = subscribing.await.unwrap();
4818 producer_clean.append_group().unwrap();
4819 assert_eq!(sub.assert_group().sequence, 0);
4820
4821 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
4828 tainted.advertised = 1;
4829 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
4830 let mut dynamic_tainted = source_tainted.dynamic();
4831 settle().await;
4832 settle().await;
4833 dynamic_tainted.assert_no_request();
4834 producer_clean.append_group().unwrap();
4835 assert_eq!(sub.assert_group().sequence, 1);
4836 sub.assert_not_closed();
4837
4838 drop(sub);
4841 drop(broadcast);
4842 drop(scoped);
4843 settle().await;
4844 let mut bumped = announce().with_hops(via_peer).with_cost(1);
4845 bumped.advertised = 1;
4846 source_tainted.set_route(bumped).unwrap();
4847 settle().await;
4848 let plain = consumer.request_broadcast("test").await.unwrap();
4849 let _plain_track = plain.track("video").unwrap().subscribe(None);
4850 settle().await;
4851 settle().await;
4852 assert!(
4853 dynamic_tainted.requested_track().now_or_never().is_some(),
4854 "the front must be free to use the route again once the peer is gone"
4855 );
4856 }
4857
4858 #[tokio::test]
4862 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
4863 tokio::time::pause();
4864
4865 let origin = Origin::random().produce();
4866 let consumer = origin.consume();
4867 let mut dynamic = origin.dynamic();
4868
4869 let peer = Origin::new(5).unwrap();
4870 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
4871 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4872 settle().await;
4873 settle().await;
4874
4875 let scoped = consumer.clone().excluding(peer);
4876 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
4877 assert!(
4878 dynamic.requested_broadcast().now_or_never().is_none(),
4879 "the dynamic handler was asked to route around the exclusion"
4880 );
4881
4882 let _pending = scoped.request_broadcast("other");
4884 settle().await;
4885 assert!(
4886 dynamic.requested_broadcast().now_or_never().is_some(),
4887 "a genuinely missing path must still fall back"
4888 );
4889 }
4890
4891 #[tokio::test]
4894 async fn test_dispatch_all_tainted_unroutable() {
4895 tokio::time::pause();
4896
4897 let origin = Origin::random().produce();
4898 let consumer = origin.consume();
4899
4900 let peer = Origin::new(5).unwrap();
4901 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
4902 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4903 settle().await;
4904 settle().await;
4905
4906 let scoped = consumer.clone().excluding(peer);
4907 match scoped.request_broadcast("test").await {
4908 Err(Error::Unroutable) => {}
4909 Err(err) => panic!("expected Unroutable, got {err:?}"),
4910 Ok(_) => panic!("expected Unroutable, got a broadcast"),
4911 }
4912
4913 consumer.request_broadcast("test").await.unwrap();
4915 }
4916
4917 #[tokio::test]
4918 async fn test_duplicate_reverse() {
4919 tokio::time::pause();
4920
4921 let origin = Origin::random().produce();
4922
4923 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4924 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4925 settle().await;
4926 assert!(origin.consume().get_broadcast("test").is_some());
4927
4928 broadcast2.finish();
4930 settle().await;
4931 assert!(origin.consume().get_broadcast("test").is_some());
4932
4933 broadcast1.finish();
4934 settle().await;
4935 assert!(origin.consume().get_broadcast("test").is_none());
4936 }
4937
4938 #[tokio::test]
4939 async fn test_deterministic_tiebreak() {
4940 tokio::time::pause();
4941
4942 fn hops(ids: &[u64]) -> OriginList {
4943 OriginList::try_from(
4944 ids.iter()
4945 .copied()
4946 .map(|id| Origin::new(id).unwrap())
4947 .collect::<Vec<_>>(),
4948 )
4949 .unwrap()
4950 }
4951
4952 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
4955 let origin = Origin::random().produce();
4956 let _a = origin
4957 .create_broadcast("test", announce().with_hops(hops(first)))
4958 .unwrap();
4959 let _b = origin
4960 .create_broadcast("test", announce().with_hops(hops(second)))
4961 .unwrap();
4962 settle().await;
4963 origin.consume().get_broadcast("test").unwrap().route().hops
4964 }
4965
4966 let forward = winner(&[5, 20], &[5, 40]).await;
4970 let reverse = winner(&[5, 40], &[5, 20]).await;
4971 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
4972
4973 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
4975 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
4976 }
4977
4978 #[tokio::test]
4983 async fn test_many_announces() {
4984 let origin = Origin::random().produce();
4985
4986 let mut consumer = origin.consume().announced();
4987 let mut broadcasts = Vec::new();
4989 for i in 0..256 {
4990 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
4991 settle().await;
4992 }
4993
4994 for i in 0..256 {
4995 consumer.assert_next_some(format!("test{i:03}"));
4996 }
4997 consumer.assert_next_wait();
4998 }
4999
5000 #[tokio::test]
5001 async fn test_many_announces_try() {
5002 let origin = Origin::random().produce();
5003
5004 let mut consumer = origin.consume().announced();
5005 let mut broadcasts = Vec::new();
5007 for i in 0..256 {
5008 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5009 settle().await;
5010 }
5011
5012 for i in 0..256 {
5013 consumer.assert_try_next_some(format!("test{i:03}"));
5014 }
5015 }
5016
5017 #[tokio::test]
5018 async fn test_with_root_basic() {
5019 let origin = Origin::random().produce();
5020
5021 let foo_producer = origin.with_root("foo").expect("should create root");
5023 assert_eq!(foo_producer.root().as_str(), "foo");
5024
5025 let mut consumer = origin.consume().announced();
5026
5027 let _broadcast = foo_producer
5029 .create_broadcast("bar/baz", announce())
5030 .expect("publish allowed");
5031 settle().await;
5032 consumer.assert_next_some("foo/bar/baz");
5034
5035 let mut foo_consumer = foo_producer.consume().announced();
5037 foo_consumer.assert_next_some("bar/baz");
5038 }
5039
5040 #[tokio::test]
5041 async fn test_with_root_nested() {
5042 let origin = Origin::random().produce();
5043
5044 let foo_producer = origin.with_root("foo").expect("should create foo root");
5046 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
5047 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
5048
5049 let mut consumer = origin.consume().announced();
5050
5051 let _broadcast = foo_bar_producer
5053 .create_broadcast("baz", announce())
5054 .expect("publish allowed");
5055 settle().await;
5056 consumer.assert_next_some("foo/bar/baz");
5058
5059 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
5061 foo_bar_consumer.assert_next_some("baz");
5062 }
5063
5064 #[tokio::test]
5065 async fn test_publish_scope_allows() {
5066 let origin = Origin::random().produce();
5067
5068 let limited_producer = origin
5070 .scope(&["allowed/path1".into(), "allowed/path2".into()])
5071 .expect("should create limited producer");
5072
5073 let _broadcast = limited_producer
5075 .create_broadcast("allowed/path1", announce())
5076 .expect("publish allowed");
5077 let _keep2 = limited_producer
5078 .create_broadcast("allowed/path1/nested", announce())
5079 .expect("publish allowed");
5080 let _keep3 = limited_producer
5081 .create_broadcast("allowed/path2", announce())
5082 .expect("publish allowed");
5083 settle().await;
5084
5085 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
5087 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
5089 }
5090
5091 #[tokio::test]
5092 async fn test_publish_max_parts() {
5093 let origin = Origin::random().produce();
5094
5095 let at_limit = (0..Path::MAX_PARTS)
5096 .map(|i| i.to_string())
5097 .collect::<Vec<_>>()
5098 .join("/");
5099 let _broadcast = origin
5100 .create_broadcast(at_limit.as_str(), announce())
5101 .expect("publish allowed");
5102 settle().await;
5103
5104 let too_deep = format!("{at_limit}/extra");
5105 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
5106
5107 let rooted = origin.with_root("root").expect("wildcard allows any root");
5109 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
5110 }
5111
5112 #[tokio::test]
5113 async fn test_publish_scope_empty() {
5114 let origin = Origin::random().produce();
5115
5116 assert!(origin.scope(&[]).is_none());
5118 }
5119
5120 #[tokio::test]
5121 async fn test_consume_scope_filters() {
5122 let origin = Origin::random().produce();
5123
5124 let mut consumer = origin.consume().announced();
5125
5126 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
5128 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
5129 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
5130 settle().await;
5131
5132 let mut limited_consumer = origin
5134 .consume()
5135 .scope(&["allowed".into()])
5136 .expect("should create limited consumer")
5137 .announced();
5138
5139 limited_consumer.assert_next_some("allowed");
5141 limited_consumer.assert_next_some("allowed/nested");
5142 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
5146 consumer.assert_next_some("allowed/nested");
5147 consumer.assert_next_some("notallowed");
5148 }
5149
5150 #[tokio::test]
5151 async fn test_consume_scope_multiple_prefixes() {
5152 let origin = Origin::random().produce();
5153
5154 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5155 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5156 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5157 settle().await;
5158
5159 let mut limited_consumer = origin
5161 .consume()
5162 .scope(&["foo".into(), "bar".into()])
5163 .expect("should create limited consumer")
5164 .announced();
5165
5166 limited_consumer.assert_next_some("bar/test");
5168 limited_consumer.assert_next_some("foo/test");
5169 limited_consumer.assert_next_wait(); }
5171
5172 #[tokio::test]
5173 async fn test_with_root_and_publish_scope() {
5174 let origin = Origin::random().produce();
5175
5176 let foo_producer = origin.with_root("foo").expect("should create foo root");
5178
5179 let limited_producer = foo_producer
5181 .scope(&["bar".into(), "goop/pee".into()])
5182 .expect("should create limited producer");
5183
5184 let mut consumer = origin.consume().announced();
5185
5186 let _broadcast = limited_producer
5188 .create_broadcast("bar", announce())
5189 .expect("publish allowed");
5190 let _keep2 = limited_producer
5191 .create_broadcast("bar/nested", announce())
5192 .expect("publish allowed");
5193 let _keep3 = limited_producer
5194 .create_broadcast("goop/pee", announce())
5195 .expect("publish allowed");
5196 let _keep4 = limited_producer
5197 .create_broadcast("goop/pee/nested", announce())
5198 .expect("publish allowed");
5199 settle().await;
5200
5201 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
5203 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
5205
5206 consumer.assert_next_some("foo/bar");
5208 consumer.assert_next_some("foo/bar/nested");
5209 consumer.assert_next_some("foo/goop/pee");
5210 consumer.assert_next_some("foo/goop/pee/nested");
5211 }
5212
5213 #[tokio::test]
5214 async fn test_with_root_and_consume_scope() {
5215 let origin = Origin::random().produce();
5216
5217 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
5219 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
5220 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
5221 settle().await;
5222
5223 let foo_producer = origin.with_root("foo").expect("should create foo root");
5225
5226 let mut limited_consumer = foo_producer
5228 .consume()
5229 .scope(&["bar".into(), "goop/pee".into()])
5230 .expect("should create limited consumer")
5231 .announced();
5232
5233 limited_consumer.assert_next_some("bar/test");
5235 limited_consumer.assert_next_some("goop/pee/test");
5236 limited_consumer.assert_next_wait(); }
5238
5239 #[tokio::test]
5240 async fn test_with_root_unauthorized() {
5241 let origin = Origin::random().produce();
5242
5243 let limited_producer = origin
5245 .scope(&["allowed".into()])
5246 .expect("should create limited producer");
5247
5248 assert!(limited_producer.with_root("notallowed").is_none());
5250
5251 let allowed_root = limited_producer
5253 .with_root("allowed")
5254 .expect("should create allowed root");
5255 assert_eq!(allowed_root.root().as_str(), "allowed");
5256 }
5257
5258 #[tokio::test]
5259 async fn test_wildcard_permission() {
5260 let origin = Origin::random().produce();
5261
5262 let root_producer = origin.clone();
5264
5265 let _broadcast = root_producer
5267 .create_broadcast("any/path", announce())
5268 .expect("publish allowed");
5269 let _keep2 = root_producer
5270 .create_broadcast("other/path", announce())
5271 .expect("publish allowed");
5272 settle().await;
5273
5274 let foo_producer = root_producer.with_root("foo").expect("should create any root");
5276 assert_eq!(foo_producer.root().as_str(), "foo");
5277 }
5278
5279 #[tokio::test]
5280 async fn test_consume_broadcast_with_permissions() {
5281 let origin = Origin::random().produce();
5282
5283 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
5284 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
5285 settle().await;
5286
5287 let limited_consumer = origin
5289 .consume()
5290 .scope(&["allowed".into()])
5291 .expect("should create limited consumer");
5292
5293 let result = limited_consumer.get_broadcast("allowed/test");
5295 assert!(result.is_some());
5296 assert!(
5297 result
5298 .unwrap()
5299 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
5300 );
5301
5302 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
5304
5305 let consumer = origin.consume();
5307 assert!(consumer.get_broadcast("allowed/test").is_some());
5308 assert!(consumer.get_broadcast("notallowed/test").is_some());
5309 }
5310
5311 #[tokio::test]
5312 async fn test_nested_paths_with_permissions() {
5313 let origin = Origin::random().produce();
5314
5315 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
5317
5318 let _broadcast = limited_producer
5320 .create_broadcast("a/b/c", announce())
5321 .expect("publish allowed");
5322 let _keep2 = limited_producer
5323 .create_broadcast("a/b/c/d", announce())
5324 .expect("publish allowed");
5325 let _keep3 = limited_producer
5326 .create_broadcast("a/b/c/d/e", announce())
5327 .expect("publish allowed");
5328 settle().await;
5329
5330 assert!(limited_producer.create_broadcast("a", announce()).is_err());
5332 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
5333 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
5334 }
5335
5336 #[tokio::test]
5337 async fn test_multiple_consumers_with_different_permissions() {
5338 let origin = Origin::random().produce();
5339
5340 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5342 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5343 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5344 settle().await;
5345
5346 let mut foo_consumer = origin
5348 .consume()
5349 .scope(&["foo".into()])
5350 .expect("should create foo consumer")
5351 .announced();
5352
5353 let mut bar_consumer = origin
5354 .consume()
5355 .scope(&["bar".into()])
5356 .expect("should create bar consumer")
5357 .announced();
5358
5359 let mut foobar_consumer = origin
5360 .consume()
5361 .scope(&["foo".into(), "bar".into()])
5362 .expect("should create foobar consumer")
5363 .announced();
5364
5365 foo_consumer.assert_next_some("foo/test");
5367 foo_consumer.assert_next_wait();
5368
5369 bar_consumer.assert_next_some("bar/test");
5370 bar_consumer.assert_next_wait();
5371
5372 foobar_consumer.assert_next_some("bar/test");
5373 foobar_consumer.assert_next_some("foo/test");
5374 foobar_consumer.assert_next_wait();
5375 }
5376
5377 #[tokio::test]
5378 async fn test_select_with_empty_prefix() {
5379 let origin = Origin::random().produce();
5380
5381 let demo_producer = origin.with_root("demo").expect("should create demo root");
5383 let limited_producer = demo_producer
5384 .scope(&["worm-node".into(), "foobar".into()])
5385 .expect("should create limited producer");
5386
5387 let _broadcast1 = limited_producer
5389 .create_broadcast("worm-node/test", announce())
5390 .expect("publish allowed");
5391 let _broadcast2 = limited_producer
5392 .create_broadcast("foobar/test", announce())
5393 .expect("publish allowed");
5394 settle().await;
5395
5396 let mut consumer = limited_producer
5398 .consume()
5399 .scope(&["".into()])
5400 .expect("should create consumer with empty prefix")
5401 .announced();
5402
5403 let a1 = consumer.try_next().expect("expected first announcement");
5405 let a2 = consumer.try_next().expect("expected second announcement");
5406 consumer.assert_next_wait();
5407
5408 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5409 paths.sort();
5410 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
5411 }
5412
5413 #[tokio::test]
5414 async fn test_select_narrowing_scope() {
5415 let origin = Origin::random().produce();
5416
5417 let demo_producer = origin.with_root("demo").expect("should create demo root");
5419 let limited_producer = demo_producer
5420 .scope(&["worm-node".into(), "foobar".into()])
5421 .expect("should create limited producer");
5422
5423 let _broadcast1 = limited_producer
5425 .create_broadcast("worm-node", announce())
5426 .expect("publish allowed");
5427 let _broadcast2 = limited_producer
5428 .create_broadcast("worm-node/foo", announce())
5429 .expect("publish allowed");
5430 let _broadcast3 = limited_producer
5431 .create_broadcast("foobar/bar", announce())
5432 .expect("publish allowed");
5433 settle().await;
5434
5435 let mut worm_consumer = limited_producer
5437 .consume()
5438 .scope(&["worm-node".into()])
5439 .expect("should create worm-node consumer")
5440 .announced();
5441
5442 worm_consumer.assert_next_some("worm-node");
5444 worm_consumer.assert_next_some("worm-node/foo");
5445 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
5449 .consume()
5450 .scope(&["worm-node/foo".into()])
5451 .expect("should create worm-node/foo consumer")
5452 .announced();
5453
5454 foo_consumer.assert_next_some("worm-node/foo");
5455 foo_consumer.assert_next_wait(); }
5457
5458 #[tokio::test]
5459 async fn test_select_multiple_roots_with_empty_prefix() {
5460 let origin = Origin::random().produce();
5461
5462 let limited_producer = origin
5464 .scope(&["app1".into(), "app2".into(), "shared".into()])
5465 .expect("should create limited producer");
5466
5467 let _broadcast1 = limited_producer
5469 .create_broadcast("app1/data", announce())
5470 .expect("publish allowed");
5471 let _broadcast2 = limited_producer
5472 .create_broadcast("app2/config", announce())
5473 .expect("publish allowed");
5474 let _broadcast3 = limited_producer
5475 .create_broadcast("shared/resource", announce())
5476 .expect("publish allowed");
5477 settle().await;
5478
5479 let mut consumer = limited_producer
5481 .consume()
5482 .scope(&["".into()])
5483 .expect("should create consumer with empty prefix")
5484 .announced();
5485
5486 consumer.assert_next_some("app1/data");
5488 consumer.assert_next_some("app2/config");
5489 consumer.assert_next_some("shared/resource");
5490 consumer.assert_next_wait();
5491 }
5492
5493 #[tokio::test]
5494 async fn test_publish_scope_with_empty_prefix() {
5495 let origin = Origin::random().produce();
5496
5497 let limited_producer = origin
5499 .scope(&["services/api".into(), "services/web".into()])
5500 .expect("should create limited producer");
5501
5502 let same_producer = limited_producer
5504 .scope(&["".into()])
5505 .expect("should create producer with empty prefix");
5506
5507 let _broadcast = same_producer
5509 .create_broadcast("services/api", announce())
5510 .expect("publish allowed");
5511 let _keep2 = same_producer
5512 .create_broadcast("services/web", announce())
5513 .expect("publish allowed");
5514 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
5515 assert!(same_producer.create_broadcast("other", announce()).is_err());
5516 }
5517
5518 #[tokio::test]
5519 async fn test_select_narrowing_to_deeper_path() {
5520 let origin = Origin::random().produce();
5521
5522 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
5524
5525 let _broadcast1 = limited_producer
5527 .create_broadcast("org/team1/project1", announce())
5528 .expect("publish allowed");
5529 let _broadcast2 = limited_producer
5530 .create_broadcast("org/team1/project2", announce())
5531 .expect("publish allowed");
5532 let _broadcast3 = limited_producer
5533 .create_broadcast("org/team2/project1", announce())
5534 .expect("publish allowed");
5535 settle().await;
5536
5537 let mut team2_consumer = limited_producer
5539 .consume()
5540 .scope(&["org/team2".into()])
5541 .expect("should create team2 consumer")
5542 .announced();
5543
5544 team2_consumer.assert_next_some("org/team2/project1");
5545 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
5549 .consume()
5550 .scope(&["org/team1/project1".into()])
5551 .expect("should create project1 consumer")
5552 .announced();
5553
5554 project1_consumer.assert_next_some("org/team1/project1");
5556 project1_consumer.assert_next_wait();
5557 }
5558
5559 #[tokio::test]
5560 async fn test_select_with_non_matching_prefix() {
5561 let origin = Origin::random().produce();
5562
5563 let limited_producer = origin
5565 .scope(&["allowed/path".into()])
5566 .expect("should create limited producer");
5567
5568 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
5570
5571 assert!(limited_producer.scope(&["other/path".into()]).is_none());
5573 }
5574
5575 #[tokio::test]
5578 async fn test_with_root_trailing_slash_consumer() {
5579 let origin = Origin::random().produce();
5580
5581 let prefix = "some_prefix/".to_string();
5583 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5584
5585 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5586 settle().await;
5587 consumer.assert_next_some("test");
5588 }
5589
5590 #[tokio::test]
5592 async fn test_with_root_trailing_slash_producer() {
5593 let origin = Origin::random().produce();
5594
5595 let prefix = "some_prefix/".to_string();
5597 let rooted = origin.with_root(prefix).unwrap();
5598
5599 let _b = rooted.create_broadcast("test", announce()).unwrap();
5600 settle().await;
5601
5602 let mut consumer = rooted.consume().announced();
5603 consumer.assert_next_some("test");
5604 }
5605
5606 #[tokio::test]
5608 async fn test_with_root_trailing_slash_unannounce() {
5609 tokio::time::pause();
5610
5611 let origin = Origin::random().produce();
5612
5613 let prefix = "some_prefix/".to_string();
5614 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5615
5616 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5617 settle().await;
5618 consumer.assert_next_some("test");
5619
5620 b.finish();
5622 settle().await;
5623
5624 consumer.assert_next_none("test");
5626 }
5627
5628 #[tokio::test]
5629 async fn test_select_maintains_access_with_wider_prefix() {
5630 let origin = Origin::random().produce();
5631
5632 let demo_producer = origin.with_root("demo").expect("should create demo root");
5634 let user_producer = demo_producer
5635 .scope(&["worm-node".into(), "foobar".into()])
5636 .expect("should create user producer");
5637
5638 let _broadcast1 = user_producer
5640 .create_broadcast("worm-node/data", announce())
5641 .expect("publish allowed");
5642 let _broadcast2 = user_producer
5643 .create_broadcast("foobar", announce())
5644 .expect("publish allowed");
5645 settle().await;
5646
5647 let mut consumer = user_producer
5649 .consume()
5650 .scope(&["".into()])
5651 .expect("scope with empty prefix should not fail when user has specific permissions")
5652 .announced();
5653
5654 let a1 = consumer.try_next().expect("expected first announcement");
5656 let a2 = consumer.try_next().expect("expected second announcement");
5657 consumer.assert_next_wait();
5658
5659 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5660 paths.sort();
5661 assert_eq!(paths, ["foobar", "worm-node/data"]);
5662
5663 let mut narrow_consumer = user_producer
5665 .consume()
5666 .scope(&["worm-node".into()])
5667 .expect("should be able to narrow scope to worm-node")
5668 .announced();
5669
5670 narrow_consumer.assert_next_some("worm-node/data");
5671 narrow_consumer.assert_next_wait(); }
5673
5674 #[tokio::test]
5675 async fn test_duplicate_prefixes_deduped() {
5676 let origin = Origin::random().produce();
5677
5678 let producer = origin
5680 .scope(&["demo".into(), "demo".into()])
5681 .expect("should create producer");
5682
5683 let _broadcast = producer
5684 .create_broadcast("demo/stream", announce())
5685 .expect("publish allowed");
5686 settle().await;
5687
5688 let mut consumer = producer.consume().announced();
5689 consumer.assert_next_some("demo/stream");
5690 consumer.assert_next_wait();
5691 }
5692
5693 #[tokio::test]
5694 async fn test_overlapping_prefixes_deduped() {
5695 let origin = Origin::random().produce();
5696
5697 let producer = origin
5699 .scope(&["demo".into(), "demo/foo".into()])
5700 .expect("should create producer");
5701
5702 let _broadcast = producer
5704 .create_broadcast("demo/bar/stream", announce())
5705 .expect("publish allowed");
5706 settle().await;
5707
5708 let mut consumer = producer.consume().announced();
5709 consumer.assert_next_some("demo/bar/stream");
5710 consumer.assert_next_wait();
5711 }
5712
5713 #[tokio::test]
5714 async fn test_overlapping_prefixes_no_duplicate_announcements() {
5715 let origin = Origin::random().produce();
5716
5717 let producer = origin
5719 .scope(&["demo".into(), "demo/foo".into()])
5720 .expect("should create producer");
5721
5722 let _broadcast = producer
5723 .create_broadcast("demo/foo/stream", announce())
5724 .expect("publish allowed");
5725 settle().await;
5726
5727 let mut consumer = producer.consume().announced();
5728 consumer.assert_next_some("demo/foo/stream");
5730 consumer.assert_next_wait();
5731 }
5732
5733 #[tokio::test]
5734 async fn test_allowed_returns_deduped_prefixes() {
5735 let origin = Origin::random().produce();
5736
5737 let producer = origin
5738 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
5739 .expect("should create producer");
5740
5741 let allowed: Vec<_> = producer.allowed().collect();
5742 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
5743 }
5744
5745 #[tokio::test]
5746 async fn test_announced_broadcast_already_announced() {
5747 let origin = Origin::random().produce();
5748
5749 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
5750 settle().await;
5751
5752 let consumer = origin.consume();
5753 let result = consumer.announced_broadcast("test").await.expect("should find it");
5754 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
5755 }
5756
5757 #[tokio::test]
5758 async fn test_announced_broadcast_delayed() {
5759 tokio::time::pause();
5760
5761 let origin = Origin::random().produce();
5762
5763 let consumer = origin.consume();
5764
5765 let wait = tokio::spawn({
5767 let consumer = consumer.clone();
5768 async move { consumer.announced_broadcast("test").await }
5769 });
5770
5771 tokio::task::yield_now().await;
5773
5774 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
5775 settle().await;
5776
5777 let result = wait.await.unwrap().expect("should find it");
5778 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
5779 }
5780
5781 #[tokio::test]
5782 async fn test_announced_broadcast_ignores_unrelated_paths() {
5783 tokio::time::pause();
5784
5785 let origin = Origin::random().produce();
5786
5787 let consumer = origin.consume();
5788
5789 let wait = tokio::spawn({
5790 let consumer = consumer.clone();
5791 async move { consumer.announced_broadcast("target").await }
5792 });
5793
5794 tokio::task::yield_now().await;
5795
5796 let _other = origin.create_broadcast("other", announce()).unwrap();
5798 settle().await;
5799 tokio::task::yield_now().await;
5800 assert!(!wait.is_finished(), "must not resolve on unrelated path");
5801
5802 let _target = origin.create_broadcast("target", announce()).unwrap();
5803 settle().await;
5804 let result = wait.await.unwrap().expect("should find target");
5805 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
5806 }
5807
5808 #[tokio::test]
5809 async fn test_announced_broadcast_skips_nested_paths() {
5810 tokio::time::pause();
5811
5812 let origin = Origin::random().produce();
5813
5814 let consumer = origin.consume();
5815
5816 let wait = tokio::spawn({
5817 let consumer = consumer.clone();
5818 async move { consumer.announced_broadcast("foo").await }
5819 });
5820
5821 tokio::task::yield_now().await;
5822
5823 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
5825 settle().await;
5826 tokio::task::yield_now().await;
5827 assert!(!wait.is_finished(), "must not resolve on a nested path");
5828
5829 let _exact = origin.create_broadcast("foo", announce()).unwrap();
5830 settle().await;
5831 let result = wait.await.unwrap().expect("should find foo exactly");
5832 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
5833 }
5834
5835 #[tokio::test]
5836 async fn test_announced_broadcast_disallowed() {
5837 let origin = Origin::random().produce();
5838 let limited = origin
5839 .consume()
5840 .scope(&["allowed".into()])
5841 .expect("should create limited");
5842
5843 assert!(limited.announced_broadcast("notallowed").await.is_none());
5845 }
5846
5847 #[tokio::test]
5848 async fn test_announced_broadcast_scope_too_narrow() {
5849 let origin = Origin::random().produce();
5852 let limited = origin
5853 .consume()
5854 .scope(&["foo/specific".into()])
5855 .expect("should create limited");
5856
5857 let result = limited
5859 .announced_broadcast("foo")
5860 .now_or_never()
5861 .expect("must not block");
5862 assert!(result.is_none());
5863 }
5864
5865 #[tokio::test]
5869 async fn test_coalesce_announce_then_unannounce() {
5870 tokio::time::pause();
5872
5873 let origin = Origin::random().produce();
5874 let mut announced = origin.consume().announced();
5875
5876 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
5877 settle().await;
5878 broadcast.finish();
5879
5880 settle().await;
5881
5882 announced.assert_next_wait();
5883 }
5884
5885 #[tokio::test]
5886 async fn test_coalesce_announce_unannounce_announce() {
5887 tokio::time::pause();
5890
5891 let origin = Origin::random().produce();
5892 let mut announced = origin.consume().announced();
5893
5894 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5895 settle().await;
5896 broadcast1.finish();
5897 settle().await;
5898 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5899 settle().await;
5900
5901 announced.assert_next_some("test");
5902 announced.assert_next_wait();
5903 }
5904
5905 #[tokio::test]
5906 async fn test_coalesce_unannounce_announce_preserved() {
5907 tokio::time::pause();
5910
5911 let origin = Origin::random().produce();
5912 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5913 settle().await;
5914
5915 let mut announced = origin.consume().announced();
5916 announced.assert_next_some("test");
5917
5918 broadcast1.finish();
5920 settle().await;
5921
5922 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5923 settle().await;
5924
5925 announced.assert_next_none("test");
5927 announced.assert_next_some("test");
5928 announced.assert_next_wait();
5929 }
5930
5931 #[tokio::test]
5932 async fn test_coalesce_unannounce_announce_unannounce() {
5933 tokio::time::pause();
5936
5937 let origin = Origin::random().produce();
5938 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5939 settle().await;
5940
5941 let mut announced = origin.consume().announced();
5942 announced.assert_next_some("test");
5943
5944 broadcast1.finish();
5945 settle().await;
5946
5947 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5948 settle().await;
5949 broadcast2.finish();
5950 settle().await;
5951
5952 announced.assert_next_none("test");
5953 announced.assert_next_wait();
5954 }
5955
5956 #[tokio::test]
5957 async fn test_coalesce_churn_bounded() {
5958 tokio::time::pause();
5963
5964 let origin = Origin::random().produce();
5965 let mut announced = origin.consume().announced();
5966
5967 for _ in 0..1000 {
5968 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
5969 settle().await;
5970 broadcast.finish();
5971 }
5972 settle().await;
5973
5974 let mut collected = Vec::new();
5975 while let Some(update) = announced.try_next() {
5976 collected.push(update);
5977 }
5978 assert!(
5979 collected.len() <= 1,
5980 "expected at most one pending update, got {}",
5981 collected.len()
5982 );
5983 assert!(
5984 collected.iter().all(|a| a.path == Path::new("test")),
5985 "unexpected path in pending updates",
5986 );
5987 }
5988
5989 #[tokio::test]
5993 async fn test_consumer_clone_is_side_effect_free() {
5994 let origin = Origin::random().produce();
5995
5996 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
5997 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
5998 settle().await;
5999
6000 let consumer = origin.consume();
6001 let mut announced = consumer.announced();
6002
6003 for _ in 0..16 {
6006 let cloned = consumer.clone();
6007 assert!(cloned.get_broadcast("test1").is_some());
6008 assert!(cloned.get_broadcast("test2").is_some());
6009 }
6010
6011 let a1 = announced.try_next().expect("first announcement");
6014 let a2 = announced.try_next().expect("second announcement");
6015 announced.assert_next_wait();
6016
6017 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6018 paths.sort();
6019 assert_eq!(paths, ["test1", "test2"]);
6020
6021 let mut fresh = consumer.announced();
6023 let b1 = fresh.try_next().expect("backlog: first");
6024 let b2 = fresh.try_next().expect("backlog: second");
6025 fresh.assert_next_wait();
6026
6027 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
6028 paths.sort();
6029 assert_eq!(paths, ["test1", "test2"]);
6030 }
6031
6032 #[tokio::test]
6034 async fn dynamic_request_unroutable_without_handler() {
6035 let origin = Origin::random().produce();
6036 let consumer = origin.consume();
6037 assert!(matches!(
6038 consumer.request_broadcast("missing").await,
6039 Err(Error::Unroutable)
6040 ));
6041 }
6042
6043 #[tokio::test(start_paused = true)]
6046 async fn dynamic_request_served_not_announced() {
6047 let origin = Origin::random().produce();
6048 let mut dynamic = origin.dynamic();
6049 let consumer = origin.consume();
6050
6051 let mut announced = origin.consume().announced();
6053 announced.assert_next_wait();
6054
6055 let served = broadcast::Info::new().produce();
6056 let request_fut = consumer.request_broadcast("fallback");
6059
6060 let mut served_dynamic = served.dynamic();
6062
6063 let request = dynamic.requested_broadcast().await.unwrap();
6064 assert_eq!(request.path(), &Path::new("fallback"));
6065 request.accept(&served);
6066
6067 let broadcast = request_fut.await.unwrap();
6068 assert!(broadcast.is_clone(&served.consume()));
6069
6070 let track_fut = broadcast.track("video").unwrap().subscribe(None);
6072 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
6073 let mut track = track_fut.await.unwrap();
6074 producer.append_group().unwrap();
6075 track.assert_group();
6076
6077 announced.assert_next_wait();
6079 }
6080
6081 #[tokio::test(start_paused = true)]
6083 async fn dynamic_request_coalesces() {
6084 let origin = Origin::random().produce();
6085 let mut dynamic = origin.dynamic();
6086 let consumer = origin.consume();
6087
6088 let f1 = consumer.request_broadcast("dup");
6090 let f2 = consumer.request_broadcast("dup");
6091
6092 let request = dynamic.requested_broadcast().await.unwrap();
6094 assert_eq!(request.path(), &Path::new("dup"));
6095 assert!(
6096 dynamic.requested_broadcast().now_or_never().is_none(),
6097 "a coalesced request must not be served twice"
6098 );
6099
6100 let served = broadcast::Info::new().produce();
6102 request.accept(&served);
6103 assert!(f1.await.unwrap().is_clone(&served.consume()));
6104 assert!(f2.await.unwrap().is_clone(&served.consume()));
6105 }
6106
6107 #[tokio::test(start_paused = true)]
6110 async fn dynamic_request_dedups_served() {
6111 let origin = Origin::random().produce();
6112 let mut dynamic = origin.dynamic();
6113 let consumer = origin.consume();
6114
6115 let request_fut = consumer.request_broadcast("fallback");
6116 let request = dynamic.requested_broadcast().await.unwrap();
6117 let served = broadcast::Info::new().produce();
6118 request.accept(&served);
6119 let first = request_fut.await.unwrap();
6120 assert!(first.is_clone(&served.consume()));
6121
6122 let second = consumer.request_broadcast("fallback").await.unwrap();
6124 assert!(second.is_clone(&served.consume()));
6125
6126 assert!(
6128 dynamic.requested_broadcast().now_or_never().is_none(),
6129 "a still-live served broadcast must not be re-requested from the handler"
6130 );
6131 }
6132
6133 #[tokio::test(start_paused = true)]
6135 async fn dynamic_request_reserves_after_close() {
6136 let origin = Origin::random().produce();
6137 let mut dynamic = origin.dynamic();
6138 let consumer = origin.consume();
6139
6140 let request_fut = consumer.request_broadcast("fallback");
6141 let request = dynamic.requested_broadcast().await.unwrap();
6142 let served = broadcast::Info::new().produce();
6143 request.accept(&served);
6144 request_fut.await.unwrap();
6145
6146 drop(served);
6148
6149 let request_fut = consumer.request_broadcast("fallback");
6151 let request = dynamic.requested_broadcast().await.unwrap();
6152 assert_eq!(request.path(), &Path::new("fallback"));
6153 let served = broadcast::Info::new().produce();
6154 request.accept(&served);
6155 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6156 }
6157
6158 #[tokio::test(start_paused = true)]
6161 async fn dynamic_request_served_cache_bounded() {
6162 let origin = Origin::random().produce();
6163 let mut dynamic = origin.dynamic();
6164 let consumer = origin.consume();
6165
6166 for i in 0..100 {
6167 let path = format!("one-shot/{i}");
6168 let request_fut = consumer.request_broadcast(&path);
6169 let request = dynamic.requested_broadcast().await.unwrap();
6170 let served = broadcast::Info::new().produce();
6171 request.accept(&served);
6172 request_fut.await.unwrap();
6173 drop(served);
6175 }
6176
6177 assert!(
6180 origin.dynamic.read().served.len() <= 4,
6181 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
6182 origin.dynamic.read().served.len()
6183 );
6184 }
6185
6186 #[tokio::test(start_paused = true)]
6189 async fn dynamic_request_coalesces_after_handoff() {
6190 let origin = Origin::random().produce();
6191 let mut dynamic = origin.dynamic();
6192 let consumer = origin.consume();
6193
6194 let f1 = consumer.request_broadcast("fallback");
6195 let request = dynamic.requested_broadcast().await.unwrap();
6197
6198 let f2 = consumer.request_broadcast("fallback");
6200 assert!(
6201 dynamic.requested_broadcast().now_or_never().is_none(),
6202 "a repeat request during hand-off must coalesce, not re-queue"
6203 );
6204
6205 let served = broadcast::Info::new().produce();
6207 request.accept(&served);
6208 assert!(f1.await.unwrap().is_clone(&served.consume()));
6209 assert!(f2.await.unwrap().is_clone(&served.consume()));
6210 }
6211
6212 #[tokio::test(start_paused = true)]
6214 async fn dynamic_request_dropped_after_handoff() {
6215 let origin = Origin::random().produce();
6216 let mut dynamic = origin.dynamic();
6217 let consumer = origin.consume();
6218
6219 let f1 = consumer.request_broadcast("fallback");
6220 let request = dynamic.requested_broadcast().await.unwrap();
6221 let f2 = consumer.request_broadcast("fallback");
6222
6223 drop(request);
6225 assert!(matches!(f1.await, Err(Error::Unroutable)));
6226 assert!(matches!(f2.await, Err(Error::Unroutable)));
6227 }
6228
6229 #[tokio::test(start_paused = true)]
6231 async fn dynamic_request_rejected() {
6232 let origin = Origin::random().produce();
6233 let mut dynamic = origin.dynamic();
6234 let consumer = origin.consume();
6235
6236 let request_fut = consumer.request_broadcast("fallback");
6237
6238 let request = dynamic.requested_broadcast().await.unwrap();
6239 request.reject(Error::Cancel);
6240
6241 assert!(matches!(request_fut.await, Err(Error::Cancel)));
6242 }
6243
6244 #[tokio::test(start_paused = true)]
6248 async fn dynamic_request_rerequest_after_reject() {
6249 let origin = Origin::random().produce();
6250 let mut dynamic = origin.dynamic();
6251 let consumer = origin.consume();
6252
6253 let f1 = consumer.request_broadcast("fallback");
6254 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
6255 assert!(matches!(f1.await, Err(Error::Unroutable)));
6256
6257 let served = broadcast::Info::new().produce();
6258 let f2 = consumer.request_broadcast("fallback");
6260 let request = dynamic.requested_broadcast().await.unwrap();
6261 assert_eq!(request.path(), &Path::new("fallback"));
6262 request.accept(&served);
6263 assert!(f2.await.unwrap().is_clone(&served.consume()));
6264 }
6265
6266 #[tokio::test(start_paused = true)]
6269 async fn dynamic_request_handler_dropped() {
6270 let origin = Origin::random().produce();
6271 let dynamic = origin.dynamic();
6272 let consumer = origin.consume();
6273
6274 let request_fut = consumer.request_broadcast("fallback");
6275 drop(dynamic);
6276 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
6277
6278 assert!(matches!(
6280 consumer.request_broadcast("again").await,
6281 Err(Error::Unroutable)
6282 ));
6283 }
6284
6285 #[tokio::test(start_paused = true)]
6289 async fn dynamic_request_accept_after_handler_dropped() {
6290 let origin = Origin::random().produce();
6291 let mut dynamic = origin.dynamic();
6292 let consumer = origin.consume();
6293
6294 let request_fut = consumer.request_broadcast("fallback");
6295
6296 let request = dynamic.requested_broadcast().await.unwrap();
6298 drop(dynamic);
6299
6300 let served = broadcast::Info::new().produce();
6301 request.accept(&served);
6303 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6304 }
6305
6306 #[tokio::test(start_paused = true)]
6308 async fn dynamic_request_prefers_announced() {
6309 let origin = Origin::random().produce();
6310 let mut dynamic = origin.dynamic();
6311 let consumer = origin.consume();
6312
6313 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
6314 settle().await;
6315
6316 let got = consumer.request_broadcast("live").await.unwrap();
6317 assert!(
6318 got.is_clone(&consumer.get_broadcast("live").unwrap()),
6319 "should return the published broadcast"
6320 );
6321 assert!(
6322 dynamic.requested_broadcast().now_or_never().is_none(),
6323 "a published path must not queue a fallback request"
6324 );
6325 }
6326
6327 #[tokio::test(start_paused = true)]
6329 async fn dynamic_clone_keeps_alive() {
6330 let origin = Origin::random().produce();
6331 let dynamic = origin.dynamic();
6332 let consumer = origin.consume();
6333
6334 drop(dynamic.clone());
6335
6336 let request_fut = consumer.request_broadcast("fallback");
6339 assert!(
6340 request_fut.now_or_never().is_none(),
6341 "request should stay pending until served"
6342 );
6343 }
6344}