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();
1919 idle_since = match (resume.is_spliced(), used) {
1920 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
1921 _ => None,
1922 };
1923 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
1924
1925 let step = {
1926 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
1927 kio::wait(|waiter| {
1928 match state.poll(waiter, |s| {
1934 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
1935 if s.closed
1936 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
1937 {
1938 Poll::Ready(())
1939 } else {
1940 Poll::Pending
1941 }
1942 }) {
1943 Poll::Ready(Ok(guard)) => {
1944 if guard.closed {
1945 return Poll::Ready(Step::Closed);
1946 }
1947 let Some(next) = guard.serve_route(skip) else {
1948 return Poll::Ready(Step::Resweep);
1949 };
1950 let source = guard
1951 .routes
1952 .iter()
1953 .find(|r| r.id == next)
1954 .expect("servable source in table")
1955 .source
1956 .clone();
1957 return Poll::Ready(Step::Splice(next, source));
1958 }
1959 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1960 Poll::Pending => {}
1961 }
1962
1963 let edge = match used {
1968 true => resume.poll_unused(waiter),
1969 false => resume.poll_used(waiter),
1970 };
1971 if edge.is_ready() {
1972 return Poll::Ready(Step::Demand);
1973 }
1974
1975 if let Some((_, track)) = &serving
1978 && let Poll::Ready(result) = track.poll_complete(waiter)
1979 {
1980 return Poll::Ready(match result {
1981 Ok(()) => Step::Complete,
1982 Err(_) => Step::Failed,
1983 });
1984 }
1985
1986 deadline.poll(waiter).map(|_| Step::Idle)
1987 })
1988 .await
1989 };
1990
1991 match step {
1992 Step::Closed => return,
1994 Step::Complete => {
1995 let _ = resume.finish();
1996 return;
1997 }
1998 Step::Failed => {
1999 serving = None;
2002 }
2003 Step::Demand => {}
2005 Step::Resweep => serving = None,
2014 Step::Idle => {
2015 if resume.release().is_err() {
2020 return;
2022 }
2023 serving = None;
2024 }
2025 Step::Splice(id, source) => {
2026 let attempt = match source.track(&name) {
2030 Ok(track) => {
2031 let query = track.info().into_inner();
2034 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2035 let info = kio::wait(|waiter| {
2036 if let Poll::Ready(result) = query.poll(waiter) {
2037 return Poll::Ready(Some(result));
2038 }
2039 match state.poll(waiter, |s| {
2040 if s.closed || s.serve_route(skip) != Some(id) {
2041 Poll::Ready(())
2042 } else {
2043 Poll::Pending
2044 }
2045 }) {
2046 Poll::Ready(_) => Poll::Ready(None),
2047 Poll::Pending => Poll::Pending,
2048 }
2049 })
2050 .await;
2051 match info {
2052 None => continue,
2055 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2059 Poll::Ready(Err(err)) => Err(err),
2060 _ => Ok(track),
2061 },
2062 Some(Err(err)) => Err(err),
2063 }
2064 }
2065 Err(err) => Err(err),
2066 };
2067
2068 match attempt {
2069 Ok(track) => {
2070 if resume.takeover(&track).is_err() {
2071 return;
2074 }
2075 fails = 0;
2080 dead.clear();
2081 serving = Some((id, track));
2082 }
2083 Err(_) if source.is_closing() => {
2087 dead.insert(id);
2088 serving = None;
2089 }
2090 Err(err) => {
2094 tracing::debug!(name = %name, source = id, %err, "source refused track");
2095 refused.insert(id);
2096 serving = None;
2097 }
2098 }
2099 }
2100 }
2101 }
2102}
2103
2104#[derive(Default)]
2110struct OriginDynamicState {
2111 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2114
2115 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2121}
2122
2123#[derive(Default)]
2130struct PendingBroadcast {
2131 resolved: Option<Result<broadcast::Consumer, Error>>,
2132}
2133
2134pub struct Dynamic {
2145 info: Origin,
2146 root: PathOwned,
2147 state: kio::Shared<OriginDynamicState>,
2148}
2149
2150impl Clone for Dynamic {
2151 fn clone(&self) -> Self {
2152 self.state.lock().requests.add_handler();
2156
2157 Self {
2158 info: self.info,
2159 root: self.root.clone(),
2160 state: self.state.clone(),
2161 }
2162 }
2163}
2164
2165impl Dynamic {
2166 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2167 state.lock().requests.add_handler();
2168
2169 Self { info, root, state }
2170 }
2171
2172 pub fn info(&self) -> &Origin {
2174 &self.info
2175 }
2176
2177 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2179 let mut state = ready!(self.state.poll(waiter, |state| {
2180 if state.requests.has_queued() {
2181 Poll::Ready(())
2182 } else {
2183 Poll::Pending
2184 }
2185 }));
2186
2187 let path = state.requests.pop().expect("predicate guaranteed a request");
2188 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2194 Poll::Ready(Ok(Request {
2195 path,
2196 producer,
2197 state: self.state.clone(),
2198 }))
2199 }
2200
2201 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2204 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2205 }
2206
2207 pub fn root(&self) -> &Path<'_> {
2209 &self.root
2210 }
2211}
2212
2213impl Drop for Dynamic {
2214 fn drop(&mut self) {
2215 let mut state = self.state.lock();
2218 if state.requests.remove_handler() {
2219 state.requests.drain_queued();
2223 }
2224 }
2225}
2226
2227pub struct Request {
2234 path: PathOwned,
2236
2237 producer: kio::Producer<PendingBroadcast>,
2240
2241 state: kio::Shared<OriginDynamicState>,
2243}
2244
2245impl Request {
2246 pub fn path(&self) -> &Path<'_> {
2248 &self.path
2249 }
2250
2251 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2257 let broadcast = broadcast.consume();
2258
2259 let resolved = {
2265 let mut state = self.state.lock();
2266 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2267 state
2268 .requests
2269 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2270 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2271 };
2272
2273 if let Ok(mut pending) = self.producer.write() {
2274 pending.resolved = Some(Ok(resolved));
2275 }
2276 }
2278
2279 pub fn reject(self, err: Error) {
2281 self.state
2282 .lock()
2283 .requests
2284 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2285 if let Ok(mut state) = self.producer.write() {
2286 state.resolved = Some(Err(err));
2287 }
2288 }
2289}
2290
2291impl Drop for Request {
2292 fn drop(&mut self) {
2293 self.state
2301 .lock()
2302 .requests
2303 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2304 }
2305}
2306
2307pub struct Requesting {
2314 inner: RequestState,
2315 stats: stats::Scope,
2318}
2319
2320enum RequestState {
2321 Ready(broadcast::Consumer),
2323 Failed(Error),
2326 Pending(kio::Consumer<PendingBroadcast>),
2328}
2329
2330impl Requesting {
2331 fn ready(broadcast: broadcast::Consumer) -> Self {
2332 Self {
2333 inner: RequestState::Ready(broadcast),
2334 stats: stats::Scope::default(),
2335 }
2336 }
2337
2338 fn failed(error: Error) -> Self {
2339 Self {
2340 inner: RequestState::Failed(error),
2341 stats: stats::Scope::default(),
2342 }
2343 }
2344
2345 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2346 Self {
2347 inner: RequestState::Pending(consumer),
2348 stats: stats::Scope::default(),
2349 }
2350 }
2351
2352 fn with_stats(mut self, scope: stats::Scope) -> Self {
2353 self.stats = scope;
2354 self
2355 }
2356
2357 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2359 match &self.inner {
2360 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2361 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2362 RequestState::Pending(consumer) => Poll::Ready(
2363 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2364 Some(result) => Poll::Ready(result.clone()),
2365 None => Poll::Pending,
2366 })) {
2367 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2368 Err(_closed) => Err(Error::Unroutable),
2370 },
2371 ),
2372 }
2373 }
2374}
2375
2376impl kio::Pollable for Requesting {
2377 type Output = Result<broadcast::Consumer, Error>;
2378
2379 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2380 self.poll_ok(waiter)
2381 }
2382}
2383
2384pub trait Consume<T> {
2392 fn consume(&self) -> T;
2394}
2395
2396impl<T, U: Consume<T>> Consume<T> for &U {
2397 fn consume(&self) -> T {
2398 (**self).consume()
2399 }
2400}
2401
2402impl Consume<Consumer> for Producer {
2403 fn consume(&self) -> Consumer {
2404 Consumer::new(
2408 self.info,
2409 self.root.clone(),
2410 self.nodes.clone(),
2411 self.dynamic.clone(),
2412 stats::Session::default(),
2413 )
2414 }
2415}
2416
2417impl Consume<Consumer> for Consumer {
2418 fn consume(&self) -> Consumer {
2419 self.clone()
2420 }
2421}
2422
2423impl Consume<broadcast::Consumer> for broadcast::Producer {
2424 fn consume(&self) -> broadcast::Consumer {
2425 self.consume()
2427 }
2428}
2429
2430impl Consume<broadcast::Consumer> for broadcast::Consumer {
2431 fn consume(&self) -> broadcast::Consumer {
2432 self.clone()
2433 }
2434}
2435
2436impl Consume<track::Consumer> for track::Producer {
2437 fn consume(&self) -> track::Consumer {
2438 self.consume()
2439 }
2440}
2441
2442impl Consume<track::Consumer> for track::Consumer {
2443 fn consume(&self) -> track::Consumer {
2444 self.clone()
2445 }
2446}
2447
2448#[derive(Clone)]
2454pub struct Consumer {
2455 info: Origin,
2457 nodes: OriginNodes,
2458
2459 root: PathOwned,
2461
2462 dynamic: kio::Shared<OriginDynamicState>,
2465
2466 stats: stats::Session,
2470
2471 exclude: Option<Origin>,
2475}
2476
2477impl std::ops::Deref for Consumer {
2478 type Target = Origin;
2479
2480 fn deref(&self) -> &Self::Target {
2481 &self.info
2482 }
2483}
2484
2485impl Consumer {
2486 fn new(
2487 info: Origin,
2488 root: PathOwned,
2489 nodes: OriginNodes,
2490 dynamic: kio::Shared<OriginDynamicState>,
2491 stats: stats::Session,
2492 ) -> Self {
2493 Self {
2494 info,
2495 nodes,
2496 root,
2497 dynamic,
2498 stats,
2499 exclude: None,
2500 }
2501 }
2502
2503 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2508 self.exclude = Some(peer);
2509 self
2510 }
2511
2512 pub fn with_stats(mut self, session: stats::Session) -> Self {
2516 self.stats = session;
2517 self
2518 }
2519
2520 fn untagged(&self) -> Self {
2524 Self {
2525 stats: stats::Session::default(),
2526 ..self.clone()
2527 }
2528 }
2529
2530 pub(crate) fn empty(&self) -> Self {
2535 Self {
2536 info: self.info,
2537 nodes: OriginNodes { nodes: Vec::new() },
2538 root: self.root.clone(),
2539 dynamic: self.dynamic.clone(),
2540 stats: self.stats.clone(),
2541 exclude: self.exclude,
2542 }
2543 }
2544
2545 pub fn announced(&self) -> AnnounceConsumer {
2552 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2553 }
2554
2555 pub fn consume(&self) -> Self {
2557 self.clone()
2558 }
2559
2560 fn resolve(&self, path: impl AsPath) -> Resolved {
2569 let path = path.as_path();
2570 let Some((root, rest)) = self.nodes.get(&path) else {
2571 return Resolved::Missing;
2572 };
2573 let state = root.lock();
2574 state.resolve_broadcast(&rest, self.exclude)
2575 }
2576
2577 #[cfg(test)]
2579 fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2580 match self.resolve(path) {
2581 Resolved::Found(broadcast) => Some(broadcast),
2582 Resolved::Excluded | Resolved::Missing => None,
2583 }
2584 }
2585
2586 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2598 let path = path.as_path();
2599
2600 let consumer = self.scope(std::slice::from_ref(&path))?;
2602
2603 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2607 return None;
2608 }
2609
2610 let mut announced = consumer.untagged().announced();
2614 let scope = self.stats.egress(self.root.join(&path).to_owned());
2615 loop {
2616 let OriginAnnounce {
2617 path: announced_path,
2618 broadcast,
2619 } = announced.next().await?;
2620 if announced_path.as_path() == path
2622 && let Some(broadcast) = broadcast
2623 {
2624 return Some(broadcast.with_stats(scope));
2625 }
2626 }
2627 }
2628
2629 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2635 let prefixes = PathPrefixes::new(prefixes);
2636 Some(Consumer {
2637 info: self.info,
2638 root: self.root.clone(),
2639 nodes: self.nodes.select(&prefixes)?,
2640 dynamic: self.dynamic.clone(),
2641 stats: self.stats.clone(),
2642 exclude: self.exclude,
2643 })
2644 }
2645
2646 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2665 let path = path.as_path();
2666
2667 let absolute = self.root.join(&path).to_owned();
2671 let scope = self.stats.egress(&absolute);
2672
2673 match self.resolve(&path) {
2679 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
2680 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
2681 Resolved::Missing => {}
2682 }
2683
2684 let mut state = self.dynamic.lock();
2685
2686 if let Some(weak) = state.served.get(&absolute) {
2690 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2691 }
2692
2693 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2696 producer.consume()
2697 } else {
2698 let producer = kio::Producer::<PendingBroadcast>::default();
2699 let consumer = producer.consume();
2700 if state.requests.insert(absolute, producer).is_err() {
2701 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2702 }
2703 consumer
2704 };
2705
2706 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2707 }
2708
2709 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2714 let prefix = prefix.as_path();
2715
2716 Some(Self {
2717 info: self.info,
2718 root: self.root.join(&prefix).to_owned(),
2719 nodes: self.nodes.root(&prefix)?,
2720 dynamic: self.dynamic.clone(),
2721 stats: self.stats.clone(),
2722 exclude: self.exclude,
2723 })
2724 }
2725
2726 pub fn root(&self) -> &Path<'_> {
2728 &self.root
2729 }
2730
2731 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2734 self.nodes.nodes.iter().map(|(root, _)| root)
2735 }
2736
2737 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2739 self.root.join(path)
2740 }
2741}
2742
2743#[derive(Clone)]
2748pub struct AnnounceProducer {
2749 nodes: OriginNodes,
2750 root: PathOwned,
2751}
2752
2753impl AnnounceProducer {
2754 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2755 Self { nodes, root }
2756 }
2757
2758 pub fn consume(&self) -> AnnounceConsumer {
2764 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2767 }
2768
2769 pub fn root(&self) -> &Path<'_> {
2771 &self.root
2772 }
2773}
2774
2775pub struct AnnounceConsumer {
2780 id: ConsumerId,
2781 nodes: OriginNodes,
2782 root: PathOwned,
2783
2784 state: kio::Producer<OriginConsumerState>,
2787
2788 stats: stats::Session,
2791
2792 guards: HashMap<PathOwned, stats::Announce>,
2796}
2797
2798impl AnnounceConsumer {
2799 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2800 let state = kio::Producer::<OriginConsumerState>::default();
2801 let id = ConsumerId::new();
2802
2803 for (_, node) in &nodes.nodes {
2804 let notify = AnnounceConsumerNotify {
2805 root: root.clone(),
2806 state: state.clone(),
2807 };
2808 node.lock().consume(id, notify);
2809 }
2810
2811 Self {
2812 id,
2813 nodes,
2814 root,
2815 state,
2816 stats,
2817 guards: HashMap::new(),
2818 }
2819 }
2820
2821 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2827 let OriginAnnounce { path, broadcast } = update;
2828 let absolute = self.root.join(&path).to_owned();
2829 match broadcast {
2830 Some(broadcast) => {
2831 let scope = self.stats.egress(&absolute);
2832 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2833 OriginAnnounce {
2834 path,
2835 broadcast: Some(broadcast.with_stats(scope)),
2836 }
2837 }
2838 None => {
2839 self.guards.remove(&absolute);
2840 OriginAnnounce { path, broadcast: None }
2841 }
2842 }
2843 }
2844
2845 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2852 kio::wait(|waiter| self.poll_next(waiter)).await
2853 }
2854
2855 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2861 let update = {
2862 let mut state = match ready!(self.state.poll(waiter, |state| {
2863 if state.pending.is_empty() {
2864 Poll::Pending
2865 } else {
2866 Poll::Ready(())
2867 }
2868 })) {
2869 Ok(state) => state,
2870 Err(_) => return Poll::Ready(None),
2872 };
2873 state.take().expect("predicate guaranteed an update")
2874 };
2875 Poll::Ready(Some(self.attribute(update)))
2876 }
2877
2878 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2883 let update = self.state.write().ok()?.take()?;
2884 Some(self.attribute(update))
2885 }
2886
2887 pub fn is_closed(&self) -> bool {
2889 self.state.write().is_err()
2890 }
2891
2892 pub fn root(&self) -> &Path<'_> {
2894 &self.root
2895 }
2896
2897 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2899 self.root.join(path)
2900 }
2901}
2902
2903impl Drop for AnnounceConsumer {
2904 fn drop(&mut self) {
2905 for (_, root) in &self.nodes.nodes {
2906 root.lock().unconsume(self.id);
2907 }
2908 }
2909}
2910
2911#[cfg(test)]
2912use futures::FutureExt;
2913
2914#[cfg(test)]
2915#[allow(missing_docs)] impl AnnounceConsumer {
2917 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2918 let expected = expected.as_path();
2919 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2920 assert_eq!(announce.path, expected, "wrong path");
2921 let announced = announce.broadcast.expect("should be an active announce");
2922 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2923 }
2924
2925 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2929 let expected = expected.as_path();
2930 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2931 assert_eq!(announce.path, expected, "wrong path");
2932 announce.broadcast.expect("should be an active announce")
2933 }
2934
2935 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2936 let expected = expected.as_path();
2937 let announce = self.try_next().expect("no next");
2938 assert_eq!(announce.path, expected, "wrong path");
2939 let announced = announce.broadcast.expect("should be an active announce");
2940 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2941 }
2942
2943 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2945 let expected = expected.as_path();
2946 let announce = self.try_next().expect("no next");
2947 assert_eq!(announce.path, expected, "wrong path");
2948 announce.broadcast.expect("should be an active announce")
2949 }
2950
2951 pub fn assert_next_none(&mut self, expected: impl AsPath) {
2952 let expected = expected.as_path();
2953 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2954 assert_eq!(announce.path, expected, "wrong path");
2955 assert!(announce.broadcast.is_none(), "should be unannounced");
2956 }
2957
2958 pub fn assert_next_wait(&mut self) {
2959 if let Some(res) = self.next().now_or_never() {
2960 panic!("next should block: got {:?}", res.map(|a| a.path));
2961 }
2962 }
2963
2964 }
2973
2974#[cfg(test)]
2975mod tests {
2976 use crate::coding::Decode;
2977 use crate::group;
2978
2979 use super::*;
2980
2981 fn announce() -> broadcast::Route {
2983 broadcast::Route::new().with_announce(true)
2984 }
2985
2986 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
2992 let name = Path::new(name);
2993 let peer_key = fnv_key(&name, [peer]);
2994 (100u64..)
2995 .map(|id| Origin::new(id).unwrap())
2996 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
2997 .unwrap()
2998 }
2999
3000 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
3003 let source = broadcast::Info::new().produce().consume();
3004 FrontState {
3005 path: Path::new("test").to_owned(),
3006 self_origin,
3007 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
3008 next_route: routes.len() as u64,
3009 excluded: HashMap::new(),
3010 routes: routes
3011 .into_iter()
3012 .enumerate()
3013 .map(|(id, route)| FrontRoute {
3014 id: id as u64,
3015 route,
3016 source: source.clone(),
3017 })
3018 .collect(),
3019 active: Some(0),
3020 linger: Duration::ZERO,
3021 closed: false,
3022 }
3023 }
3024
3025 fn sibling_route(peer: Origin) -> broadcast::Route {
3028 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3029 announce().with_hops(hops)
3030 }
3031
3032 fn upstream_route(cost: u64) -> broadcast::Route {
3034 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3035 announce().with_hops(hops).with_cost(cost)
3036 }
3037
3038 #[test]
3042 fn test_carrying_gate_keys() {
3043 let peer = Origin::new(3).unwrap();
3044
3045 let mut lost = front_state(
3047 origin_keyed("test", peer, false),
3048 vec![upstream_route(10), sibling_route(peer)],
3049 );
3050 lost.reselect(true);
3051 assert_eq!(
3052 lost.active,
3053 Some(0),
3054 "carrying front re-parented onto a higher-keyed peer"
3055 );
3056 lost.reselect(false);
3057 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3058
3059 let mut won = front_state(
3061 origin_keyed("test", peer, true),
3062 vec![upstream_route(10), sibling_route(peer)],
3063 );
3064 won.reselect(true);
3065 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3066 }
3067
3068 #[test]
3073 fn test_carrying_gate_symmetric_race() {
3074 let a = Origin::new(1).unwrap();
3075 let b = Origin::new(2).unwrap();
3076
3077 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3078 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3079 a_view.reselect(true);
3080 b_view.reselect(true);
3081
3082 let a_moved = a_view.active == Some(1);
3083 let b_moved = b_view.active == Some(1);
3084 assert!(
3085 a_moved != b_moved,
3086 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3087 );
3088 }
3089
3090 #[test]
3095 fn test_carrying_switches_to_benign_routes() {
3096 let peer = Origin::new(3).unwrap();
3097 let lost = origin_keyed("test", peer, false);
3098
3099 let mut forwarder = sibling_route(peer).with_cost(4);
3101 forwarder.advertised = 4;
3102 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3103 state.reselect(true);
3104 assert_eq!(
3105 state.active,
3106 Some(1),
3107 "a cheaper forwarder path must win while carrying"
3108 );
3109
3110 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3112 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3113 state.reselect(true);
3114 assert_eq!(
3115 state.active,
3116 Some(1),
3117 "a direct publisher route must win while carrying"
3118 );
3119
3120 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3125 state.reselect(true);
3126 assert_eq!(
3127 state.active,
3128 Some(1),
3129 "a reconnect on an identical chain must win while carrying"
3130 );
3131 }
3132
3133 #[test]
3136 fn test_carrying_gate_ignores_unannounced_incumbent() {
3137 let peer = Origin::new(3).unwrap();
3138 let unannounced = upstream_route(10).with_announce(false);
3139 let mut state = front_state(
3140 origin_keyed("test", peer, false),
3141 vec![unannounced, sibling_route(peer)],
3142 );
3143 state.reselect(true);
3144 assert_eq!(
3145 state.active,
3146 Some(1),
3147 "an unannounced incumbent must always be displaced"
3148 );
3149 }
3150
3151 async fn settle() {
3154 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3155 }
3156
3157 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3160 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3161 .await
3162 .expect("timed out waiting for a track request")
3163 .expect("source closed");
3164 assert_eq!(request.name(), name, "unexpected track dispatched");
3165 request.accept(None)
3166 }
3167
3168 #[tokio::test]
3172 async fn test_stats_tagged_end_to_end() {
3173 use crate::Timestamp;
3174 use crate::stats::{Config, Registry, Tier};
3175 use bytes::Bytes;
3176
3177 tokio::time::pause();
3178
3179 let registry = Registry::new(Config::new());
3180 let ctx = registry.tier(Tier::default()).session("acme");
3181
3182 let origin = Origin::random().produce();
3183 let ingress = origin.clone().with_stats(ctx.clone());
3184 let egress = origin.consume().with_stats(ctx.clone());
3185
3186 let mut announced = egress.announced();
3189
3190 let source = ingress.create_broadcast("demo", announce()).unwrap();
3192 let mut dynamic = source.dynamic();
3193 settle().await;
3194 settle().await;
3195
3196 let update = announced.next().await.unwrap();
3198 assert_eq!(update.path.as_str(), "demo");
3199 let broadcast = update.broadcast.unwrap();
3200
3201 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3203 let mut producer = accept_track(&mut dynamic, "video").await;
3204 settle().await;
3205 let mut sub = subscribing.await.unwrap();
3206
3207 let mut group = producer.append_group().unwrap();
3209 group
3210 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3211 .unwrap();
3212 group
3213 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3214 .unwrap();
3215 group.finish().unwrap();
3216
3217 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3219 let mut frames = 0;
3220 while let Some(frame) = group_c.read_frame().await.unwrap() {
3221 assert_eq!(frame.payload.len(), 5);
3222 frames += 1;
3223 }
3224 assert_eq!(frames, 2);
3225 settle().await;
3226
3227 let report = registry.report();
3228 let entry = report
3229 .traffic
3230 .iter()
3231 .find(|e| e.path.as_str() == "demo")
3232 .expect("demo tracked");
3233 let path_len = "demo".len() as u64;
3234
3235 let egress = &entry.publisher;
3237 assert_eq!(egress.announced, 1, "one egress announce");
3238 assert_eq!(egress.announced_bytes, path_len);
3239 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3240 assert_eq!(egress.broadcasts, 1, "one viewer");
3241 assert_eq!(egress.groups, 1);
3242 assert_eq!(egress.frames, 2);
3243 assert_eq!(egress.bytes, 10);
3244 assert_eq!(egress.fetches, 0);
3245
3246 let ingress = &entry.subscriber;
3248 assert_eq!(ingress.announced, 1, "one ingress announce");
3249 assert_eq!(ingress.announced_bytes, path_len);
3250 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3251 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3252 assert_eq!(ingress.groups, 1);
3253 assert_eq!(ingress.frames, 2);
3254 assert_eq!(ingress.bytes, 10);
3255
3256 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3258 let _ = fetched;
3259 settle().await;
3260 let report = registry.report();
3261 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3262 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3263 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3264 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3265 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3269 }
3270
3271 #[tokio::test]
3276 async fn test_stats_read_frame_counts_once() {
3277 use crate::Timestamp;
3278 use crate::stats::{Config, Registry, Tier};
3279 use bytes::Bytes;
3280
3281 tokio::time::pause();
3282
3283 let registry = Registry::new(Config::new());
3284 let ctx = registry.tier(Tier::default()).session("acme");
3285
3286 let origin = Origin::random().produce();
3287 let ingress = origin.clone().with_stats(ctx.clone());
3288 let egress = origin.consume().with_stats(ctx.clone());
3289
3290 let mut announced = egress.announced();
3291 let source = ingress.create_broadcast("demo", announce()).unwrap();
3292 let mut dynamic = source.dynamic();
3293 settle().await;
3294 settle().await;
3295
3296 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3297 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3298 let mut producer = accept_track(&mut dynamic, "video").await;
3299 settle().await;
3300 let mut sub = subscribing.await.unwrap();
3301
3302 producer
3304 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3305 .unwrap();
3306
3307 let frame = sub.read_frame().await.unwrap().expect("frame");
3308 assert_eq!(frame.payload.len(), 5);
3309 settle().await;
3310
3311 let report = registry.report();
3312 let entry = report
3313 .traffic
3314 .iter()
3315 .find(|e| e.path.as_str() == "demo")
3316 .expect("demo tracked");
3317 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3318 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3319 assert_eq!(
3320 entry.publisher.bytes, 5,
3321 "payload counted once, not zero and not doubled"
3322 );
3323 }
3324
3325 #[tokio::test]
3329 async fn test_stats_datagrams_counted_both_sides() {
3330 use crate::Timestamp;
3331 use crate::stats::{Config, Registry, Tier};
3332
3333 tokio::time::pause();
3334
3335 let registry = Registry::new(Config::new());
3336 let ctx = registry.tier(Tier::default()).session("acme");
3337
3338 let origin = Origin::random().produce();
3339 let ingress = origin.clone().with_stats(ctx.clone());
3340 let egress = origin.consume().with_stats(ctx.clone());
3341
3342 let mut announced = egress.announced();
3343 let source = ingress.create_broadcast("demo", announce()).unwrap();
3344 let mut dynamic = source.dynamic();
3345 settle().await;
3346 settle().await;
3347
3348 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3349 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3350 let mut producer = accept_track(&mut dynamic, "video").await;
3351 settle().await;
3352 let mut sub = subscribing.await.unwrap();
3353
3354 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3355 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3356 assert_eq!(&datagram.payload[..], b"hello");
3357 settle().await;
3358
3359 let report = registry.report();
3360 let entry = report
3361 .traffic
3362 .iter()
3363 .find(|e| e.path.as_str() == "demo")
3364 .expect("demo tracked");
3365
3366 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3367 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3368 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3369 assert_eq!(traffic.frames, 1, "{side}: one frame");
3370 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3371 }
3372 }
3373
3374 #[test]
3375 fn origin_rejects_reserved_ids() {
3376 assert!(Origin::new(0).is_err());
3377 assert!(Origin::new(1u64 << 62).is_err());
3378 assert_eq!(Origin::new(1).unwrap().id(), 1);
3379
3380 let mut zero = [0u8].as_slice();
3381 assert_eq!(
3382 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3383 Origin::UNKNOWN
3384 );
3385 }
3386
3387 #[test]
3388 fn origin_list_push_fails_at_limit() {
3389 let mut list = OriginList::new();
3390 for _ in 0..MAX_HOPS {
3391 list.push(Origin::random()).unwrap();
3392 }
3393 assert_eq!(list.len(), MAX_HOPS);
3394 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3395 }
3396
3397 #[test]
3398 fn origin_list_replace_first() {
3399 let mut list = OriginList::new();
3400 for _ in 0..3 {
3401 list.push(Origin::UNKNOWN).unwrap();
3402 }
3403
3404 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3406 assert_eq!(
3407 list.as_slice(),
3408 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3409 );
3410
3411 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3413 assert_eq!(list.len(), 3);
3414 }
3415
3416 #[test]
3417 fn origin_list_try_from_vec_enforces_limit() {
3418 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3419 assert!(OriginList::try_from(under).is_ok());
3420
3421 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3422 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3423 }
3424
3425 #[tokio::test]
3426 async fn test_announce() {
3427 tokio::time::pause();
3428
3429 let origin = Origin::random().produce();
3430
3431 let mut consumer1 = origin.consume().announced();
3432 consumer1.assert_next_wait();
3433
3434 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3436 settle().await;
3437
3438 consumer1.assert_next_some("test1");
3439 consumer1.assert_next_wait();
3440
3441 let mut consumer2 = origin.consume().announced();
3444
3445 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3447 settle().await;
3448
3449 consumer1.assert_next_some("test2");
3450 consumer1.assert_next_wait();
3451
3452 consumer2.assert_next_some("test1");
3453 consumer2.assert_next_some("test2");
3454 consumer2.assert_next_wait();
3455
3456 broadcast1.finish();
3458 settle().await;
3459
3460 consumer1.assert_next_none("test1");
3462 consumer2.assert_next_none("test1");
3463 consumer1.assert_next_wait();
3464 consumer2.assert_next_wait();
3465
3466 let mut consumer3 = origin.consume().announced();
3468 consumer3.assert_next_some("test2");
3469 consumer3.assert_next_wait();
3470
3471 broadcast2.finish();
3472 settle().await;
3473
3474 consumer1.assert_next_none("test2");
3475 consumer2.assert_next_none("test2");
3476 consumer3.assert_next_none("test2");
3477 }
3478
3479 #[tokio::test]
3483 async fn test_duplicate() {
3484 tokio::time::pause();
3485
3486 let origin = Origin::random().produce();
3487 let consumer = origin.consume();
3488 let mut announced = consumer.announced();
3489
3490 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3491 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3492 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3493 settle().await;
3494 assert!(consumer.get_broadcast("test").is_some());
3495
3496 announced.assert_next_some("test");
3497 announced.assert_next_wait();
3498
3499 broadcast2.finish();
3501 settle().await;
3502 assert!(consumer.get_broadcast("test").is_some());
3503 announced.assert_next_wait();
3504
3505 broadcast1.finish();
3507 settle().await;
3508 assert!(consumer.get_broadcast("test").is_some());
3509 announced.assert_next_wait();
3510
3511 broadcast3.finish();
3513 settle().await;
3514 assert!(consumer.get_broadcast("test").is_none());
3515
3516 announced.assert_next_none("test");
3517 announced.assert_next_wait();
3518 }
3519
3520 #[tokio::test]
3523 async fn test_route_failover() {
3524 tokio::time::pause();
3525
3526 let origin = Origin::random().produce();
3527 let consumer = origin.consume();
3528 let mut announced = consumer.announced();
3529
3530 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3533 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3534
3535 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3537 let mut dynamic_a = source_a.dynamic();
3538 settle().await;
3539 settle().await;
3540 let broadcast = consumer.request_broadcast("test").await.unwrap();
3541 announced.assert_next_some("test");
3542
3543 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3545 let mut dynamic_b = source_b.dynamic();
3546 settle().await;
3547 settle().await;
3548 announced.assert_next_wait();
3549
3550 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3552 let mut producer = accept_track(&mut dynamic_a, "video").await;
3553 settle().await;
3554 dynamic_b.assert_no_request();
3555
3556 let mut sub = subscribing.await.unwrap();
3557 sub.assert_no_group();
3560 assert_eq!(producer.subscription().unwrap().group_start, None);
3561
3562 producer.append_group().unwrap();
3563 producer.append_group().unwrap();
3564 assert_eq!(sub.assert_group().sequence, 0);
3565 assert_eq!(sub.assert_group().sequence, 1);
3566
3567 producer.abort(Error::Dropped).unwrap();
3571 source_a.abort(Error::Dropped).unwrap();
3572 drop(dynamic_a);
3573 settle().await;
3574 announced.assert_next_wait();
3575
3576 let mut producer = accept_track(&mut dynamic_b, "video").await;
3579 settle().await;
3580 sub.assert_no_group();
3581 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3582 producer.create_group(group::Info { sequence: 1 }).unwrap();
3583 producer.create_group(group::Info { sequence: 2 }).unwrap();
3584 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3585 sub.assert_not_closed();
3586 }
3587
3588 #[tokio::test]
3591 async fn test_broadcast_route_watch() {
3592 let mut producer = broadcast::Info::new().produce();
3593 let mut consumer = producer.consume();
3594
3595 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3597
3598 producer.set_route(broadcast::Route::default()).unwrap();
3600 assert!(consumer.route_changed().now_or_never().is_none());
3601
3602 let mut hops = OriginList::new();
3603 hops.push(Origin::new(7).unwrap()).unwrap();
3604 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3605 producer.set_route(route.clone()).unwrap();
3606 assert_eq!(consumer.route_changed().await.unwrap(), route);
3607
3608 let mut fresh = producer.consume();
3610 assert_eq!(fresh.route_changed().await.unwrap(), route);
3611
3612 drop(producer);
3613 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3614 }
3615
3616 #[tokio::test]
3620 async fn test_route_cost_update() {
3621 tokio::time::pause();
3622
3623 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3627 let consumer = origin.consume();
3628 let mut announced = consumer.announced();
3629
3630 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3633 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3634
3635 let mut source_a = origin
3637 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3638 .unwrap();
3639 let mut dynamic_a = source_a.dynamic();
3640 settle().await;
3641 let broadcast = consumer.request_broadcast("test").await.unwrap();
3642 announced.assert_next_some("test");
3643
3644 let mut watch = broadcast.clone();
3645 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3646
3647 let mut source_b = origin
3648 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3649 .unwrap();
3650 let mut dynamic_b = source_b.dynamic();
3651 settle().await;
3652 assert!(
3653 watch.route_changed().now_or_never().is_none(),
3654 "a losing standby must not change the advertised route"
3655 );
3656
3657 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3659 let mut producer = accept_track(&mut dynamic_a, "video").await;
3660 settle().await;
3661 let mut sub = subscribing.await.unwrap();
3662 producer.append_group().unwrap();
3663 assert_eq!(sub.assert_group().sequence, 0);
3664
3665 source_a
3668 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3669 .unwrap();
3670 settle().await;
3671 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3672 announced.assert_next_wait();
3673
3674 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3675 settle().await;
3676 sub.assert_no_group();
3679 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3680 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3681 assert_eq!(sub.assert_group().sequence, 1);
3682 sub.assert_not_closed();
3683
3684 source_b
3686 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3687 .unwrap();
3688 settle().await;
3689 let advertised = watch.route_changed().await.unwrap();
3690 assert_eq!(advertised.hops, hops_b);
3691 assert_eq!(advertised.cost, 5);
3692 announced.assert_next_wait();
3693 }
3694
3695 #[tokio::test]
3698 async fn test_completed_track_survives_route_churn() {
3699 tokio::time::pause();
3700
3701 let origin = Origin::random().produce();
3702 let consumer = origin.consume();
3703
3704 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3706 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3707
3708 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3709 let mut dynamic_a = source_a.dynamic();
3710 settle().await;
3711 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3712 let mut dynamic_b = source_b.dynamic();
3713 settle().await;
3714 settle().await;
3715 let broadcast = consumer.request_broadcast("test").await.unwrap();
3716
3717 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3719 let mut producer = accept_track(&mut dynamic_a, "video").await;
3720 settle().await;
3721 let mut sub = subscribing.await.unwrap();
3722 producer.append_group().unwrap();
3723 assert_eq!(sub.assert_group().sequence, 0);
3724 producer.finish().unwrap();
3725 drop(producer);
3726 settle().await;
3727 sub.assert_closed();
3728
3729 source_a.abort(Error::Dropped).unwrap();
3731 drop(dynamic_a);
3732 settle().await;
3733 dynamic_b.assert_no_request();
3734
3735 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3737 late.assert_closed();
3738 }
3739
3740 #[tokio::test]
3743 async fn test_serve_resets_retry_budget() {
3744 tokio::time::pause();
3745
3746 let origin = Origin::random().produce();
3747 let consumer = origin.consume();
3748
3749 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3750 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3751 let mut dynamic = source.dynamic();
3752 settle().await;
3753 settle().await;
3754 let broadcast = consumer.request_broadcast("test").await.unwrap();
3755
3756 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3758
3759 for _ in 0..2 * MAX_TRACK_RETRIES {
3762 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3763 .await
3764 .expect("timed out waiting for a retry")
3765 .unwrap();
3766 request.reject(Error::NotFound);
3767 let producer = accept_track(&mut dynamic, "video").await;
3768 settle().await;
3769 drop(producer);
3770 }
3771
3772 let _producer = accept_track(&mut dynamic, "video").await;
3773 settle().await;
3774 let mut sub = subscribing.await.unwrap();
3775 sub.assert_not_closed();
3776 }
3777
3778 #[tokio::test]
3782 async fn test_route_handover() {
3783 tokio::time::pause();
3784
3785 let origin = Origin::random().produce();
3786 let consumer = origin.consume();
3787 let mut announced = consumer.announced();
3788
3789 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3791 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3792
3793 let source_a = origin
3794 .create_broadcast("test", announce().with_hops(hops_long))
3795 .unwrap();
3796 let mut dynamic_a = source_a.dynamic();
3797 settle().await;
3798 settle().await;
3799 let broadcast = consumer.request_broadcast("test").await.unwrap();
3800 announced.assert_next_some("test");
3801
3802 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3803 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3804 settle().await;
3805 let mut sub = subscribing.await.unwrap();
3806 producer_a.append_group().unwrap();
3807 producer_a.append_group().unwrap();
3808 assert_eq!(sub.assert_group().sequence, 0);
3809 assert_eq!(sub.assert_group().sequence, 1);
3810
3811 let source_b = origin
3814 .create_broadcast("test", announce().with_hops(hops_short))
3815 .unwrap();
3816 let mut dynamic_b = source_b.dynamic();
3817 settle().await;
3818 settle().await;
3819 announced.assert_next_wait();
3820
3821 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3822 settle().await;
3823
3824 sub.assert_no_group();
3827 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3828 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3829
3830 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3832 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3833 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3834 assert_eq!(sub.assert_group().sequence, 2);
3835 assert_eq!(sub.assert_group().sequence, 3);
3836 sub.assert_no_group();
3837 sub.assert_not_closed();
3838 }
3839
3840 #[tokio::test(start_paused = true)]
3843 async fn test_route_unannounce_immediate() {
3844 let origin = Origin::random().produce();
3845 let consumer = origin.consume();
3846 let mut announced = consumer.announced();
3847
3848 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3849 let mut source = origin
3850 .create_broadcast("test", announce().with_hops(hops.clone()))
3851 .unwrap();
3852 settle().await;
3853 let broadcast = consumer.request_broadcast("test").await.unwrap();
3854 announced.assert_next_some("test");
3855
3856 source.finish();
3859 settle().await;
3860 announced.assert_next_none("test");
3861
3862 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3864 settle().await;
3865 let fresh = consumer.request_broadcast("test").await.unwrap();
3866 announced.assert_next_some("test");
3867 assert!(
3868 !fresh.is_clone(&broadcast),
3869 "re-create must not splice the old broadcast"
3870 );
3871 }
3872
3873 #[tokio::test(start_paused = true)]
3878 async fn test_route_detach_immediate() {
3879 let origin = Origin::random().produce();
3880 let consumer = origin.consume();
3881 let mut announced = consumer.announced();
3882
3883 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3884 let source = origin
3885 .create_broadcast("test", announce().with_hops(hops.clone()))
3886 .unwrap();
3887 let mut dynamic = source.dynamic();
3888 settle().await;
3889 settle().await;
3890 let broadcast = consumer.request_broadcast("test").await.unwrap();
3891 announced.assert_next_some("test");
3892
3893 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3894 let producer = accept_track(&mut dynamic, "video").await;
3895 settle().await;
3896 let mut sub = subscribing.await.unwrap();
3897
3898 drop(producer);
3900 source.abort(Error::Dropped).unwrap();
3901 drop(dynamic);
3902
3903 settle().await;
3904 announced.assert_next_none("test");
3905 sub.assert_error();
3906
3907 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3910 settle().await;
3911 settle().await;
3912 let fresh = consumer.request_broadcast("test").await.unwrap();
3913 announced.assert_next_some("test");
3914 assert!(
3915 !fresh.is_clone(&broadcast),
3916 "re-create must not splice the old broadcast"
3917 );
3918 }
3919
3920 #[tokio::test(start_paused = true)]
3925 async fn test_idle_track_releases_without_respinning() {
3926 let origin = Info::new(Origin::random()).produce();
3927 let consumer = origin.consume();
3928
3929 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3930 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3931 let mut dynamic = source.dynamic();
3932 settle().await;
3933 let broadcast = consumer.request_broadcast("test").await.unwrap();
3934
3935 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3936 let producer = accept_track(&mut dynamic, "video").await;
3937 settle().await;
3938 let sub = subscribing.await.unwrap();
3939
3940 drop(sub);
3943 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
3944 settle().await;
3945 assert!(
3946 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
3947 "the copy must stay spliced inside the linger",
3948 );
3949
3950 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3953 settle().await;
3954 assert!(
3955 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3956 "an idle copy must be released after the linger",
3957 );
3958
3959 for _ in 0..3 {
3963 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3964 settle().await;
3965 assert!(
3966 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3967 "an unread copy must stay released, not be re-spliced",
3968 );
3969 }
3970 assert!(
3971 dynamic.requested_track().now_or_never().is_none(),
3972 "an unread track must not be re-requested",
3973 );
3974 drop(producer);
3975
3976 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3978 let mut producer = accept_track(&mut dynamic, "video").await;
3979 settle().await;
3980 let mut sub = subscribing.await.unwrap();
3981 producer.append_group().unwrap();
3982 assert_eq!(sub.assert_group().sequence, 0);
3983 }
3984
3985 #[tokio::test(start_paused = true)]
3989 async fn test_back_to_back_fetches_reuse_the_track() {
3990 let origin = Info::new(Origin::random()).produce();
3991 let consumer = origin.consume();
3992
3993 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3994 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3995 let mut dynamic = source.dynamic();
3996 settle().await;
3997 let broadcast = consumer.request_broadcast("test").await.unwrap();
3998
3999 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4001 let mut producer = accept_track(&mut dynamic, "video").await;
4002 producer.append_group().unwrap().finish().unwrap();
4003 settle().await;
4004 let first = fetching.await.expect("first fetch");
4005 drop(first);
4006
4007 settle().await;
4009 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4010 settle().await;
4011 assert!(
4012 dynamic.requested_track().now_or_never().is_none(),
4013 "a fetch inside the linger must reuse the track, not re-request it",
4014 );
4015 drop(fetching.await.expect("second fetch"));
4016
4017 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4019 settle().await;
4020 assert!(
4021 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4022 "the copy must be released once the fetches stop",
4023 );
4024 drop(producer);
4025
4026 settle().await;
4028 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4029 let mut producer = accept_track(&mut dynamic, "video").await;
4030 producer.append_group().unwrap().finish().unwrap();
4031 settle().await;
4032 fetching.await.expect("fetch after the linger");
4033 }
4034
4035 #[tokio::test(start_paused = true)]
4039 async fn test_linger_reconnect_splices() {
4040 let origin = Info::new(Origin::random())
4041 .with_linger(Duration::from_secs(5))
4042 .produce();
4043 let consumer = origin.consume();
4044 let mut announced = consumer.announced();
4045
4046 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4047 let source = origin
4048 .create_broadcast("test", announce().with_hops(hops.clone()))
4049 .unwrap();
4050 let mut dynamic = source.dynamic();
4051 settle().await;
4052 settle().await;
4053 let broadcast = consumer.request_broadcast("test").await.unwrap();
4054 announced.assert_next_some("test");
4055
4056 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4057 let mut producer = accept_track(&mut dynamic, "video").await;
4058 settle().await;
4059 let mut sub = subscribing.await.unwrap();
4060
4061 producer.append_group().unwrap();
4062 producer.append_group().unwrap();
4063 assert_eq!(sub.assert_group().sequence, 0);
4064 assert_eq!(sub.assert_group().sequence, 1);
4065
4066 drop(producer);
4069 source.abort(Error::Dropped).unwrap();
4070 drop(dynamic);
4071 settle().await;
4072
4073 announced.assert_next_wait();
4075 sub.assert_no_group();
4076 sub.assert_not_closed();
4077
4078 let during = consumer.request_broadcast("test").await.unwrap();
4080 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4081
4082 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4085 let mut dynamic = source.dynamic();
4086 settle().await;
4087 settle().await;
4088 announced.assert_next_wait();
4089 let again = consumer.request_broadcast("test").await.unwrap();
4090 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4091
4092 let mut producer = accept_track(&mut dynamic, "video").await;
4096 settle().await;
4097 sub.assert_no_group();
4098 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
4099 producer.create_group(group::Info { sequence: 2 }).unwrap();
4100 assert_eq!(sub.assert_group().sequence, 2);
4101 sub.assert_not_closed();
4102 }
4103
4104 #[tokio::test(start_paused = true)]
4107 async fn test_linger_expiry_closes() {
4108 let origin = Info::new(Origin::random())
4109 .with_linger(Duration::from_secs(5))
4110 .produce();
4111 let consumer = origin.consume();
4112 let mut announced = consumer.announced();
4113
4114 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4115 let source = origin
4116 .create_broadcast("test", announce().with_hops(hops.clone()))
4117 .unwrap();
4118 let mut dynamic = source.dynamic();
4119 settle().await;
4120 settle().await;
4121 let broadcast = consumer.request_broadcast("test").await.unwrap();
4122 announced.assert_next_some("test");
4123
4124 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4125 let producer = accept_track(&mut dynamic, "video").await;
4126 settle().await;
4127 let mut sub = subscribing.await.unwrap();
4128
4129 drop(producer);
4130 source.abort(Error::Dropped).unwrap();
4131 drop(dynamic);
4132 settle().await;
4133 announced.assert_next_wait();
4134
4135 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4137 settle().await;
4138 announced.assert_next_none("test");
4139 sub.assert_error();
4140
4141 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4143 settle().await;
4144 settle().await;
4145 let fresh = consumer.request_broadcast("test").await.unwrap();
4146 announced.assert_next_some("test");
4147 assert!(
4148 !fresh.is_clone(&broadcast),
4149 "a late re-create must not splice the expired broadcast"
4150 );
4151 }
4152
4153 #[tokio::test(start_paused = true)]
4157 async fn test_linger_forever() {
4158 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4159 let consumer = origin.consume();
4160 let mut announced = consumer.announced();
4161
4162 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4163 let source = origin
4164 .create_broadcast("test", announce().with_hops(hops.clone()))
4165 .unwrap();
4166 settle().await;
4167 let broadcast = consumer.request_broadcast("test").await.unwrap();
4168 announced.assert_next_some("test");
4169
4170 source.abort(Error::Dropped).unwrap();
4171 settle().await;
4172
4173 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4175 announced.assert_next_wait();
4176 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4177 settle().await;
4178 settle().await;
4179 let again = consumer.request_broadcast("test").await.unwrap();
4180 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4181 drop(source);
4182 }
4183
4184 #[tokio::test(start_paused = true)]
4199 async fn test_linger_parks_a_live_subscription() {
4200 let origin = Info::new(Origin::random())
4201 .with_linger(Duration::from_secs(5))
4202 .produce();
4203 let consumer = origin.consume();
4204
4205 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4206 let source = origin
4207 .create_broadcast("test", announce().with_hops(hops.clone()))
4208 .unwrap();
4209 let mut dynamic = source.dynamic();
4210 settle().await;
4211 settle().await;
4212 let broadcast = consumer.request_broadcast("test").await.unwrap();
4213
4214 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4217 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4218 settle().await;
4219 let mut sub = subscribing.await.unwrap();
4220 producer.append_group().unwrap();
4221 assert_eq!(sub.assert_group().sequence, 0);
4222
4223 source.abort(Error::Dropped).unwrap();
4228 settle().await;
4229 settle().await;
4230 sub.assert_not_closed();
4231
4232 tokio::time::sleep(Duration::from_secs(4)).await;
4235 settle().await;
4236 sub.assert_not_closed();
4237
4238 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4240 let mut dynamic = source.dynamic();
4241 settle().await;
4242 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4243 settle().await;
4244 producer.create_group(group::Info { sequence: 1 }).unwrap();
4245 assert_eq!(sub.assert_group().sequence, 1);
4246 sub.assert_not_closed();
4247 }
4248
4249 #[tokio::test(start_paused = true)]
4260 async fn test_idle_release_survives_the_route_leaving() {
4261 let origin = Info::new(Origin::random())
4262 .with_linger(Duration::from_secs(600))
4263 .produce();
4264 let consumer = origin.consume();
4265
4266 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4267 let source = origin
4268 .create_broadcast("test", announce().with_hops(hops.clone()))
4269 .unwrap();
4270 let mut dynamic = source.dynamic();
4271 settle().await;
4272 settle().await;
4273 let broadcast = consumer.request_broadcast("test").await.unwrap();
4274
4275 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4276 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4277 settle().await;
4278 let mut sub = subscribing.await.unwrap();
4279 producer.append_group().unwrap();
4280 assert_eq!(sub.assert_group().sequence, 0);
4281
4282 source.abort(Error::Dropped).unwrap();
4285 settle().await;
4286 settle().await;
4287 drop(sub);
4288 drop(producer);
4289 drop(dynamic);
4290 settle().await;
4291 tokio::time::sleep(TRACK_IDLE_LINGER + Duration::from_secs(1)).await;
4292 settle().await;
4293
4294 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4297 let mut dynamic = source.dynamic();
4298 settle().await;
4299 settle().await;
4300 let broadcast = consumer.request_broadcast("test").await.unwrap();
4301 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4302 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4303 settle().await;
4304 let mut sub = subscribing.await.unwrap();
4305 producer.append_group().unwrap();
4306 assert_eq!(
4307 sub.assert_group().sequence,
4308 0,
4309 "the reconnect's first group must not be filtered by a stale boundary"
4310 );
4311 }
4312
4313 #[tokio::test(start_paused = true)]
4316 async fn test_linger_skipped_on_finish() {
4317 let origin = Info::new(Origin::random())
4318 .with_linger(Duration::from_secs(5))
4319 .produce();
4320 let consumer = origin.consume();
4321 let mut announced = consumer.announced();
4322
4323 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4324 let mut source = origin
4325 .create_broadcast("test", announce().with_hops(hops.clone()))
4326 .unwrap();
4327 settle().await;
4328 let broadcast = consumer.request_broadcast("test").await.unwrap();
4329 announced.assert_next_some("test");
4330
4331 source.finish();
4334 settle().await;
4335 announced.assert_next_none("test");
4336
4337 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4339 settle().await;
4340 let fresh = consumer.request_broadcast("test").await.unwrap();
4341 announced.assert_next_some("test");
4342 assert!(
4343 !fresh.is_clone(&broadcast),
4344 "a finish must not leave a lingering broadcast to splice into"
4345 );
4346 }
4347
4348 #[tokio::test]
4351 async fn test_announce_toggle() {
4352 tokio::time::pause();
4353
4354 let origin = Origin::random().produce();
4355 let consumer = origin.consume();
4356 let mut announced = consumer.announced();
4357
4358 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4359 settle().await;
4360
4361 announced.assert_next_wait();
4363 let broadcast = consumer
4364 .get_broadcast("test")
4365 .expect("offline broadcast is still routable");
4366 assert!(!broadcast.route().announce);
4367
4368 let requested = consumer.request_broadcast("test").await.unwrap();
4370 assert!(requested.is_clone(&broadcast));
4371
4372 source.set_route(announce()).unwrap();
4374 settle().await;
4375 let face = announced.assert_next_some("test");
4376 assert!(face.is_clone(&broadcast));
4377
4378 let mut fresh = origin.consume().announced();
4380 fresh.assert_next_some("test");
4381 fresh.assert_next_wait();
4382
4383 source.set_route(broadcast::Route::new()).unwrap();
4385 settle().await;
4386 announced.assert_next_none("test");
4387 assert!(consumer.get_broadcast("test").is_some());
4388 let mut fresh = origin.consume().announced();
4389 fresh.assert_next_wait();
4390
4391 source.finish();
4392 settle().await;
4393 assert!(consumer.get_broadcast("test").is_none());
4394 }
4395
4396 #[tokio::test]
4399 async fn test_announce_beats_offline() {
4400 tokio::time::pause();
4401
4402 let origin = Origin::random().produce();
4403 let consumer = origin.consume();
4404 let mut announced = consumer.announced();
4405
4406 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4408 settle().await;
4409 announced.assert_next_wait();
4410
4411 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
4414 settle().await;
4415 announced.assert_next_some("test");
4416 let face = consumer.get_broadcast("test").unwrap();
4417 assert!(face.route().announce);
4418 assert_eq!(face.route().cost, 10);
4419
4420 announced_source.finish();
4423 settle().await;
4424 announced.assert_next_none("test");
4425 assert!(consumer.get_broadcast("test").is_some());
4426 }
4427
4428 #[tokio::test]
4431 async fn test_better_source_no_churn() {
4432 tokio::time::pause();
4433
4434 let origin = Origin::random().produce();
4435 let mut announced = origin.consume().announced();
4436
4437 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4440 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4441 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4442 settle().await;
4443 let face = announced.assert_next_some("test");
4444
4445 let _b = origin
4446 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4447 .unwrap();
4448 settle().await;
4449 announced.assert_next_wait();
4450 let current = origin.consume().get_broadcast("test").unwrap();
4451 assert!(current.is_clone(&face), "the broadcast identity must not change");
4452 assert_eq!(current.route().hops, hops_b);
4454 }
4455
4456 #[tokio::test]
4462 async fn test_publisher_mismatch_replaces() {
4463 tokio::time::pause();
4464
4465 let origin = Origin::random().produce();
4466 let consumer = origin.consume();
4467 let mut announced = consumer.announced();
4468
4469 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4470 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4471
4472 let mut source_a = origin
4473 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4474 .unwrap();
4475 settle().await;
4476 let face_a = announced.assert_next_some("test");
4477
4478 let _source_b = origin
4481 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4482 .unwrap();
4483 settle().await;
4484 settle().await;
4485 announced.assert_next_none("test");
4486 let face_b = announced.assert_next_some("test");
4487 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
4488 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4489 assert!(face_a.is_closed(), "the displaced front must close");
4493
4494 source_a.finish();
4496 settle().await;
4497 settle().await;
4498 announced.assert_next_wait();
4499 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4500 }
4501
4502 #[tokio::test]
4507 async fn test_reconnect_wins_over_stale_route() {
4508 tokio::time::pause();
4509
4510 let origin = Origin::random().produce();
4511 let consumer = origin.consume();
4512
4513 let publisher = Origin::new(1).unwrap();
4514 let hops = OriginList::try_from(vec![publisher]).unwrap();
4515
4516 let stale = origin
4519 .create_broadcast("test", announce().with_hops(hops.clone()))
4520 .unwrap();
4521 let mut stale_dynamic = stale.dynamic();
4522 settle().await;
4523
4524 let fresh = origin
4526 .create_broadcast("test", announce().with_hops(hops.clone()))
4527 .unwrap();
4528 let mut fresh_dynamic = fresh.dynamic();
4529 settle().await;
4530 settle().await;
4531
4532 let broadcast = consumer.request_broadcast("test").await.unwrap();
4534 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4535 settle().await;
4536 let _producer = accept_track(&mut fresh_dynamic, "video").await;
4537 settle().await;
4538 subscribing.await.unwrap();
4539 stale_dynamic.assert_no_request();
4540 }
4541
4542 #[tokio::test]
4548 async fn test_carrying_reconnect_switches_immediately() {
4549 tokio::time::pause();
4550
4551 let origin = Origin::random().produce();
4552 let consumer = origin.consume();
4553
4554 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4555
4556 let stale = origin
4557 .create_broadcast("test", announce().with_hops(hops.clone()))
4558 .unwrap();
4559 let mut stale_dynamic = stale.dynamic();
4560 settle().await;
4561
4562 let broadcast = consumer.request_broadcast("test").await.unwrap();
4564 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4565 settle().await;
4566 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
4567 settle().await;
4568 let _subscription = subscribing.await.unwrap();
4570
4571 let fresh = origin
4573 .create_broadcast("test", announce().with_hops(hops.clone()))
4574 .unwrap();
4575 let mut fresh_dynamic = fresh.dynamic();
4576 settle().await;
4577 settle().await;
4578
4579 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
4582 }
4583
4584 #[tokio::test]
4594 async fn test_offline_mismatch_never_evicts_a_live_front() {
4595 tokio::time::pause();
4596
4597 let origin = Origin::random().produce();
4598 let consumer = origin.consume();
4599 let mut announced = consumer.announced();
4600
4601 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4602 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4603
4604 let mut live = origin
4605 .create_broadcast("test", announce().with_hops(hops_live.clone()))
4606 .unwrap();
4607 let mut live_dynamic = live.dynamic();
4608 settle().await;
4609 let face = announced.assert_next_some("test");
4610
4611 let broadcast = consumer.request_broadcast("test").await.unwrap();
4613 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4614 settle().await;
4615 let _producer = accept_track(&mut live_dynamic, "video").await;
4616 settle().await;
4617 subscribing.await.unwrap();
4618
4619 let cache = origin
4622 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
4623 .unwrap();
4624 settle().await;
4625 settle().await;
4626 announced.assert_next_wait();
4627 assert!(!face.is_closed(), "the live front must survive");
4628 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
4629
4630 live.finish();
4633 settle().await;
4634 settle().await;
4635 announced.assert_next_none("test");
4636 let taken = consumer
4637 .get_broadcast("test")
4638 .expect("the parked source must take over");
4639 assert_eq!(taken.route().hops, hops_cache);
4640 announced.assert_next_wait();
4642 drop(cache);
4643 }
4644
4645 #[tokio::test]
4651 async fn test_dispatch_excludes_requester() {
4652 tokio::time::pause();
4653
4654 let origin = Origin::random().produce();
4655 let consumer = origin.consume();
4656
4657 let peer = Origin::new(5).unwrap();
4658 let publisher = Origin::new(1).unwrap();
4659 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
4661 let clean = OriginList::try_from(vec![publisher]).unwrap();
4662
4663 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4664 let mut dynamic_a = source_a.dynamic();
4665 settle().await;
4666 let source_b = origin
4667 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
4668 .unwrap();
4669 let mut dynamic_b = source_b.dynamic();
4670 settle().await;
4671 settle().await;
4672
4673 let shared = consumer.request_broadcast("test").await.unwrap();
4676 let subscribing = shared.track("video").unwrap().subscribe(None);
4677 let _producer_a = accept_track(&mut dynamic_a, "video").await;
4678 settle().await;
4679 subscribing.await.unwrap();
4680
4681 let scoped = consumer.clone().excluding(peer);
4686 let pinned = scoped.request_broadcast("test").await.unwrap();
4687 let subscribing = pinned.track("video").unwrap().subscribe(None);
4688 let _producer_b = accept_track(&mut dynamic_b, "video").await;
4689 settle().await;
4690 subscribing.await.unwrap();
4691 dynamic_a.assert_no_request();
4692 }
4693
4694 #[tokio::test]
4700 async fn test_standby_join_splices_live_subscriber() {
4701 tokio::time::pause();
4702
4703 let origin = Origin::random().produce();
4704 let consumer = origin.consume();
4705
4706 let publisher = Origin::new(1).unwrap();
4707 let peer = Origin::new(5).unwrap();
4708 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4709 let local = OriginList::try_from(vec![publisher]).unwrap();
4710
4711 let source_remote = origin
4713 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4714 .unwrap();
4715 let mut dynamic_remote = source_remote.dynamic();
4716 settle().await;
4717 settle().await;
4718 let broadcast = consumer.request_broadcast("test").await.unwrap();
4719 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4720 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
4721 settle().await;
4722 let mut sub = subscribing.await.unwrap();
4723 producer_remote.append_group().unwrap();
4724 assert_eq!(sub.assert_group().sequence, 0);
4725
4726 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4729 let mut dynamic_local = source_local.dynamic();
4730 settle().await;
4731 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4732 settle().await;
4733 sub.assert_no_group();
4734 assert_eq!(producer_local.subscription().unwrap().group_start, Some(1));
4735 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
4736 assert_eq!(sub.assert_group().sequence, 1);
4737 sub.assert_not_closed();
4738 }
4739
4740 #[tokio::test]
4747 async fn test_standby_missing_track_keeps_incumbent() {
4748 tokio::time::pause();
4749
4750 let origin = Origin::random().produce();
4751 let consumer = origin.consume();
4752
4753 let publisher = Origin::new(1).unwrap();
4754 let peer = Origin::new(5).unwrap();
4755 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4756 let local = OriginList::try_from(vec![publisher]).unwrap();
4757
4758 let source_remote = origin
4760 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4761 .unwrap();
4762 let mut dynamic_remote = source_remote.dynamic();
4763 settle().await;
4764 settle().await;
4765 let broadcast = consumer.request_broadcast("test").await.unwrap();
4766 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4767 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
4768 settle().await;
4769 let mut sub = subscribing.await.unwrap();
4770 producer_remote.append_group().unwrap();
4771 assert_eq!(sub.assert_group().sequence, 0);
4772
4773 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4778 let mut dynamic_local = source_local.dynamic();
4779 settle().await;
4780 for _ in 0..2 * MAX_TRACK_RETRIES {
4781 if let Some(Ok(request)) = dynamic_local.requested_track().now_or_never() {
4782 assert_eq!(request.name(), "audio");
4783 request.reject(Error::NotFound);
4784 }
4785 settle().await;
4786 }
4787
4788 producer_remote.append_group().unwrap();
4790 assert_eq!(sub.assert_group().sequence, 1);
4791 sub.assert_not_closed();
4792
4793 source_remote.abort(Error::Dropped).unwrap();
4796 settle().await;
4797 settle().await;
4798 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
4799 settle().await;
4800 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
4801 assert_eq!(sub.assert_group().sequence, 2);
4802 sub.assert_not_closed();
4803 }
4804
4805 #[tokio::test]
4810 async fn test_unservable_track_retried_by_a_later_request() {
4811 tokio::time::pause();
4812
4813 let origin = Origin::random().produce();
4814 let consumer = origin.consume();
4815
4816 let source = origin.create_broadcast("test", announce()).unwrap();
4817 let mut dynamic = source.dynamic();
4818 settle().await;
4819 settle().await;
4820 let broadcast = consumer.request_broadcast("test").await.unwrap();
4821
4822 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4824 for _ in 0..MAX_TRACK_RETRIES {
4825 let request = dynamic.requested_track().await.unwrap();
4826 request.reject(Error::NotFound);
4827 settle().await;
4828 }
4829 assert!(matches!(subscribing.await, Err(Error::Unroutable)));
4830
4831 let retry = broadcast.track("audio").unwrap().subscribe(None);
4833 let mut producer = accept_track(&mut dynamic, "audio").await;
4834 settle().await;
4835 let mut sub = retry.await.expect("a fresh request must reach the source");
4836 producer.append_group().unwrap();
4837 assert_eq!(sub.assert_group().sequence, 0);
4838 }
4839
4840 #[tokio::test]
4846 async fn test_per_track_fallback_respects_exclusion() {
4847 tokio::time::pause();
4848
4849 let origin = Origin::random().produce();
4850 let consumer = origin.consume();
4851
4852 let publisher = Origin::new(1).unwrap();
4853 let peer = Origin::new(5).unwrap();
4854 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4855 let local = OriginList::try_from(vec![publisher]).unwrap();
4856
4857 let source_tainted = origin
4859 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4860 .unwrap();
4861 let mut dynamic_tainted = source_tainted.dynamic();
4862 settle().await;
4863 settle().await;
4864
4865 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4868 let mut dynamic_clean = source_clean.dynamic();
4869 settle().await;
4870
4871 let scoped = consumer.clone().excluding(peer);
4872 let broadcast = scoped.request_broadcast("test").await.unwrap();
4873 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
4874 settle().await;
4875
4876 for _ in 0..2 * MAX_TRACK_RETRIES {
4879 if let Some(Ok(request)) = dynamic_clean.requested_track().now_or_never() {
4880 request.reject(Error::NotFound);
4881 }
4882 settle().await;
4883 }
4884 dynamic_tainted.assert_no_request();
4885 }
4886
4887 #[tokio::test]
4891 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
4892 tokio::time::pause();
4893
4894 let origin = Origin::random().produce();
4895 let consumer = origin.consume();
4896
4897 let publisher = Origin::new(1).unwrap();
4898 let peer = Origin::new(5).unwrap();
4899 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4900 let local = OriginList::try_from(vec![publisher]).unwrap();
4901
4902 let source_tainted = origin
4903 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4904 .unwrap();
4905 let mut dynamic_tainted = source_tainted.dynamic();
4906 settle().await;
4907 settle().await;
4908 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4909 let mut dynamic_clean = source_clean.dynamic();
4910 settle().await;
4911
4912 let scoped = consumer.clone().excluding(peer);
4913 let broadcast = scoped.request_broadcast("test").await.unwrap();
4914 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
4915 let _clean = accept_track(&mut dynamic_clean, "video").await;
4916 settle().await;
4917
4918 source_clean.abort(Error::Dropped).unwrap();
4920 settle().await;
4921 settle().await;
4922 dynamic_tainted.assert_no_request();
4923
4924 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
4927 }
4928
4929 #[tokio::test]
4934 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
4935 tokio::time::pause();
4936
4937 let origin = Origin::random().produce();
4938 let consumer = origin.consume();
4939
4940 let publisher = Origin::new(1).unwrap();
4941 let peer = Origin::new(5).unwrap();
4942 let local = OriginList::try_from(vec![publisher]).unwrap();
4943 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4944
4945 let source_clean = origin
4949 .create_broadcast("test", announce().with_hops(local).with_cost(5))
4950 .unwrap();
4951 let mut dynamic_clean = source_clean.dynamic();
4952 settle().await;
4953 settle().await;
4954 let scoped = consumer.clone().excluding(peer);
4955 let broadcast = scoped.request_broadcast("test").await.unwrap();
4956 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4957 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
4958 settle().await;
4959 let mut sub = subscribing.await.unwrap();
4960 producer_clean.append_group().unwrap();
4961 assert_eq!(sub.assert_group().sequence, 0);
4962
4963 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
4970 tainted.advertised = 1;
4971 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
4972 let mut dynamic_tainted = source_tainted.dynamic();
4973 settle().await;
4974 settle().await;
4975 dynamic_tainted.assert_no_request();
4976 producer_clean.append_group().unwrap();
4977 assert_eq!(sub.assert_group().sequence, 1);
4978 sub.assert_not_closed();
4979
4980 drop(sub);
4983 drop(broadcast);
4984 drop(scoped);
4985 settle().await;
4986 let mut bumped = announce().with_hops(via_peer).with_cost(1);
4987 bumped.advertised = 1;
4988 source_tainted.set_route(bumped).unwrap();
4989 settle().await;
4990 let plain = consumer.request_broadcast("test").await.unwrap();
4991 let _plain_track = plain.track("video").unwrap().subscribe(None);
4992 settle().await;
4993 settle().await;
4994 assert!(
4995 dynamic_tainted.requested_track().now_or_never().is_some(),
4996 "the front must be free to use the route again once the peer is gone"
4997 );
4998 }
4999
5000 #[tokio::test]
5004 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
5005 tokio::time::pause();
5006
5007 let origin = Origin::random().produce();
5008 let consumer = origin.consume();
5009 let mut dynamic = origin.dynamic();
5010
5011 let peer = Origin::new(5).unwrap();
5012 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5013 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5014 settle().await;
5015 settle().await;
5016
5017 let scoped = consumer.clone().excluding(peer);
5018 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5019 assert!(
5020 dynamic.requested_broadcast().now_or_never().is_none(),
5021 "the dynamic handler was asked to route around the exclusion"
5022 );
5023
5024 let _pending = scoped.request_broadcast("other");
5026 settle().await;
5027 assert!(
5028 dynamic.requested_broadcast().now_or_never().is_some(),
5029 "a genuinely missing path must still fall back"
5030 );
5031 }
5032
5033 #[tokio::test]
5036 async fn test_dispatch_all_tainted_unroutable() {
5037 tokio::time::pause();
5038
5039 let origin = Origin::random().produce();
5040 let consumer = origin.consume();
5041
5042 let peer = Origin::new(5).unwrap();
5043 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5044 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5045 settle().await;
5046 settle().await;
5047
5048 let scoped = consumer.clone().excluding(peer);
5049 match scoped.request_broadcast("test").await {
5050 Err(Error::Unroutable) => {}
5051 Err(err) => panic!("expected Unroutable, got {err:?}"),
5052 Ok(_) => panic!("expected Unroutable, got a broadcast"),
5053 }
5054
5055 consumer.request_broadcast("test").await.unwrap();
5057 }
5058
5059 #[tokio::test]
5060 async fn test_duplicate_reverse() {
5061 tokio::time::pause();
5062
5063 let origin = Origin::random().produce();
5064
5065 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5066 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5067 settle().await;
5068 assert!(origin.consume().get_broadcast("test").is_some());
5069
5070 broadcast2.finish();
5072 settle().await;
5073 assert!(origin.consume().get_broadcast("test").is_some());
5074
5075 broadcast1.finish();
5076 settle().await;
5077 assert!(origin.consume().get_broadcast("test").is_none());
5078 }
5079
5080 #[tokio::test]
5081 async fn test_deterministic_tiebreak() {
5082 tokio::time::pause();
5083
5084 fn hops(ids: &[u64]) -> OriginList {
5085 OriginList::try_from(
5086 ids.iter()
5087 .copied()
5088 .map(|id| Origin::new(id).unwrap())
5089 .collect::<Vec<_>>(),
5090 )
5091 .unwrap()
5092 }
5093
5094 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
5097 let origin = Origin::random().produce();
5098 let _a = origin
5099 .create_broadcast("test", announce().with_hops(hops(first)))
5100 .unwrap();
5101 let _b = origin
5102 .create_broadcast("test", announce().with_hops(hops(second)))
5103 .unwrap();
5104 settle().await;
5105 origin.consume().get_broadcast("test").unwrap().route().hops
5106 }
5107
5108 let forward = winner(&[5, 20], &[5, 40]).await;
5112 let reverse = winner(&[5, 40], &[5, 20]).await;
5113 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
5114
5115 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
5117 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
5118 }
5119
5120 #[tokio::test]
5125 async fn test_many_announces() {
5126 let origin = Origin::random().produce();
5127
5128 let mut consumer = origin.consume().announced();
5129 let mut broadcasts = Vec::new();
5131 for i in 0..256 {
5132 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5133 settle().await;
5134 }
5135
5136 for i in 0..256 {
5137 consumer.assert_next_some(format!("test{i:03}"));
5138 }
5139 consumer.assert_next_wait();
5140 }
5141
5142 #[tokio::test]
5143 async fn test_many_announces_try() {
5144 let origin = Origin::random().produce();
5145
5146 let mut consumer = origin.consume().announced();
5147 let mut broadcasts = Vec::new();
5149 for i in 0..256 {
5150 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5151 settle().await;
5152 }
5153
5154 for i in 0..256 {
5155 consumer.assert_try_next_some(format!("test{i:03}"));
5156 }
5157 }
5158
5159 #[tokio::test]
5160 async fn test_with_root_basic() {
5161 let origin = Origin::random().produce();
5162
5163 let foo_producer = origin.with_root("foo").expect("should create root");
5165 assert_eq!(foo_producer.root().as_str(), "foo");
5166
5167 let mut consumer = origin.consume().announced();
5168
5169 let _broadcast = foo_producer
5171 .create_broadcast("bar/baz", announce())
5172 .expect("publish allowed");
5173 settle().await;
5174 consumer.assert_next_some("foo/bar/baz");
5176
5177 let mut foo_consumer = foo_producer.consume().announced();
5179 foo_consumer.assert_next_some("bar/baz");
5180 }
5181
5182 #[tokio::test]
5183 async fn test_with_root_nested() {
5184 let origin = Origin::random().produce();
5185
5186 let foo_producer = origin.with_root("foo").expect("should create foo root");
5188 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
5189 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
5190
5191 let mut consumer = origin.consume().announced();
5192
5193 let _broadcast = foo_bar_producer
5195 .create_broadcast("baz", announce())
5196 .expect("publish allowed");
5197 settle().await;
5198 consumer.assert_next_some("foo/bar/baz");
5200
5201 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
5203 foo_bar_consumer.assert_next_some("baz");
5204 }
5205
5206 #[tokio::test]
5207 async fn test_publish_scope_allows() {
5208 let origin = Origin::random().produce();
5209
5210 let limited_producer = origin
5212 .scope(&["allowed/path1".into(), "allowed/path2".into()])
5213 .expect("should create limited producer");
5214
5215 let _broadcast = limited_producer
5217 .create_broadcast("allowed/path1", announce())
5218 .expect("publish allowed");
5219 let _keep2 = limited_producer
5220 .create_broadcast("allowed/path1/nested", announce())
5221 .expect("publish allowed");
5222 let _keep3 = limited_producer
5223 .create_broadcast("allowed/path2", announce())
5224 .expect("publish allowed");
5225 settle().await;
5226
5227 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
5229 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
5231 }
5232
5233 #[tokio::test]
5234 async fn test_publish_max_parts() {
5235 let origin = Origin::random().produce();
5236
5237 let at_limit = (0..Path::MAX_PARTS)
5238 .map(|i| i.to_string())
5239 .collect::<Vec<_>>()
5240 .join("/");
5241 let _broadcast = origin
5242 .create_broadcast(at_limit.as_str(), announce())
5243 .expect("publish allowed");
5244 settle().await;
5245
5246 let too_deep = format!("{at_limit}/extra");
5247 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
5248
5249 let rooted = origin.with_root("root").expect("wildcard allows any root");
5251 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
5252 }
5253
5254 #[tokio::test]
5255 async fn test_publish_scope_empty() {
5256 let origin = Origin::random().produce();
5257
5258 assert!(origin.scope(&[]).is_none());
5260 }
5261
5262 #[tokio::test]
5263 async fn test_consume_scope_filters() {
5264 let origin = Origin::random().produce();
5265
5266 let mut consumer = origin.consume().announced();
5267
5268 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
5270 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
5271 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
5272 settle().await;
5273
5274 let mut limited_consumer = origin
5276 .consume()
5277 .scope(&["allowed".into()])
5278 .expect("should create limited consumer")
5279 .announced();
5280
5281 limited_consumer.assert_next_some("allowed");
5283 limited_consumer.assert_next_some("allowed/nested");
5284 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
5288 consumer.assert_next_some("allowed/nested");
5289 consumer.assert_next_some("notallowed");
5290 }
5291
5292 #[tokio::test]
5293 async fn test_consume_scope_multiple_prefixes() {
5294 let origin = Origin::random().produce();
5295
5296 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5297 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5298 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5299 settle().await;
5300
5301 let mut limited_consumer = origin
5303 .consume()
5304 .scope(&["foo".into(), "bar".into()])
5305 .expect("should create limited consumer")
5306 .announced();
5307
5308 limited_consumer.assert_next_some("bar/test");
5310 limited_consumer.assert_next_some("foo/test");
5311 limited_consumer.assert_next_wait(); }
5313
5314 #[tokio::test]
5315 async fn test_with_root_and_publish_scope() {
5316 let origin = Origin::random().produce();
5317
5318 let foo_producer = origin.with_root("foo").expect("should create foo root");
5320
5321 let limited_producer = foo_producer
5323 .scope(&["bar".into(), "goop/pee".into()])
5324 .expect("should create limited producer");
5325
5326 let mut consumer = origin.consume().announced();
5327
5328 let _broadcast = limited_producer
5330 .create_broadcast("bar", announce())
5331 .expect("publish allowed");
5332 let _keep2 = limited_producer
5333 .create_broadcast("bar/nested", announce())
5334 .expect("publish allowed");
5335 let _keep3 = limited_producer
5336 .create_broadcast("goop/pee", announce())
5337 .expect("publish allowed");
5338 let _keep4 = limited_producer
5339 .create_broadcast("goop/pee/nested", announce())
5340 .expect("publish allowed");
5341 settle().await;
5342
5343 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
5345 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
5347
5348 consumer.assert_next_some("foo/bar");
5350 consumer.assert_next_some("foo/bar/nested");
5351 consumer.assert_next_some("foo/goop/pee");
5352 consumer.assert_next_some("foo/goop/pee/nested");
5353 }
5354
5355 #[tokio::test]
5356 async fn test_with_root_and_consume_scope() {
5357 let origin = Origin::random().produce();
5358
5359 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
5361 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
5362 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
5363 settle().await;
5364
5365 let foo_producer = origin.with_root("foo").expect("should create foo root");
5367
5368 let mut limited_consumer = foo_producer
5370 .consume()
5371 .scope(&["bar".into(), "goop/pee".into()])
5372 .expect("should create limited consumer")
5373 .announced();
5374
5375 limited_consumer.assert_next_some("bar/test");
5377 limited_consumer.assert_next_some("goop/pee/test");
5378 limited_consumer.assert_next_wait(); }
5380
5381 #[tokio::test]
5382 async fn test_with_root_unauthorized() {
5383 let origin = Origin::random().produce();
5384
5385 let limited_producer = origin
5387 .scope(&["allowed".into()])
5388 .expect("should create limited producer");
5389
5390 assert!(limited_producer.with_root("notallowed").is_none());
5392
5393 let allowed_root = limited_producer
5395 .with_root("allowed")
5396 .expect("should create allowed root");
5397 assert_eq!(allowed_root.root().as_str(), "allowed");
5398 }
5399
5400 #[tokio::test]
5401 async fn test_wildcard_permission() {
5402 let origin = Origin::random().produce();
5403
5404 let root_producer = origin.clone();
5406
5407 let _broadcast = root_producer
5409 .create_broadcast("any/path", announce())
5410 .expect("publish allowed");
5411 let _keep2 = root_producer
5412 .create_broadcast("other/path", announce())
5413 .expect("publish allowed");
5414 settle().await;
5415
5416 let foo_producer = root_producer.with_root("foo").expect("should create any root");
5418 assert_eq!(foo_producer.root().as_str(), "foo");
5419 }
5420
5421 #[tokio::test]
5422 async fn test_consume_broadcast_with_permissions() {
5423 let origin = Origin::random().produce();
5424
5425 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
5426 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
5427 settle().await;
5428
5429 let limited_consumer = origin
5431 .consume()
5432 .scope(&["allowed".into()])
5433 .expect("should create limited consumer");
5434
5435 let result = limited_consumer.get_broadcast("allowed/test");
5437 assert!(result.is_some());
5438 assert!(
5439 result
5440 .unwrap()
5441 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
5442 );
5443
5444 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
5446
5447 let consumer = origin.consume();
5449 assert!(consumer.get_broadcast("allowed/test").is_some());
5450 assert!(consumer.get_broadcast("notallowed/test").is_some());
5451 }
5452
5453 #[tokio::test]
5454 async fn test_nested_paths_with_permissions() {
5455 let origin = Origin::random().produce();
5456
5457 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
5459
5460 let _broadcast = limited_producer
5462 .create_broadcast("a/b/c", announce())
5463 .expect("publish allowed");
5464 let _keep2 = limited_producer
5465 .create_broadcast("a/b/c/d", announce())
5466 .expect("publish allowed");
5467 let _keep3 = limited_producer
5468 .create_broadcast("a/b/c/d/e", announce())
5469 .expect("publish allowed");
5470 settle().await;
5471
5472 assert!(limited_producer.create_broadcast("a", announce()).is_err());
5474 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
5475 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
5476 }
5477
5478 #[tokio::test]
5479 async fn test_multiple_consumers_with_different_permissions() {
5480 let origin = Origin::random().produce();
5481
5482 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5484 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5485 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5486 settle().await;
5487
5488 let mut foo_consumer = origin
5490 .consume()
5491 .scope(&["foo".into()])
5492 .expect("should create foo consumer")
5493 .announced();
5494
5495 let mut bar_consumer = origin
5496 .consume()
5497 .scope(&["bar".into()])
5498 .expect("should create bar consumer")
5499 .announced();
5500
5501 let mut foobar_consumer = origin
5502 .consume()
5503 .scope(&["foo".into(), "bar".into()])
5504 .expect("should create foobar consumer")
5505 .announced();
5506
5507 foo_consumer.assert_next_some("foo/test");
5509 foo_consumer.assert_next_wait();
5510
5511 bar_consumer.assert_next_some("bar/test");
5512 bar_consumer.assert_next_wait();
5513
5514 foobar_consumer.assert_next_some("bar/test");
5515 foobar_consumer.assert_next_some("foo/test");
5516 foobar_consumer.assert_next_wait();
5517 }
5518
5519 #[tokio::test]
5520 async fn test_select_with_empty_prefix() {
5521 let origin = Origin::random().produce();
5522
5523 let demo_producer = origin.with_root("demo").expect("should create demo root");
5525 let limited_producer = demo_producer
5526 .scope(&["worm-node".into(), "foobar".into()])
5527 .expect("should create limited producer");
5528
5529 let _broadcast1 = limited_producer
5531 .create_broadcast("worm-node/test", announce())
5532 .expect("publish allowed");
5533 let _broadcast2 = limited_producer
5534 .create_broadcast("foobar/test", announce())
5535 .expect("publish allowed");
5536 settle().await;
5537
5538 let mut consumer = limited_producer
5540 .consume()
5541 .scope(&["".into()])
5542 .expect("should create consumer with empty prefix")
5543 .announced();
5544
5545 let a1 = consumer.try_next().expect("expected first announcement");
5547 let a2 = consumer.try_next().expect("expected second announcement");
5548 consumer.assert_next_wait();
5549
5550 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5551 paths.sort();
5552 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
5553 }
5554
5555 #[tokio::test]
5556 async fn test_select_narrowing_scope() {
5557 let origin = Origin::random().produce();
5558
5559 let demo_producer = origin.with_root("demo").expect("should create demo root");
5561 let limited_producer = demo_producer
5562 .scope(&["worm-node".into(), "foobar".into()])
5563 .expect("should create limited producer");
5564
5565 let _broadcast1 = limited_producer
5567 .create_broadcast("worm-node", announce())
5568 .expect("publish allowed");
5569 let _broadcast2 = limited_producer
5570 .create_broadcast("worm-node/foo", announce())
5571 .expect("publish allowed");
5572 let _broadcast3 = limited_producer
5573 .create_broadcast("foobar/bar", announce())
5574 .expect("publish allowed");
5575 settle().await;
5576
5577 let mut worm_consumer = limited_producer
5579 .consume()
5580 .scope(&["worm-node".into()])
5581 .expect("should create worm-node consumer")
5582 .announced();
5583
5584 worm_consumer.assert_next_some("worm-node");
5586 worm_consumer.assert_next_some("worm-node/foo");
5587 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
5591 .consume()
5592 .scope(&["worm-node/foo".into()])
5593 .expect("should create worm-node/foo consumer")
5594 .announced();
5595
5596 foo_consumer.assert_next_some("worm-node/foo");
5597 foo_consumer.assert_next_wait(); }
5599
5600 #[tokio::test]
5601 async fn test_select_multiple_roots_with_empty_prefix() {
5602 let origin = Origin::random().produce();
5603
5604 let limited_producer = origin
5606 .scope(&["app1".into(), "app2".into(), "shared".into()])
5607 .expect("should create limited producer");
5608
5609 let _broadcast1 = limited_producer
5611 .create_broadcast("app1/data", announce())
5612 .expect("publish allowed");
5613 let _broadcast2 = limited_producer
5614 .create_broadcast("app2/config", announce())
5615 .expect("publish allowed");
5616 let _broadcast3 = limited_producer
5617 .create_broadcast("shared/resource", announce())
5618 .expect("publish allowed");
5619 settle().await;
5620
5621 let mut consumer = limited_producer
5623 .consume()
5624 .scope(&["".into()])
5625 .expect("should create consumer with empty prefix")
5626 .announced();
5627
5628 consumer.assert_next_some("app1/data");
5630 consumer.assert_next_some("app2/config");
5631 consumer.assert_next_some("shared/resource");
5632 consumer.assert_next_wait();
5633 }
5634
5635 #[tokio::test]
5636 async fn test_publish_scope_with_empty_prefix() {
5637 let origin = Origin::random().produce();
5638
5639 let limited_producer = origin
5641 .scope(&["services/api".into(), "services/web".into()])
5642 .expect("should create limited producer");
5643
5644 let same_producer = limited_producer
5646 .scope(&["".into()])
5647 .expect("should create producer with empty prefix");
5648
5649 let _broadcast = same_producer
5651 .create_broadcast("services/api", announce())
5652 .expect("publish allowed");
5653 let _keep2 = same_producer
5654 .create_broadcast("services/web", announce())
5655 .expect("publish allowed");
5656 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
5657 assert!(same_producer.create_broadcast("other", announce()).is_err());
5658 }
5659
5660 #[tokio::test]
5661 async fn test_select_narrowing_to_deeper_path() {
5662 let origin = Origin::random().produce();
5663
5664 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
5666
5667 let _broadcast1 = limited_producer
5669 .create_broadcast("org/team1/project1", announce())
5670 .expect("publish allowed");
5671 let _broadcast2 = limited_producer
5672 .create_broadcast("org/team1/project2", announce())
5673 .expect("publish allowed");
5674 let _broadcast3 = limited_producer
5675 .create_broadcast("org/team2/project1", announce())
5676 .expect("publish allowed");
5677 settle().await;
5678
5679 let mut team2_consumer = limited_producer
5681 .consume()
5682 .scope(&["org/team2".into()])
5683 .expect("should create team2 consumer")
5684 .announced();
5685
5686 team2_consumer.assert_next_some("org/team2/project1");
5687 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
5691 .consume()
5692 .scope(&["org/team1/project1".into()])
5693 .expect("should create project1 consumer")
5694 .announced();
5695
5696 project1_consumer.assert_next_some("org/team1/project1");
5698 project1_consumer.assert_next_wait();
5699 }
5700
5701 #[tokio::test]
5702 async fn test_select_with_non_matching_prefix() {
5703 let origin = Origin::random().produce();
5704
5705 let limited_producer = origin
5707 .scope(&["allowed/path".into()])
5708 .expect("should create limited producer");
5709
5710 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
5712
5713 assert!(limited_producer.scope(&["other/path".into()]).is_none());
5715 }
5716
5717 #[tokio::test]
5720 async fn test_with_root_trailing_slash_consumer() {
5721 let origin = Origin::random().produce();
5722
5723 let prefix = "some_prefix/".to_string();
5725 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5726
5727 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5728 settle().await;
5729 consumer.assert_next_some("test");
5730 }
5731
5732 #[tokio::test]
5734 async fn test_with_root_trailing_slash_producer() {
5735 let origin = Origin::random().produce();
5736
5737 let prefix = "some_prefix/".to_string();
5739 let rooted = origin.with_root(prefix).unwrap();
5740
5741 let _b = rooted.create_broadcast("test", announce()).unwrap();
5742 settle().await;
5743
5744 let mut consumer = rooted.consume().announced();
5745 consumer.assert_next_some("test");
5746 }
5747
5748 #[tokio::test]
5750 async fn test_with_root_trailing_slash_unannounce() {
5751 tokio::time::pause();
5752
5753 let origin = Origin::random().produce();
5754
5755 let prefix = "some_prefix/".to_string();
5756 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5757
5758 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5759 settle().await;
5760 consumer.assert_next_some("test");
5761
5762 b.finish();
5764 settle().await;
5765
5766 consumer.assert_next_none("test");
5768 }
5769
5770 #[tokio::test]
5771 async fn test_select_maintains_access_with_wider_prefix() {
5772 let origin = Origin::random().produce();
5773
5774 let demo_producer = origin.with_root("demo").expect("should create demo root");
5776 let user_producer = demo_producer
5777 .scope(&["worm-node".into(), "foobar".into()])
5778 .expect("should create user producer");
5779
5780 let _broadcast1 = user_producer
5782 .create_broadcast("worm-node/data", announce())
5783 .expect("publish allowed");
5784 let _broadcast2 = user_producer
5785 .create_broadcast("foobar", announce())
5786 .expect("publish allowed");
5787 settle().await;
5788
5789 let mut consumer = user_producer
5791 .consume()
5792 .scope(&["".into()])
5793 .expect("scope with empty prefix should not fail when user has specific permissions")
5794 .announced();
5795
5796 let a1 = consumer.try_next().expect("expected first announcement");
5798 let a2 = consumer.try_next().expect("expected second announcement");
5799 consumer.assert_next_wait();
5800
5801 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5802 paths.sort();
5803 assert_eq!(paths, ["foobar", "worm-node/data"]);
5804
5805 let mut narrow_consumer = user_producer
5807 .consume()
5808 .scope(&["worm-node".into()])
5809 .expect("should be able to narrow scope to worm-node")
5810 .announced();
5811
5812 narrow_consumer.assert_next_some("worm-node/data");
5813 narrow_consumer.assert_next_wait(); }
5815
5816 #[tokio::test]
5817 async fn test_duplicate_prefixes_deduped() {
5818 let origin = Origin::random().produce();
5819
5820 let producer = origin
5822 .scope(&["demo".into(), "demo".into()])
5823 .expect("should create producer");
5824
5825 let _broadcast = producer
5826 .create_broadcast("demo/stream", announce())
5827 .expect("publish allowed");
5828 settle().await;
5829
5830 let mut consumer = producer.consume().announced();
5831 consumer.assert_next_some("demo/stream");
5832 consumer.assert_next_wait();
5833 }
5834
5835 #[tokio::test]
5836 async fn test_overlapping_prefixes_deduped() {
5837 let origin = Origin::random().produce();
5838
5839 let producer = origin
5841 .scope(&["demo".into(), "demo/foo".into()])
5842 .expect("should create producer");
5843
5844 let _broadcast = producer
5846 .create_broadcast("demo/bar/stream", announce())
5847 .expect("publish allowed");
5848 settle().await;
5849
5850 let mut consumer = producer.consume().announced();
5851 consumer.assert_next_some("demo/bar/stream");
5852 consumer.assert_next_wait();
5853 }
5854
5855 #[tokio::test]
5856 async fn test_overlapping_prefixes_no_duplicate_announcements() {
5857 let origin = Origin::random().produce();
5858
5859 let producer = origin
5861 .scope(&["demo".into(), "demo/foo".into()])
5862 .expect("should create producer");
5863
5864 let _broadcast = producer
5865 .create_broadcast("demo/foo/stream", announce())
5866 .expect("publish allowed");
5867 settle().await;
5868
5869 let mut consumer = producer.consume().announced();
5870 consumer.assert_next_some("demo/foo/stream");
5872 consumer.assert_next_wait();
5873 }
5874
5875 #[tokio::test]
5876 async fn test_allowed_returns_deduped_prefixes() {
5877 let origin = Origin::random().produce();
5878
5879 let producer = origin
5880 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
5881 .expect("should create producer");
5882
5883 let allowed: Vec<_> = producer.allowed().collect();
5884 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
5885 }
5886
5887 #[tokio::test]
5888 async fn test_announced_broadcast_already_announced() {
5889 let origin = Origin::random().produce();
5890
5891 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
5892 settle().await;
5893
5894 let consumer = origin.consume();
5895 let result = consumer.announced_broadcast("test").await.expect("should find it");
5896 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
5897 }
5898
5899 #[tokio::test]
5900 async fn test_announced_broadcast_delayed() {
5901 tokio::time::pause();
5902
5903 let origin = Origin::random().produce();
5904
5905 let consumer = origin.consume();
5906
5907 let wait = tokio::spawn({
5909 let consumer = consumer.clone();
5910 async move { consumer.announced_broadcast("test").await }
5911 });
5912
5913 tokio::task::yield_now().await;
5915
5916 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
5917 settle().await;
5918
5919 let result = wait.await.unwrap().expect("should find it");
5920 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
5921 }
5922
5923 #[tokio::test]
5924 async fn test_announced_broadcast_ignores_unrelated_paths() {
5925 tokio::time::pause();
5926
5927 let origin = Origin::random().produce();
5928
5929 let consumer = origin.consume();
5930
5931 let wait = tokio::spawn({
5932 let consumer = consumer.clone();
5933 async move { consumer.announced_broadcast("target").await }
5934 });
5935
5936 tokio::task::yield_now().await;
5937
5938 let _other = origin.create_broadcast("other", announce()).unwrap();
5940 settle().await;
5941 tokio::task::yield_now().await;
5942 assert!(!wait.is_finished(), "must not resolve on unrelated path");
5943
5944 let _target = origin.create_broadcast("target", announce()).unwrap();
5945 settle().await;
5946 let result = wait.await.unwrap().expect("should find target");
5947 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
5948 }
5949
5950 #[tokio::test]
5951 async fn test_announced_broadcast_skips_nested_paths() {
5952 tokio::time::pause();
5953
5954 let origin = Origin::random().produce();
5955
5956 let consumer = origin.consume();
5957
5958 let wait = tokio::spawn({
5959 let consumer = consumer.clone();
5960 async move { consumer.announced_broadcast("foo").await }
5961 });
5962
5963 tokio::task::yield_now().await;
5964
5965 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
5967 settle().await;
5968 tokio::task::yield_now().await;
5969 assert!(!wait.is_finished(), "must not resolve on a nested path");
5970
5971 let _exact = origin.create_broadcast("foo", announce()).unwrap();
5972 settle().await;
5973 let result = wait.await.unwrap().expect("should find foo exactly");
5974 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
5975 }
5976
5977 #[tokio::test]
5978 async fn test_announced_broadcast_disallowed() {
5979 let origin = Origin::random().produce();
5980 let limited = origin
5981 .consume()
5982 .scope(&["allowed".into()])
5983 .expect("should create limited");
5984
5985 assert!(limited.announced_broadcast("notallowed").await.is_none());
5987 }
5988
5989 #[tokio::test]
5990 async fn test_announced_broadcast_scope_too_narrow() {
5991 let origin = Origin::random().produce();
5994 let limited = origin
5995 .consume()
5996 .scope(&["foo/specific".into()])
5997 .expect("should create limited");
5998
5999 let result = limited
6001 .announced_broadcast("foo")
6002 .now_or_never()
6003 .expect("must not block");
6004 assert!(result.is_none());
6005 }
6006
6007 #[tokio::test]
6011 async fn test_coalesce_announce_then_unannounce() {
6012 tokio::time::pause();
6014
6015 let origin = Origin::random().produce();
6016 let mut announced = origin.consume().announced();
6017
6018 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6019 settle().await;
6020 broadcast.finish();
6021
6022 settle().await;
6023
6024 announced.assert_next_wait();
6025 }
6026
6027 #[tokio::test]
6028 async fn test_coalesce_announce_unannounce_announce() {
6029 tokio::time::pause();
6032
6033 let origin = Origin::random().produce();
6034 let mut announced = origin.consume().announced();
6035
6036 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6037 settle().await;
6038 broadcast1.finish();
6039 settle().await;
6040 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6041 settle().await;
6042
6043 announced.assert_next_some("test");
6044 announced.assert_next_wait();
6045 }
6046
6047 #[tokio::test]
6048 async fn test_coalesce_unannounce_announce_preserved() {
6049 tokio::time::pause();
6052
6053 let origin = Origin::random().produce();
6054 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6055 settle().await;
6056
6057 let mut announced = origin.consume().announced();
6058 announced.assert_next_some("test");
6059
6060 broadcast1.finish();
6062 settle().await;
6063
6064 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6065 settle().await;
6066
6067 announced.assert_next_none("test");
6069 announced.assert_next_some("test");
6070 announced.assert_next_wait();
6071 }
6072
6073 #[tokio::test]
6074 async fn test_coalesce_unannounce_announce_unannounce() {
6075 tokio::time::pause();
6078
6079 let origin = Origin::random().produce();
6080 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6081 settle().await;
6082
6083 let mut announced = origin.consume().announced();
6084 announced.assert_next_some("test");
6085
6086 broadcast1.finish();
6087 settle().await;
6088
6089 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6090 settle().await;
6091 broadcast2.finish();
6092 settle().await;
6093
6094 announced.assert_next_none("test");
6095 announced.assert_next_wait();
6096 }
6097
6098 #[tokio::test]
6099 async fn test_coalesce_churn_bounded() {
6100 tokio::time::pause();
6105
6106 let origin = Origin::random().produce();
6107 let mut announced = origin.consume().announced();
6108
6109 for _ in 0..1000 {
6110 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6111 settle().await;
6112 broadcast.finish();
6113 }
6114 settle().await;
6115
6116 let mut collected = Vec::new();
6117 while let Some(update) = announced.try_next() {
6118 collected.push(update);
6119 }
6120 assert!(
6121 collected.len() <= 1,
6122 "expected at most one pending update, got {}",
6123 collected.len()
6124 );
6125 assert!(
6126 collected.iter().all(|a| a.path == Path::new("test")),
6127 "unexpected path in pending updates",
6128 );
6129 }
6130
6131 #[tokio::test]
6135 async fn test_consumer_clone_is_side_effect_free() {
6136 let origin = Origin::random().produce();
6137
6138 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
6139 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
6140 settle().await;
6141
6142 let consumer = origin.consume();
6143 let mut announced = consumer.announced();
6144
6145 for _ in 0..16 {
6148 let cloned = consumer.clone();
6149 assert!(cloned.get_broadcast("test1").is_some());
6150 assert!(cloned.get_broadcast("test2").is_some());
6151 }
6152
6153 let a1 = announced.try_next().expect("first announcement");
6156 let a2 = announced.try_next().expect("second announcement");
6157 announced.assert_next_wait();
6158
6159 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6160 paths.sort();
6161 assert_eq!(paths, ["test1", "test2"]);
6162
6163 let mut fresh = consumer.announced();
6165 let b1 = fresh.try_next().expect("backlog: first");
6166 let b2 = fresh.try_next().expect("backlog: second");
6167 fresh.assert_next_wait();
6168
6169 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
6170 paths.sort();
6171 assert_eq!(paths, ["test1", "test2"]);
6172 }
6173
6174 #[tokio::test]
6176 async fn dynamic_request_unroutable_without_handler() {
6177 let origin = Origin::random().produce();
6178 let consumer = origin.consume();
6179 assert!(matches!(
6180 consumer.request_broadcast("missing").await,
6181 Err(Error::Unroutable)
6182 ));
6183 }
6184
6185 #[tokio::test(start_paused = true)]
6188 async fn dynamic_request_served_not_announced() {
6189 let origin = Origin::random().produce();
6190 let mut dynamic = origin.dynamic();
6191 let consumer = origin.consume();
6192
6193 let mut announced = origin.consume().announced();
6195 announced.assert_next_wait();
6196
6197 let served = broadcast::Info::new().produce();
6198 let request_fut = consumer.request_broadcast("fallback");
6201
6202 let mut served_dynamic = served.dynamic();
6204
6205 let request = dynamic.requested_broadcast().await.unwrap();
6206 assert_eq!(request.path(), &Path::new("fallback"));
6207 request.accept(&served);
6208
6209 let broadcast = request_fut.await.unwrap();
6210 assert!(broadcast.is_clone(&served.consume()));
6211
6212 let track_fut = broadcast.track("video").unwrap().subscribe(None);
6214 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
6215 let mut track = track_fut.await.unwrap();
6216 producer.append_group().unwrap();
6217 track.assert_group();
6218
6219 announced.assert_next_wait();
6221 }
6222
6223 #[tokio::test(start_paused = true)]
6225 async fn dynamic_request_coalesces() {
6226 let origin = Origin::random().produce();
6227 let mut dynamic = origin.dynamic();
6228 let consumer = origin.consume();
6229
6230 let f1 = consumer.request_broadcast("dup");
6232 let f2 = consumer.request_broadcast("dup");
6233
6234 let request = dynamic.requested_broadcast().await.unwrap();
6236 assert_eq!(request.path(), &Path::new("dup"));
6237 assert!(
6238 dynamic.requested_broadcast().now_or_never().is_none(),
6239 "a coalesced request must not be served twice"
6240 );
6241
6242 let served = broadcast::Info::new().produce();
6244 request.accept(&served);
6245 assert!(f1.await.unwrap().is_clone(&served.consume()));
6246 assert!(f2.await.unwrap().is_clone(&served.consume()));
6247 }
6248
6249 #[tokio::test(start_paused = true)]
6252 async fn dynamic_request_dedups_served() {
6253 let origin = Origin::random().produce();
6254 let mut dynamic = origin.dynamic();
6255 let consumer = origin.consume();
6256
6257 let request_fut = consumer.request_broadcast("fallback");
6258 let request = dynamic.requested_broadcast().await.unwrap();
6259 let served = broadcast::Info::new().produce();
6260 request.accept(&served);
6261 let first = request_fut.await.unwrap();
6262 assert!(first.is_clone(&served.consume()));
6263
6264 let second = consumer.request_broadcast("fallback").await.unwrap();
6266 assert!(second.is_clone(&served.consume()));
6267
6268 assert!(
6270 dynamic.requested_broadcast().now_or_never().is_none(),
6271 "a still-live served broadcast must not be re-requested from the handler"
6272 );
6273 }
6274
6275 #[tokio::test(start_paused = true)]
6277 async fn dynamic_request_reserves_after_close() {
6278 let origin = Origin::random().produce();
6279 let mut dynamic = origin.dynamic();
6280 let consumer = origin.consume();
6281
6282 let request_fut = consumer.request_broadcast("fallback");
6283 let request = dynamic.requested_broadcast().await.unwrap();
6284 let served = broadcast::Info::new().produce();
6285 request.accept(&served);
6286 request_fut.await.unwrap();
6287
6288 drop(served);
6290
6291 let request_fut = consumer.request_broadcast("fallback");
6293 let request = dynamic.requested_broadcast().await.unwrap();
6294 assert_eq!(request.path(), &Path::new("fallback"));
6295 let served = broadcast::Info::new().produce();
6296 request.accept(&served);
6297 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6298 }
6299
6300 #[tokio::test(start_paused = true)]
6303 async fn dynamic_request_served_cache_bounded() {
6304 let origin = Origin::random().produce();
6305 let mut dynamic = origin.dynamic();
6306 let consumer = origin.consume();
6307
6308 for i in 0..100 {
6309 let path = format!("one-shot/{i}");
6310 let request_fut = consumer.request_broadcast(&path);
6311 let request = dynamic.requested_broadcast().await.unwrap();
6312 let served = broadcast::Info::new().produce();
6313 request.accept(&served);
6314 request_fut.await.unwrap();
6315 drop(served);
6317 }
6318
6319 assert!(
6322 origin.dynamic.read().served.len() <= 4,
6323 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
6324 origin.dynamic.read().served.len()
6325 );
6326 }
6327
6328 #[tokio::test(start_paused = true)]
6331 async fn dynamic_request_coalesces_after_handoff() {
6332 let origin = Origin::random().produce();
6333 let mut dynamic = origin.dynamic();
6334 let consumer = origin.consume();
6335
6336 let f1 = consumer.request_broadcast("fallback");
6337 let request = dynamic.requested_broadcast().await.unwrap();
6339
6340 let f2 = consumer.request_broadcast("fallback");
6342 assert!(
6343 dynamic.requested_broadcast().now_or_never().is_none(),
6344 "a repeat request during hand-off must coalesce, not re-queue"
6345 );
6346
6347 let served = broadcast::Info::new().produce();
6349 request.accept(&served);
6350 assert!(f1.await.unwrap().is_clone(&served.consume()));
6351 assert!(f2.await.unwrap().is_clone(&served.consume()));
6352 }
6353
6354 #[tokio::test(start_paused = true)]
6356 async fn dynamic_request_dropped_after_handoff() {
6357 let origin = Origin::random().produce();
6358 let mut dynamic = origin.dynamic();
6359 let consumer = origin.consume();
6360
6361 let f1 = consumer.request_broadcast("fallback");
6362 let request = dynamic.requested_broadcast().await.unwrap();
6363 let f2 = consumer.request_broadcast("fallback");
6364
6365 drop(request);
6367 assert!(matches!(f1.await, Err(Error::Unroutable)));
6368 assert!(matches!(f2.await, Err(Error::Unroutable)));
6369 }
6370
6371 #[tokio::test(start_paused = true)]
6373 async fn dynamic_request_rejected() {
6374 let origin = Origin::random().produce();
6375 let mut dynamic = origin.dynamic();
6376 let consumer = origin.consume();
6377
6378 let request_fut = consumer.request_broadcast("fallback");
6379
6380 let request = dynamic.requested_broadcast().await.unwrap();
6381 request.reject(Error::Cancel);
6382
6383 assert!(matches!(request_fut.await, Err(Error::Cancel)));
6384 }
6385
6386 #[tokio::test(start_paused = true)]
6390 async fn dynamic_request_rerequest_after_reject() {
6391 let origin = Origin::random().produce();
6392 let mut dynamic = origin.dynamic();
6393 let consumer = origin.consume();
6394
6395 let f1 = consumer.request_broadcast("fallback");
6396 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
6397 assert!(matches!(f1.await, Err(Error::Unroutable)));
6398
6399 let served = broadcast::Info::new().produce();
6400 let f2 = consumer.request_broadcast("fallback");
6402 let request = dynamic.requested_broadcast().await.unwrap();
6403 assert_eq!(request.path(), &Path::new("fallback"));
6404 request.accept(&served);
6405 assert!(f2.await.unwrap().is_clone(&served.consume()));
6406 }
6407
6408 #[tokio::test(start_paused = true)]
6411 async fn dynamic_request_handler_dropped() {
6412 let origin = Origin::random().produce();
6413 let dynamic = origin.dynamic();
6414 let consumer = origin.consume();
6415
6416 let request_fut = consumer.request_broadcast("fallback");
6417 drop(dynamic);
6418 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
6419
6420 assert!(matches!(
6422 consumer.request_broadcast("again").await,
6423 Err(Error::Unroutable)
6424 ));
6425 }
6426
6427 #[tokio::test(start_paused = true)]
6431 async fn dynamic_request_accept_after_handler_dropped() {
6432 let origin = Origin::random().produce();
6433 let mut dynamic = origin.dynamic();
6434 let consumer = origin.consume();
6435
6436 let request_fut = consumer.request_broadcast("fallback");
6437
6438 let request = dynamic.requested_broadcast().await.unwrap();
6440 drop(dynamic);
6441
6442 let served = broadcast::Info::new().produce();
6443 request.accept(&served);
6445 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6446 }
6447
6448 #[tokio::test(start_paused = true)]
6450 async fn dynamic_request_prefers_announced() {
6451 let origin = Origin::random().produce();
6452 let mut dynamic = origin.dynamic();
6453 let consumer = origin.consume();
6454
6455 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
6456 settle().await;
6457
6458 let got = consumer.request_broadcast("live").await.unwrap();
6459 assert!(
6460 got.is_clone(&consumer.get_broadcast("live").unwrap()),
6461 "should return the published broadcast"
6462 );
6463 assert!(
6464 dynamic.requested_broadcast().now_or_never().is_none(),
6465 "a published path must not queue a fallback request"
6466 );
6467 }
6468
6469 #[tokio::test(start_paused = true)]
6471 async fn dynamic_clone_keeps_alive() {
6472 let origin = Origin::random().produce();
6473 let dynamic = origin.dynamic();
6474 let consumer = origin.consume();
6475
6476 drop(dynamic.clone());
6477
6478 let request_fut = consumer.request_broadcast("fallback");
6481 assert!(
6482 request_fut.now_or_never().is_none(),
6483 "request should stay pending until served"
6484 );
6485 }
6486}