1use crate::{broadcast, cache, stats, track};
2use kio::Pollable;
3use std::{
4 collections::{BTreeMap, HashMap},
5 fmt,
6 sync::Arc,
7 sync::atomic::{AtomicU64, Ordering},
8 task::{Poll, ready},
9 time::Duration,
10};
11
12use rand::RngExt;
13use web_async::Lock;
14
15use super::{Requests, WeakCache};
16use crate::{
17 AsPath, Error, Path, PathOwned, PathPrefixes,
18 coding::{BoundsExceeded, Decode, DecodeError, Encode, EncodeError},
19};
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
27pub struct Origin {
28 id: u64,
30}
31
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34#[non_exhaustive]
35pub struct InvalidOrigin;
36
37impl fmt::Display for InvalidOrigin {
38 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
39 write!(f, "local origin id must be non-zero and below 2^62")
40 }
41}
42
43impl std::error::Error for InvalidOrigin {}
44
45impl Origin {
46 pub(crate) const UNKNOWN: Self = Self { id: 0 };
49
50 pub fn new(id: u64) -> Result<Self, InvalidOrigin> {
56 if id == 0 || id >= 1u64 << 62 {
57 return Err(InvalidOrigin);
58 }
59 Ok(Self { id })
60 }
61
62 pub fn random() -> Self {
71 let mut rng = rand::rng();
72 let id = rng.random_range(1..(1u64 << 53));
73 Self { id }
74 }
75
76 pub fn id(self) -> u64 {
78 self.id
79 }
80
81 pub fn produce(self) -> Producer {
84 Info::new(self).produce()
85 }
86}
87
88#[derive(Clone, Debug)]
98#[non_exhaustive]
99pub struct Info {
100 pub id: Origin,
103
104 pub pool: cache::Pool,
110
111 pub cache_duration: Duration,
118
119 pub linger: Duration,
128}
129
130impl Default for Info {
131 fn default() -> Self {
134 Self {
135 id: Origin::UNKNOWN,
136 pool: cache::Pool::default(),
137 cache_duration: Duration::MAX,
138 linger: Duration::ZERO,
139 }
140 }
141}
142
143impl Info {
144 pub fn new(id: Origin) -> Self {
146 Self { id, ..Self::default() }
147 }
148
149 pub fn with_pool(mut self, pool: cache::Pool) -> Self {
151 self.pool = pool;
152 self
153 }
154
155 pub fn with_cache_duration(mut self, cache_duration: Duration) -> Self {
158 self.cache_duration = cache_duration;
159 self
160 }
161
162 pub fn with_linger(mut self, linger: Duration) -> Self {
165 self.linger = linger;
166 self
167 }
168
169 pub fn produce(self) -> Producer {
171 Producer::new(self)
172 }
173}
174
175impl TryFrom<u64> for Origin {
176 type Error = InvalidOrigin;
177
178 fn try_from(id: u64) -> Result<Self, Self::Error> {
179 Self::new(id)
180 }
181}
182
183impl fmt::Display for Origin {
184 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
185 self.id.fmt(f)
186 }
187}
188
189impl<V: Copy> Encode<V> for Origin
190where
191 u64: Encode<V>,
192{
193 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
194 self.id.encode(w, version)
195 }
196}
197
198impl<V: Copy> Decode<V> for Origin
199where
200 u64: Decode<V>,
201{
202 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
203 let id = u64::decode(r, version)?;
204 if id >= 1u64 << 62 {
205 return Err(DecodeError::InvalidValue);
206 }
207 Ok(Self { id })
208 }
209}
210
211pub(crate) const MAX_HOPS: usize = 32;
217
218#[derive(Debug, Clone, Default, PartialEq, Eq)]
223pub struct OriginList(Vec<Origin>);
224
225#[derive(Debug, Clone, Copy, PartialEq, Eq)]
227#[non_exhaustive]
228pub struct TooManyOrigins;
229
230impl fmt::Display for TooManyOrigins {
231 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
232 write!(f, "too many origins (max {MAX_HOPS})")
233 }
234}
235
236impl std::error::Error for TooManyOrigins {}
237
238impl From<TooManyOrigins> for DecodeError {
239 fn from(_: TooManyOrigins) -> Self {
240 DecodeError::BoundsExceeded
241 }
242}
243
244impl OriginList {
245 pub fn new() -> Self {
247 Self(Vec::new())
248 }
249
250 pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
252 if self.0.len() >= MAX_HOPS {
253 return Err(TooManyOrigins);
254 }
255 self.0.push(origin);
256 Ok(())
257 }
258
259 pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
262 for entry in &mut self.0 {
263 if *entry == target {
264 *entry = replacement;
265 return true;
266 }
267 }
268 false
269 }
270
271 pub fn contains(&self, origin: &Origin) -> bool {
273 self.0.contains(origin)
274 }
275
276 pub fn len(&self) -> usize {
278 self.0.len()
279 }
280
281 pub fn is_empty(&self) -> bool {
283 self.0.is_empty()
284 }
285
286 pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
288 self.0.iter()
289 }
290
291 pub fn as_slice(&self) -> &[Origin] {
293 &self.0
294 }
295}
296
297impl TryFrom<Vec<Origin>> for OriginList {
298 type Error = TooManyOrigins;
299
300 fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
301 if v.len() > MAX_HOPS {
302 return Err(TooManyOrigins);
303 }
304 Ok(Self(v))
305 }
306}
307
308impl<'a> IntoIterator for &'a OriginList {
309 type Item = &'a Origin;
310 type IntoIter = std::slice::Iter<'a, Origin>;
311
312 fn into_iter(self) -> Self::IntoIter {
313 self.iter()
314 }
315}
316
317impl<V: Copy> Encode<V> for OriginList
318where
319 u64: Encode<V>,
320 Origin: Encode<V>,
321{
322 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
323 (self.0.len() as u64).encode(w, version)?;
324 for origin in &self.0 {
325 origin.encode(w, version)?;
326 }
327 Ok(())
328 }
329}
330
331impl<V: Copy> Decode<V> for OriginList
332where
333 u64: Decode<V>,
334 Origin: Decode<V>,
335{
336 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
337 let count = u64::decode(r, version)? as usize;
338 if count > MAX_HOPS {
339 return Err(DecodeError::BoundsExceeded);
340 }
341 let mut list = Vec::with_capacity(count);
342 for _ in 0..count {
343 list.push(Origin::decode(r, version)?);
344 }
345 Ok(Self(list))
346 }
347}
348
349static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
350
351#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
352struct ConsumerId(u64);
353
354impl ConsumerId {
355 fn new() -> Self {
356 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
357 }
358}
359
360struct OriginBroadcast {
366 path: PathOwned,
367 broadcast: broadcast::Producer,
369 state: kio::Producer<FrontState>,
372 announced: bool,
373}
374
375fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
384 (hops.len(), fnv_key(name, hops.iter().copied()))
385}
386
387fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
400 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
402
403 let mut hash = SEED;
404 for &byte in name.as_str().as_bytes() {
405 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
406 }
407 for origin in origins {
408 for &byte in &origin.id().to_le_bytes() {
409 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
410 }
411 }
412
413 hash
414}
415
416fn route_order(name: &Path, route: &broadcast::Route) -> (bool, u64, usize, u64) {
425 let (len, hash) = route_key(name, &route.hops);
426 (!route.announce, route.cost, len, hash)
427}
428
429enum PendingUpdate {
438 Announce(broadcast::Consumer),
439 Unannounce,
440 UnannounceAnnounce(broadcast::Consumer),
441}
442
443#[derive(Default)]
448struct OriginConsumerState {
449 pending: BTreeMap<PathOwned, PendingUpdate>,
450}
451
452impl OriginConsumerState {
453 fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
454 let new = match self.pending.remove(&path) {
455 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
457 Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
459 PendingUpdate::UnannounceAnnounce(broadcast)
460 }
461 };
462 self.pending.insert(path, new);
463 }
464
465 fn apply_unannounce(&mut self, path: PathOwned) {
466 match self.pending.remove(&path) {
467 Some(PendingUpdate::Announce(_)) => {}
469 None | Some(PendingUpdate::Unannounce) => {
470 self.pending.insert(path, PendingUpdate::Unannounce);
471 }
472 Some(PendingUpdate::UnannounceAnnounce(_)) => {
475 self.pending.insert(path, PendingUpdate::Unannounce);
476 }
477 }
478 }
479
480 fn take(&mut self) -> Option<OriginAnnounce> {
482 let path = self.pending.keys().next()?.clone();
483 let broadcast = match self.pending.remove(&path).unwrap() {
484 PendingUpdate::Announce(broadcast) => Some(broadcast),
485 PendingUpdate::Unannounce => None,
486 PendingUpdate::UnannounceAnnounce(broadcast) => {
487 self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
490 None
491 }
492 };
493 Some(OriginAnnounce { path, broadcast })
494 }
495}
496
497#[derive(Clone)]
498struct AnnounceConsumerNotify {
499 root: PathOwned,
500 state: kio::Producer<OriginConsumerState>,
501}
502
503impl AnnounceConsumerNotify {
504 fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer) {
505 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
506 self.state
507 .write()
508 .ok()
509 .expect("consumer closed")
510 .apply_announce(path, broadcast);
511 }
512
513 fn unannounce(&self, path: impl AsPath) {
514 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
515 self.state.write().ok().expect("consumer closed").apply_unannounce(path);
516 }
517}
518
519struct NotifyNode {
520 parent: Option<Lock<NotifyNode>>,
521
522 consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
525}
526
527impl NotifyNode {
528 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
529 Self {
530 parent,
531 consumers: HashMap::new(),
532 }
533 }
534
535 fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer) {
536 for consumer in self.consumers.values() {
537 consumer.announce(path.as_path(), broadcast.clone());
538 }
539
540 if let Some(parent) = &self.parent {
541 parent.lock().announce(path, broadcast);
542 }
543 }
544
545 fn unannounce(&mut self, path: impl AsPath) {
546 for consumer in self.consumers.values() {
547 consumer.unannounce(path.as_path());
548 }
549
550 if let Some(parent) = &self.parent {
551 parent.lock().unannounce(path);
552 }
553 }
554}
555
556struct OriginNode {
557 broadcast: Option<OriginBroadcast>,
560
561 nested: HashMap<String, Lock<OriginNode>>,
563
564 notify: Lock<NotifyNode>,
566}
567
568impl OriginNode {
569 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
570 Self {
571 broadcast: None,
572 nested: HashMap::new(),
573 notify: Lock::new(NotifyNode::new(parent)),
574 }
575 }
576
577 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
578 let (dir, rest) = path.next_part().expect("leaf called with empty path");
579
580 let next = self.entry(dir);
581 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
582 }
583
584 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
585 match self.nested.get(dir) {
586 Some(next) => next.clone(),
587 None => {
588 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
589 self.nested.insert(dir.to_string(), next.clone());
590 next
591 }
592 }
593 }
594
595 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
599 let Some(existing) = &mut self.broadcast else { return };
600 if !existing.state.same_channel(expect) || existing.announced == announce {
601 return;
602 }
603 existing.announced = announce;
604 let path = existing.path.clone();
605 let consumer = existing.broadcast.consume();
606 let mut notify = self.notify.lock();
607 if announce {
608 notify.announce(&path, &consumer);
609 } else {
610 notify.unannounce(&path);
611 }
612 }
613
614 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
615 self.consume_initial(&mut notify);
616 self.notify.lock().consumers.insert(id, notify);
617 }
618
619 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
620 if let Some(broadcast) = &self.broadcast
623 && broadcast.announced
624 {
625 notify.announce(&broadcast.path, broadcast.broadcast.consume());
626 }
627
628 for nested in self.nested.values() {
630 nested.lock().consume_initial(notify);
631 }
632 }
633
634 fn consume_broadcast(&self, rest: impl AsPath) -> Option<broadcast::Consumer> {
635 let rest = rest.as_path();
636
637 if let Some((dir, rest)) = rest.next_part() {
638 let node = self.nested.get(dir)?.lock();
639 node.consume_broadcast(&rest)
640 } else {
641 self.broadcast.as_ref().map(|b| b.broadcast.consume())
642 }
643 }
644
645 fn unconsume(&mut self, id: ConsumerId) {
646 self.notify.lock().consumers.remove(&id).expect("consumer not found");
647 if self.is_empty() {
648 }
651 }
652
653 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
657 let relative = relative.as_path();
658
659 if let Some((dir, relative)) = relative.next_part() {
660 let Some(nested) = self.nested.get(dir) else { return };
661 let nested = nested.clone();
662 let mut locked = nested.lock();
663 locked.remove(expect, &relative);
664
665 if locked.is_empty() {
666 drop(locked);
667 self.nested.remove(dir);
668 }
669 } else if let Some(existing) = &self.broadcast
670 && existing.state.same_channel(expect)
671 {
672 let existing = self.broadcast.take().expect("checked above");
673 if existing.announced {
674 self.notify.lock().unannounce(&existing.path);
675 }
676 }
677 }
678
679 fn is_empty(&self) -> bool {
680 self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
681 }
682}
683
684#[derive(Clone)]
685struct OriginNodes {
686 nodes: Vec<(PathOwned, Lock<OriginNode>)>,
687}
688
689impl OriginNodes {
690 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
693 let mut roots = Vec::new();
694
695 for (root, state) in &self.nodes {
696 for prefix in prefixes {
697 if root.has_prefix(prefix) {
698 roots.push((root.to_owned(), state.clone()));
700 continue;
701 }
702
703 if let Some(suffix) = prefix.strip_prefix(root) {
704 let nested = state.lock().leaf(&suffix);
706 roots.push((prefix.to_owned(), nested));
707 }
708 }
709 }
710
711 if roots.is_empty() {
712 None
713 } else {
714 Some(Self { nodes: roots })
715 }
716 }
717
718 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
719 let new_root = new_root.as_path();
720 let mut roots = Vec::new();
721
722 if new_root.is_empty() {
723 return Some(self.clone());
724 }
725
726 for (root, state) in &self.nodes {
727 if let Some(suffix) = root.strip_prefix(&new_root) {
728 roots.push((suffix.to_owned(), state.clone()));
730 } else if let Some(suffix) = new_root.strip_prefix(root) {
731 let nested = state.lock().leaf(&suffix);
734 roots.push(("".into(), nested));
735 }
736 }
737
738 if roots.is_empty() {
739 None
740 } else {
741 Some(Self { nodes: roots })
742 }
743 }
744
745 pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
747 let path = path.as_path();
748
749 for (root, state) in &self.nodes {
750 if let Some(suffix) = path.strip_prefix(root) {
751 return Some((state.clone(), suffix.to_owned()));
752 }
753 }
754
755 None
756 }
757}
758
759impl Default for OriginNodes {
760 fn default() -> Self {
761 Self {
762 nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
763 }
764 }
765}
766
767#[derive(Clone)]
769pub struct OriginAnnounce {
770 pub path: PathOwned,
772 pub broadcast: Option<broadcast::Consumer>,
778}
779
780#[derive(Clone)]
782pub struct Producer {
783 info: Origin,
787
788 nodes: OriginNodes,
791
792 root: PathOwned,
794
795 dynamic: kio::Shared<OriginDynamicState>,
799
800 pool: cache::Pool,
803
804 cache_duration: Duration,
807
808 linger: Duration,
811
812 stats: stats::Session,
816}
817
818impl std::ops::Deref for Producer {
819 type Target = Origin;
820
821 fn deref(&self) -> &Self::Target {
822 &self.info
823 }
824}
825
826impl Producer {
827 pub fn new(info: Info) -> Self {
831 Self {
832 info: info.id,
833 nodes: OriginNodes::default(),
834 root: PathOwned::default(),
835 dynamic: kio::Shared::default(),
836 pool: info.pool,
837 cache_duration: info.cache_duration,
838 linger: info.linger,
839 stats: stats::Session::default(),
840 }
841 }
842
843 pub fn with_stats(mut self, session: stats::Session) -> Self {
847 self.stats = session;
848 self
849 }
850
851 pub fn with_linger(mut self, linger: Duration) -> Self {
860 self.linger = linger;
861 self
862 }
863
864 pub fn info(&self) -> Info {
867 Info {
868 id: self.info,
869 pool: self.pool.clone(),
870 cache_duration: self.cache_duration,
871 linger: self.linger,
872 }
873 }
874
875 pub(crate) fn empty(info: Origin) -> Self {
880 Self {
881 info,
882 nodes: OriginNodes { nodes: Vec::new() },
883 root: PathOwned::default(),
884 dynamic: kio::Shared::default(),
885 pool: cache::Pool::default(),
886 cache_duration: Duration::MAX,
887 linger: Duration::ZERO,
888 stats: stats::Session::default(),
889 }
890 }
891
892 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
931 let path = path.as_path();
932
933 debug_assert!(
934 !route.hops.contains(&self.info),
935 "create_broadcast called with a looping hop chain",
936 );
937
938 let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
939 let full = self.root.join(&path).to_owned();
940
941 if full.parts().count() > Path::MAX_PARTS {
945 return Err(BoundsExceeded.into());
946 }
947
948 let ingress = self.stats.ingress(&full);
952
953 let mut source = broadcast::Info { origin: self.info() }
954 .produce()
955 .with_stats(ingress.clone());
956 source.set_route(route).expect("fresh producer");
957
958 web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
959
960 Ok(source)
961 }
962
963 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
969 let prefixes = PathPrefixes::new(prefixes);
970 Some(Producer {
971 info: self.info,
972 nodes: self.nodes.select(&prefixes)?,
973 root: self.root.clone(),
974 dynamic: self.dynamic.clone(),
975 pool: self.pool.clone(),
976 cache_duration: self.cache_duration,
977 linger: self.linger,
978 stats: self.stats.clone(),
979 })
980 }
981
982 pub fn dynamic(&self) -> Dynamic {
991 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
992 }
993
994 pub fn consume(&self) -> Consumer {
999 Consumer::new(
1002 self.info,
1003 self.root.clone(),
1004 self.nodes.clone(),
1005 self.dynamic.clone(),
1006 stats::Session::default(),
1007 )
1008 }
1009
1010 pub fn announces(&self) -> AnnounceProducer {
1016 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
1017 }
1018
1019 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1024 let prefix = prefix.as_path();
1025
1026 Some(Self {
1027 info: self.info,
1028 root: self.root.join(&prefix).to_owned(),
1029 nodes: self.nodes.root(&prefix)?,
1030 dynamic: self.dynamic.clone(),
1031 pool: self.pool.clone(),
1032 cache_duration: self.cache_duration,
1033 linger: self.linger,
1034 stats: self.stats.clone(),
1035 })
1036 }
1037
1038 pub fn root(&self) -> &Path<'_> {
1040 &self.root
1041 }
1042
1043 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1046 self.nodes.nodes.iter().map(|(root, _)| root)
1047 }
1048
1049 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1051 self.root.join(path)
1052 }
1053}
1054
1055const MAX_TRACK_RETRIES: u32 = 3;
1060
1061const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1074
1075struct FrontRoute {
1077 id: u64,
1078 route: broadcast::Route,
1081 source: broadcast::Consumer,
1083}
1084
1085struct FrontState {
1087 path: PathOwned,
1089 self_origin: Origin,
1091 next_route: u64,
1092 routes: Vec<FrontRoute>,
1093 active: Option<u64>,
1095 linger: Duration,
1098 closed: bool,
1102}
1103
1104impl FrontState {
1105 fn best_route(&self) -> Option<u64> {
1108 self.routes
1109 .iter()
1110 .min_by_key(|r| route_order(&self.path.as_path(), &r.route))
1111 .map(|r| r.id)
1112 }
1113
1114 fn reselect(&mut self, carrying: bool) {
1132 let best = self.best_route();
1133 if carrying
1134 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1135 && best_id != cur_id
1136 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1137 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1138 && incumbent.route.announce
1139 && candidate.route.cost < incumbent.route.cost
1140 && candidate.route.advertised == 0
1141 && candidate.route.hops.len() >= 2
1142 && !self.handover_allowed(&candidate.route)
1143 {
1144 return;
1146 }
1147 self.active = best;
1148 }
1149
1150 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1162 let name = self.path.as_path();
1163 match route.hops.iter().last() {
1164 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1165 None => true,
1166 }
1167 }
1168
1169 fn active_route(&self) -> Option<broadcast::Route> {
1172 let id = self.active?;
1173 self.routes.iter().find(|r| r.id == id).map(|r| r.route.clone())
1174 }
1175}
1176
1177fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1187 let mut leaf_guard = leaf.lock();
1192 let advert = state.read().active_route();
1193 if let Some(advert) = advert {
1194 let announce = advert.announce;
1195 let _ = broadcast.clone().set_route(advert);
1196 leaf_guard.set_announced(state, announce);
1197 }
1198}
1199
1200fn detach_source(
1213 state: &kio::Producer<FrontState>,
1214 broadcast: &broadcast::Producer,
1215 leaf: &Lock<OriginNode>,
1216 id: u64,
1217 graceful: bool,
1218) {
1219 let close = {
1220 let carrying = broadcast.demand().is_used();
1221 let Ok(mut s) = state.write() else { return };
1222 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1223 return;
1224 };
1225 s.routes.remove(pos);
1226 s.reselect(carrying);
1227 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1228 s.closed = true;
1231 true
1232 } else {
1233 false
1234 }
1235 };
1236 if close {
1237 broadcast.abort_spliced(Error::Dropped);
1238 }
1239 sync_front(state, broadcast, leaf);
1240}
1241
1242async fn run_source(
1246 origin: Info,
1247 node: Lock<OriginNode>,
1248 full: PathOwned,
1249 rest: PathOwned,
1250 mut source: broadcast::Consumer,
1251 ingress: stats::Scope,
1252) {
1253 let Ok(route) = source.route_changed().await else {
1257 return;
1259 };
1260
1261 let mut announce = route.announce.then(|| ingress.announce());
1266
1267 let leaf = if rest.is_empty() {
1268 node.clone()
1269 } else {
1270 node.lock().leaf(&rest)
1271 };
1272
1273 let (state, broadcast, id) = attach_source(&origin, &node, &leaf, &full, &rest, &source, route);
1274
1275 loop {
1276 match source.route_changed().await {
1277 Ok(route) => {
1278 let announced = route.announce;
1279 {
1280 let carrying = broadcast.demand().is_used();
1281 let Ok(mut s) = state.write() else { return };
1282 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1283 return;
1284 };
1285 if entry.route == route {
1286 continue;
1287 }
1288 entry.route = route;
1289 s.reselect(carrying);
1290 }
1291 match (announced, announce.is_some()) {
1293 (true, false) => announce = Some(ingress.announce()),
1294 (false, true) => announce = None,
1295 _ => {}
1296 }
1297 sync_front(&state, &broadcast, &leaf);
1298 }
1299 Err(_) => {
1300 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1303 return;
1304 }
1305 }
1306 }
1307}
1308
1309fn attach_source(
1314 origin: &Info,
1315 node: &Lock<OriginNode>,
1316 leaf: &Lock<OriginNode>,
1317 full: &PathOwned,
1318 rest: &PathOwned,
1319 source: &broadcast::Consumer,
1320 route: broadcast::Route,
1321) -> (kio::Producer<FrontState>, broadcast::Producer, u64) {
1322 let mut leaf_guard = leaf.lock();
1323
1324 if let Some(existing) = &leaf_guard.broadcast {
1327 let mut joined = None;
1328 let carrying = existing.broadcast.demand().is_used();
1329 if let Ok(mut s) = existing.state.write()
1330 && !s.closed
1331 {
1332 let id = s.next_route;
1333 s.next_route += 1;
1334 s.routes.push(FrontRoute {
1335 id,
1336 route: route.clone(),
1337 source: source.clone(),
1338 });
1339 s.reselect(carrying);
1340 joined = Some(id);
1341 }
1342 if let Some(id) = joined {
1343 let state = existing.state.clone();
1344 let broadcast = existing.broadcast.clone();
1345 drop(leaf_guard);
1346 sync_front(&state, &broadcast, leaf);
1347 return (state, broadcast, id);
1348 }
1349 }
1350
1351 let announce = route.announce;
1353 let broadcast = broadcast::Producer::new_spliced(broadcast::Info { origin: origin.clone() });
1354 let _ = broadcast.clone().set_route(route.clone());
1355 let state = kio::Producer::new(FrontState {
1356 path: full.clone(),
1357 self_origin: origin.id,
1358 next_route: 1,
1359 routes: vec![FrontRoute {
1360 id: 0,
1361 route,
1362 source: source.clone(),
1363 }],
1364 active: Some(0),
1365 linger: origin.linger,
1366 closed: false,
1367 });
1368
1369 if let Some(stale) = leaf_guard.broadcast.take()
1373 && stale.announced
1374 {
1375 leaf_guard.notify.lock().unannounce(&stale.path);
1376 }
1377 let entry = OriginBroadcast {
1378 path: full.clone(),
1379 broadcast: broadcast.clone(),
1380 state: state.clone(),
1381 announced: announce,
1382 };
1383 if entry.announced {
1384 leaf_guard.notify.lock().announce(full, &broadcast.consume());
1385 }
1386 leaf_guard.broadcast = Some(entry);
1387 drop(leaf_guard);
1388
1389 web_async::spawn(run_front(state.clone(), broadcast.clone(), node.clone(), rest.clone()));
1390
1391 (state, broadcast, 0)
1392}
1393
1394async fn run_front(
1397 state: kio::Producer<FrontState>,
1398 mut broadcast: broadcast::Producer,
1399 node: Lock<OriginNode>,
1400 rest: PathOwned,
1401) {
1402 enum Step {
1403 Serve(Arc<str>, super::resume::Producer),
1404 Changed,
1406 Expired,
1408 Closed,
1409 }
1410
1411 let linger = state.read().linger;
1412 let mut deadline = kio::time::Deadline::new();
1417
1418 loop {
1419 let empty = {
1420 let s = state.read();
1421 !s.closed && s.routes.is_empty()
1422 };
1423 deadline.set(match (empty, deadline.deadline()) {
1424 (true, None) => web_async::time::Instant::now().checked_add(linger),
1427 (true, at) => at,
1428 (false, _) => None,
1429 });
1430
1431 let step = {
1432 kio::wait(|waiter| {
1433 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1434 return Poll::Ready(Step::Serve(name, resume));
1435 }
1436 match state.poll(waiter, |s| {
1439 if s.closed || s.routes.is_empty() != empty {
1440 Poll::Ready(())
1441 } else {
1442 Poll::Pending
1443 }
1444 }) {
1445 Poll::Ready(Ok(guard)) => {
1446 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1447 }
1448 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1449 Poll::Pending => {}
1450 }
1451 deadline.poll(waiter).map(|_| Step::Expired)
1452 })
1453 .await
1454 };
1455
1456 match step {
1457 Step::Serve(name, resume) => {
1458 web_async::spawn(serve_track(state.clone(), name, resume));
1461 }
1462 Step::Changed => {}
1463 Step::Expired => {
1464 let close = {
1468 let Ok(mut s) = state.write() else { break };
1469 if !s.closed && s.routes.is_empty() {
1470 s.closed = true;
1471 true
1472 } else {
1473 false
1474 }
1475 };
1476 if close {
1477 break;
1478 }
1479 }
1480 Step::Closed => break,
1481 }
1482 }
1483
1484 broadcast.abort_spliced(Error::Dropped);
1486
1487 broadcast.finish();
1489
1490 node.lock().remove(&state, &rest);
1493}
1494
1495async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1503 enum Step {
1504 Closed,
1505 Splice(u64, broadcast::Consumer),
1506 Complete,
1507 Failed,
1508 Idle,
1510 Demand,
1512 }
1513
1514 let mut fails = 0u32;
1515 let mut serving: Option<(u64, track::Consumer)> = None;
1517 let mut dead: Option<u64> = None;
1521 let mut idle_since: Option<web_async::time::Instant> = None;
1523 let mut deadline = kio::time::Deadline::new();
1524
1525 loop {
1526 let serving_id = serving.as_ref().map(|(id, _)| *id);
1527
1528 let used = resume.is_used();
1532 idle_since = match (serving.is_some(), used) {
1533 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
1534 _ => None,
1535 };
1536 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
1537
1538 let step = {
1539 kio::wait(|waiter| {
1540 match state.poll(waiter, |s| {
1544 if s.closed
1545 || (used
1546 && matches!(s.active, Some(active) if Some(active) != serving_id && Some(active) != dead))
1547 {
1548 Poll::Ready(())
1549 } else {
1550 Poll::Pending
1551 }
1552 }) {
1553 Poll::Ready(Ok(guard)) => {
1554 if guard.closed {
1555 return Poll::Ready(Step::Closed);
1556 }
1557 let active = guard.active.expect("predicate guaranteed an active source");
1558 let source = guard
1559 .routes
1560 .iter()
1561 .find(|r| r.id == active)
1562 .expect("active source in table")
1563 .source
1564 .clone();
1565 return Poll::Ready(Step::Splice(active, source));
1566 }
1567 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1568 Poll::Pending => {}
1569 }
1570
1571 let edge = match used {
1576 true => resume.poll_unused(waiter),
1577 false => resume.poll_used(waiter),
1578 };
1579 if edge.is_ready() {
1580 return Poll::Ready(Step::Demand);
1581 }
1582
1583 if let Some((_, track)) = &serving
1586 && let Poll::Ready(result) = track.poll_complete(waiter)
1587 {
1588 return Poll::Ready(match result {
1589 Ok(()) => Step::Complete,
1590 Err(_) => Step::Failed,
1591 });
1592 }
1593
1594 deadline.poll(waiter).map(|_| Step::Idle)
1595 })
1596 .await
1597 };
1598
1599 match step {
1600 Step::Closed => return,
1602 Step::Complete => {
1603 let _ = resume.finish();
1604 return;
1605 }
1606 Step::Failed => {
1607 serving = None;
1610 }
1611 Step::Demand => {}
1613 Step::Idle => {
1614 if resume.release().is_err() {
1619 return;
1621 }
1622 serving = None;
1623 }
1624 Step::Splice(id, source) => {
1625 let attempt = match source.track(&name) {
1629 Ok(track) => {
1630 let query = track.info().into_inner();
1633 let info = kio::wait(|waiter| {
1634 if let Poll::Ready(result) = query.poll(waiter) {
1635 return Poll::Ready(Some(result));
1636 }
1637 match state.poll(waiter, |s| {
1638 if s.closed || s.active != Some(id) {
1639 Poll::Ready(())
1640 } else {
1641 Poll::Pending
1642 }
1643 }) {
1644 Poll::Ready(_) => Poll::Ready(None),
1645 Poll::Pending => Poll::Pending,
1646 }
1647 })
1648 .await;
1649 match info {
1650 None => continue,
1653 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
1657 Poll::Ready(Err(err)) => Err(err),
1658 _ => Ok(track),
1659 },
1660 Some(Err(err)) => Err(err),
1661 }
1662 }
1663 Err(err) => Err(err),
1664 };
1665
1666 match attempt {
1667 Ok(track) => {
1668 if resume.takeover(&track).is_err() {
1669 return;
1672 }
1673 fails = 0;
1676 dead = None;
1677 serving = Some((id, track));
1678 }
1679 Err(_) if source.is_closing() => {
1683 dead = Some(id);
1684 serving = None;
1685 }
1686 Err(err) => {
1687 fails += 1;
1688 if fails >= MAX_TRACK_RETRIES {
1689 tracing::debug!(name = %name, %err, "aborting unservable track");
1690 let _ = resume.abort(Error::Unroutable);
1691 return;
1692 }
1693 serving = None;
1694 }
1695 }
1696 }
1697 }
1698 }
1699}
1700
1701#[derive(Default)]
1707struct OriginDynamicState {
1708 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
1711
1712 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
1718}
1719
1720#[derive(Default)]
1727struct PendingBroadcast {
1728 resolved: Option<Result<broadcast::Consumer, Error>>,
1729}
1730
1731pub struct Dynamic {
1742 info: Origin,
1743 root: PathOwned,
1744 state: kio::Shared<OriginDynamicState>,
1745}
1746
1747impl Clone for Dynamic {
1748 fn clone(&self) -> Self {
1749 self.state.lock().requests.add_handler();
1753
1754 Self {
1755 info: self.info,
1756 root: self.root.clone(),
1757 state: self.state.clone(),
1758 }
1759 }
1760}
1761
1762impl Dynamic {
1763 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
1764 state.lock().requests.add_handler();
1765
1766 Self { info, root, state }
1767 }
1768
1769 pub fn info(&self) -> &Origin {
1771 &self.info
1772 }
1773
1774 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
1776 let mut state = ready!(self.state.poll(waiter, |state| {
1777 if state.requests.has_queued() {
1778 Poll::Ready(())
1779 } else {
1780 Poll::Pending
1781 }
1782 }));
1783
1784 let path = state.requests.pop().expect("predicate guaranteed a request");
1785 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
1791 Poll::Ready(Ok(Request {
1792 path,
1793 producer,
1794 state: self.state.clone(),
1795 }))
1796 }
1797
1798 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
1801 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
1802 }
1803
1804 pub fn root(&self) -> &Path<'_> {
1806 &self.root
1807 }
1808}
1809
1810impl Drop for Dynamic {
1811 fn drop(&mut self) {
1812 let mut state = self.state.lock();
1815 if state.requests.remove_handler() {
1816 state.requests.drain_queued();
1820 }
1821 }
1822}
1823
1824pub struct Request {
1831 path: PathOwned,
1833
1834 producer: kio::Producer<PendingBroadcast>,
1837
1838 state: kio::Shared<OriginDynamicState>,
1840}
1841
1842impl Request {
1843 pub fn path(&self) -> &Path<'_> {
1845 &self.path
1846 }
1847
1848 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
1854 let broadcast = broadcast.consume();
1855
1856 let resolved = {
1862 let mut state = self.state.lock();
1863 let existing = state.served.insert(self.path.clone(), broadcast.weak());
1864 state
1865 .requests
1866 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1867 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
1868 };
1869
1870 if let Ok(mut pending) = self.producer.write() {
1871 pending.resolved = Some(Ok(resolved));
1872 }
1873 }
1875
1876 pub fn reject(self, err: Error) {
1878 self.state
1879 .lock()
1880 .requests
1881 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1882 if let Ok(mut state) = self.producer.write() {
1883 state.resolved = Some(Err(err));
1884 }
1885 }
1886}
1887
1888impl Drop for Request {
1889 fn drop(&mut self) {
1890 self.state
1898 .lock()
1899 .requests
1900 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1901 }
1902}
1903
1904pub struct Requesting {
1911 inner: RequestState,
1912 stats: stats::Scope,
1915}
1916
1917enum RequestState {
1918 Ready(broadcast::Consumer),
1920 Failed(Error),
1923 Pending(kio::Consumer<PendingBroadcast>),
1925}
1926
1927impl Requesting {
1928 fn ready(broadcast: broadcast::Consumer) -> Self {
1929 Self {
1930 inner: RequestState::Ready(broadcast),
1931 stats: stats::Scope::default(),
1932 }
1933 }
1934
1935 fn failed(error: Error) -> Self {
1936 Self {
1937 inner: RequestState::Failed(error),
1938 stats: stats::Scope::default(),
1939 }
1940 }
1941
1942 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
1943 Self {
1944 inner: RequestState::Pending(consumer),
1945 stats: stats::Scope::default(),
1946 }
1947 }
1948
1949 fn with_stats(mut self, scope: stats::Scope) -> Self {
1950 self.stats = scope;
1951 self
1952 }
1953
1954 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
1956 match &self.inner {
1957 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
1958 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
1959 RequestState::Pending(consumer) => Poll::Ready(
1960 match ready!(consumer.poll(waiter, |state| match &state.resolved {
1961 Some(result) => Poll::Ready(result.clone()),
1962 None => Poll::Pending,
1963 })) {
1964 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
1965 Err(_closed) => Err(Error::Unroutable),
1967 },
1968 ),
1969 }
1970 }
1971}
1972
1973impl kio::Pollable for Requesting {
1974 type Output = Result<broadcast::Consumer, Error>;
1975
1976 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1977 self.poll_ok(waiter)
1978 }
1979}
1980
1981pub trait Consume<T> {
1989 fn consume(&self) -> T;
1991}
1992
1993impl<T, U: Consume<T>> Consume<T> for &U {
1994 fn consume(&self) -> T {
1995 (**self).consume()
1996 }
1997}
1998
1999impl Consume<Consumer> for Producer {
2000 fn consume(&self) -> Consumer {
2001 Consumer::new(
2005 self.info,
2006 self.root.clone(),
2007 self.nodes.clone(),
2008 self.dynamic.clone(),
2009 stats::Session::default(),
2010 )
2011 }
2012}
2013
2014impl Consume<Consumer> for Consumer {
2015 fn consume(&self) -> Consumer {
2016 self.clone()
2017 }
2018}
2019
2020impl Consume<broadcast::Consumer> for broadcast::Producer {
2021 fn consume(&self) -> broadcast::Consumer {
2022 self.consume()
2024 }
2025}
2026
2027impl Consume<broadcast::Consumer> for broadcast::Consumer {
2028 fn consume(&self) -> broadcast::Consumer {
2029 self.clone()
2030 }
2031}
2032
2033impl Consume<track::Consumer> for track::Producer {
2034 fn consume(&self) -> track::Consumer {
2035 self.consume()
2036 }
2037}
2038
2039impl Consume<track::Consumer> for track::Consumer {
2040 fn consume(&self) -> track::Consumer {
2041 self.clone()
2042 }
2043}
2044
2045#[derive(Clone)]
2051pub struct Consumer {
2052 info: Origin,
2054 nodes: OriginNodes,
2055
2056 root: PathOwned,
2058
2059 dynamic: kio::Shared<OriginDynamicState>,
2062
2063 stats: stats::Session,
2067}
2068
2069impl std::ops::Deref for Consumer {
2070 type Target = Origin;
2071
2072 fn deref(&self) -> &Self::Target {
2073 &self.info
2074 }
2075}
2076
2077impl Consumer {
2078 fn new(
2079 info: Origin,
2080 root: PathOwned,
2081 nodes: OriginNodes,
2082 dynamic: kio::Shared<OriginDynamicState>,
2083 stats: stats::Session,
2084 ) -> Self {
2085 Self {
2086 info,
2087 nodes,
2088 root,
2089 dynamic,
2090 stats,
2091 }
2092 }
2093
2094 pub fn with_stats(mut self, session: stats::Session) -> Self {
2098 self.stats = session;
2099 self
2100 }
2101
2102 fn untagged(&self) -> Self {
2106 Self {
2107 stats: stats::Session::default(),
2108 ..self.clone()
2109 }
2110 }
2111
2112 pub(crate) fn empty(&self) -> Self {
2117 Self {
2118 info: self.info,
2119 nodes: OriginNodes { nodes: Vec::new() },
2120 root: self.root.clone(),
2121 dynamic: self.dynamic.clone(),
2122 stats: self.stats.clone(),
2123 }
2124 }
2125
2126 pub fn announced(&self) -> AnnounceConsumer {
2133 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2134 }
2135
2136 pub fn consume(&self) -> Self {
2138 self.clone()
2139 }
2140
2141 fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2148 let path = path.as_path();
2149 let (root, rest) = self.nodes.get(&path)?;
2150 let state = root.lock();
2151 state.consume_broadcast(&rest)
2152 }
2153
2154 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2166 let path = path.as_path();
2167
2168 let consumer = self.scope(std::slice::from_ref(&path))?;
2170
2171 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2175 return None;
2176 }
2177
2178 let mut announced = consumer.untagged().announced();
2182 let scope = self.stats.egress(self.root.join(&path).to_owned());
2183 loop {
2184 let OriginAnnounce {
2185 path: announced_path,
2186 broadcast,
2187 } = announced.next().await?;
2188 if announced_path.as_path() == path
2190 && let Some(broadcast) = broadcast
2191 {
2192 return Some(broadcast.with_stats(scope));
2193 }
2194 }
2195 }
2196
2197 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2203 let prefixes = PathPrefixes::new(prefixes);
2204 Some(Consumer::new(
2205 self.info,
2206 self.root.clone(),
2207 self.nodes.select(&prefixes)?,
2208 self.dynamic.clone(),
2209 self.stats.clone(),
2210 ))
2211 }
2212
2213 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2232 let path = path.as_path();
2233
2234 let absolute = self.root.join(&path).to_owned();
2238 let scope = self.stats.egress(&absolute);
2239
2240 if let Some(broadcast) = self.get_broadcast(&path) {
2242 return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope));
2243 }
2244
2245 let mut state = self.dynamic.lock();
2246
2247 if let Some(weak) = state.served.get(&absolute) {
2251 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2252 }
2253
2254 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2257 producer.consume()
2258 } else {
2259 let producer = kio::Producer::<PendingBroadcast>::default();
2260 let consumer = producer.consume();
2261 if state.requests.insert(absolute, producer).is_err() {
2262 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2263 }
2264 consumer
2265 };
2266
2267 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2268 }
2269
2270 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2275 let prefix = prefix.as_path();
2276
2277 Some(Self::new(
2278 self.info,
2279 self.root.join(&prefix).to_owned(),
2280 self.nodes.root(&prefix)?,
2281 self.dynamic.clone(),
2282 self.stats.clone(),
2283 ))
2284 }
2285
2286 pub fn root(&self) -> &Path<'_> {
2288 &self.root
2289 }
2290
2291 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2294 self.nodes.nodes.iter().map(|(root, _)| root)
2295 }
2296
2297 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2299 self.root.join(path)
2300 }
2301}
2302
2303#[derive(Clone)]
2308pub struct AnnounceProducer {
2309 nodes: OriginNodes,
2310 root: PathOwned,
2311}
2312
2313impl AnnounceProducer {
2314 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2315 Self { nodes, root }
2316 }
2317
2318 pub fn consume(&self) -> AnnounceConsumer {
2324 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2327 }
2328
2329 pub fn root(&self) -> &Path<'_> {
2331 &self.root
2332 }
2333}
2334
2335pub struct AnnounceConsumer {
2340 id: ConsumerId,
2341 nodes: OriginNodes,
2342 root: PathOwned,
2343
2344 state: kio::Producer<OriginConsumerState>,
2347
2348 stats: stats::Session,
2351
2352 guards: HashMap<PathOwned, stats::Announce>,
2356}
2357
2358impl AnnounceConsumer {
2359 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2360 let state = kio::Producer::<OriginConsumerState>::default();
2361 let id = ConsumerId::new();
2362
2363 for (_, node) in &nodes.nodes {
2364 let notify = AnnounceConsumerNotify {
2365 root: root.clone(),
2366 state: state.clone(),
2367 };
2368 node.lock().consume(id, notify);
2369 }
2370
2371 Self {
2372 id,
2373 nodes,
2374 root,
2375 state,
2376 stats,
2377 guards: HashMap::new(),
2378 }
2379 }
2380
2381 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2387 let OriginAnnounce { path, broadcast } = update;
2388 let absolute = self.root.join(&path).to_owned();
2389 match broadcast {
2390 Some(broadcast) => {
2391 let scope = self.stats.egress(&absolute);
2392 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2393 OriginAnnounce {
2394 path,
2395 broadcast: Some(broadcast.with_stats(scope)),
2396 }
2397 }
2398 None => {
2399 self.guards.remove(&absolute);
2400 OriginAnnounce { path, broadcast: None }
2401 }
2402 }
2403 }
2404
2405 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2412 kio::wait(|waiter| self.poll_next(waiter)).await
2413 }
2414
2415 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2421 let update = {
2422 let mut state = match ready!(self.state.poll(waiter, |state| {
2423 if state.pending.is_empty() {
2424 Poll::Pending
2425 } else {
2426 Poll::Ready(())
2427 }
2428 })) {
2429 Ok(state) => state,
2430 Err(_) => return Poll::Ready(None),
2432 };
2433 state.take().expect("predicate guaranteed an update")
2434 };
2435 Poll::Ready(Some(self.attribute(update)))
2436 }
2437
2438 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2443 let update = self.state.write().ok()?.take()?;
2444 Some(self.attribute(update))
2445 }
2446
2447 pub fn is_closed(&self) -> bool {
2449 self.state.write().is_err()
2450 }
2451
2452 pub fn root(&self) -> &Path<'_> {
2454 &self.root
2455 }
2456
2457 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2459 self.root.join(path)
2460 }
2461}
2462
2463impl Drop for AnnounceConsumer {
2464 fn drop(&mut self) {
2465 for (_, root) in &self.nodes.nodes {
2466 root.lock().unconsume(self.id);
2467 }
2468 }
2469}
2470
2471#[cfg(test)]
2472use futures::FutureExt;
2473
2474#[cfg(test)]
2475#[allow(missing_docs)] impl AnnounceConsumer {
2477 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2478 let expected = expected.as_path();
2479 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2480 assert_eq!(announce.path, expected, "wrong path");
2481 let announced = announce.broadcast.expect("should be an active announce");
2482 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2483 }
2484
2485 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2489 let expected = expected.as_path();
2490 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2491 assert_eq!(announce.path, expected, "wrong path");
2492 announce.broadcast.expect("should be an active announce")
2493 }
2494
2495 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2496 let expected = expected.as_path();
2497 let announce = self.try_next().expect("no next");
2498 assert_eq!(announce.path, expected, "wrong path");
2499 let announced = announce.broadcast.expect("should be an active announce");
2500 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2501 }
2502
2503 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2505 let expected = expected.as_path();
2506 let announce = self.try_next().expect("no next");
2507 assert_eq!(announce.path, expected, "wrong path");
2508 announce.broadcast.expect("should be an active announce")
2509 }
2510
2511 pub fn assert_next_none(&mut self, expected: impl AsPath) {
2512 let expected = expected.as_path();
2513 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2514 assert_eq!(announce.path, expected, "wrong path");
2515 assert!(announce.broadcast.is_none(), "should be unannounced");
2516 }
2517
2518 pub fn assert_next_wait(&mut self) {
2519 if let Some(res) = self.next().now_or_never() {
2520 panic!("next should block: got {:?}", res.map(|a| a.path));
2521 }
2522 }
2523
2524 }
2533
2534#[cfg(test)]
2535mod tests {
2536 use crate::coding::Decode;
2537 use crate::group;
2538
2539 use super::*;
2540
2541 fn announce() -> broadcast::Route {
2543 broadcast::Route::new().with_announce(true)
2544 }
2545
2546 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
2552 let name = Path::new(name);
2553 let peer_key = fnv_key(&name, [peer]);
2554 (100u64..)
2555 .map(|id| Origin::new(id).unwrap())
2556 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
2557 .unwrap()
2558 }
2559
2560 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
2563 let source = broadcast::Info::new().produce().consume();
2564 FrontState {
2565 path: Path::new("test").to_owned(),
2566 self_origin,
2567 next_route: routes.len() as u64,
2568 routes: routes
2569 .into_iter()
2570 .enumerate()
2571 .map(|(id, route)| FrontRoute {
2572 id: id as u64,
2573 route,
2574 source: source.clone(),
2575 })
2576 .collect(),
2577 active: Some(0),
2578 linger: Duration::ZERO,
2579 closed: false,
2580 }
2581 }
2582
2583 fn sibling_route(peer: Origin) -> broadcast::Route {
2586 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
2587 announce().with_hops(hops)
2588 }
2589
2590 fn upstream_route(cost: u64) -> broadcast::Route {
2592 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
2593 announce().with_hops(hops).with_cost(cost)
2594 }
2595
2596 #[test]
2600 fn test_carrying_gate_keys() {
2601 let peer = Origin::new(3).unwrap();
2602
2603 let mut lost = front_state(
2605 origin_keyed("test", peer, false),
2606 vec![upstream_route(10), sibling_route(peer)],
2607 );
2608 lost.reselect(true);
2609 assert_eq!(
2610 lost.active,
2611 Some(0),
2612 "carrying front re-parented onto a higher-keyed peer"
2613 );
2614 lost.reselect(false);
2615 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
2616
2617 let mut won = front_state(
2619 origin_keyed("test", peer, true),
2620 vec![upstream_route(10), sibling_route(peer)],
2621 );
2622 won.reselect(true);
2623 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
2624 }
2625
2626 #[test]
2631 fn test_carrying_gate_symmetric_race() {
2632 let a = Origin::new(1).unwrap();
2633 let b = Origin::new(2).unwrap();
2634
2635 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
2636 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
2637 a_view.reselect(true);
2638 b_view.reselect(true);
2639
2640 let a_moved = a_view.active == Some(1);
2641 let b_moved = b_view.active == Some(1);
2642 assert!(
2643 a_moved != b_moved,
2644 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
2645 );
2646 }
2647
2648 #[test]
2653 fn test_carrying_switches_to_benign_routes() {
2654 let peer = Origin::new(3).unwrap();
2655 let lost = origin_keyed("test", peer, false);
2656
2657 let mut forwarder = sibling_route(peer).with_cost(4);
2659 forwarder.advertised = 4;
2660 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
2661 state.reselect(true);
2662 assert_eq!(
2663 state.active,
2664 Some(1),
2665 "a cheaper forwarder path must win while carrying"
2666 );
2667
2668 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
2670 let mut state = front_state(lost, vec![upstream_route(10), direct]);
2671 state.reselect(true);
2672 assert_eq!(
2673 state.active,
2674 Some(1),
2675 "a direct publisher route must win while carrying"
2676 );
2677 }
2678
2679 #[test]
2682 fn test_carrying_gate_ignores_unannounced_incumbent() {
2683 let peer = Origin::new(3).unwrap();
2684 let unannounced = upstream_route(10).with_announce(false);
2685 let mut state = front_state(
2686 origin_keyed("test", peer, false),
2687 vec![unannounced, sibling_route(peer)],
2688 );
2689 state.reselect(true);
2690 assert_eq!(
2691 state.active,
2692 Some(1),
2693 "an unannounced incumbent must always be displaced"
2694 );
2695 }
2696
2697 async fn settle() {
2700 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
2701 }
2702
2703 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
2706 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
2707 .await
2708 .expect("timed out waiting for a track request")
2709 .expect("source closed");
2710 assert_eq!(request.name(), name, "unexpected track dispatched");
2711 request.accept(None)
2712 }
2713
2714 #[tokio::test]
2718 async fn test_stats_tagged_end_to_end() {
2719 use crate::Timestamp;
2720 use crate::stats::{Config, Registry, Tier};
2721 use bytes::Bytes;
2722
2723 tokio::time::pause();
2724
2725 let registry = Registry::new(Config::new());
2726 let ctx = registry.tier(Tier::default()).session("acme");
2727
2728 let origin = Origin::random().produce();
2729 let ingress = origin.clone().with_stats(ctx.clone());
2730 let egress = origin.consume().with_stats(ctx.clone());
2731
2732 let mut announced = egress.announced();
2735
2736 let source = ingress.create_broadcast("demo", announce()).unwrap();
2738 let mut dynamic = source.dynamic();
2739 settle().await;
2740 settle().await;
2741
2742 let update = announced.next().await.unwrap();
2744 assert_eq!(update.path.as_str(), "demo");
2745 let broadcast = update.broadcast.unwrap();
2746
2747 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2749 let mut producer = accept_track(&mut dynamic, "video").await;
2750 settle().await;
2751 let mut sub = subscribing.await.unwrap();
2752
2753 let mut group = producer.append_group().unwrap();
2755 group
2756 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2757 .unwrap();
2758 group
2759 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
2760 .unwrap();
2761 group.finish().unwrap();
2762
2763 let mut group_c = sub.recv_group().await.unwrap().unwrap();
2765 let mut frames = 0;
2766 while let Some(frame) = group_c.read_frame().await.unwrap() {
2767 assert_eq!(frame.payload.len(), 5);
2768 frames += 1;
2769 }
2770 assert_eq!(frames, 2);
2771 settle().await;
2772
2773 let report = registry.report();
2774 let entry = report
2775 .traffic
2776 .iter()
2777 .find(|e| e.path.as_str() == "demo")
2778 .expect("demo tracked");
2779 let path_len = "demo".len() as u64;
2780
2781 let egress = &entry.publisher;
2783 assert_eq!(egress.announced, 1, "one egress announce");
2784 assert_eq!(egress.announced_bytes, path_len);
2785 assert_eq!(egress.subscriptions, 1, "one egress subscription");
2786 assert_eq!(egress.broadcasts, 1, "one viewer");
2787 assert_eq!(egress.groups, 1);
2788 assert_eq!(egress.frames, 2);
2789 assert_eq!(egress.bytes, 10);
2790 assert_eq!(egress.fetches, 0);
2791
2792 let ingress = &entry.subscriber;
2794 assert_eq!(ingress.announced, 1, "one ingress announce");
2795 assert_eq!(ingress.announced_bytes, path_len);
2796 assert_eq!(ingress.subscriptions, 1, "one ingress track");
2797 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
2798 assert_eq!(ingress.groups, 1);
2799 assert_eq!(ingress.frames, 2);
2800 assert_eq!(ingress.bytes, 10);
2801
2802 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
2804 let _ = fetched;
2805 settle().await;
2806 let report = registry.report();
2807 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
2808 assert_eq!(entry.publisher.fetches, 1, "one fetch");
2809 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
2810 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
2811 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
2815 }
2816
2817 #[tokio::test]
2822 async fn test_stats_read_frame_counts_once() {
2823 use crate::Timestamp;
2824 use crate::stats::{Config, Registry, Tier};
2825 use bytes::Bytes;
2826
2827 tokio::time::pause();
2828
2829 let registry = Registry::new(Config::new());
2830 let ctx = registry.tier(Tier::default()).session("acme");
2831
2832 let origin = Origin::random().produce();
2833 let ingress = origin.clone().with_stats(ctx.clone());
2834 let egress = origin.consume().with_stats(ctx.clone());
2835
2836 let mut announced = egress.announced();
2837 let source = ingress.create_broadcast("demo", announce()).unwrap();
2838 let mut dynamic = source.dynamic();
2839 settle().await;
2840 settle().await;
2841
2842 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2843 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2844 let mut producer = accept_track(&mut dynamic, "video").await;
2845 settle().await;
2846 let mut sub = subscribing.await.unwrap();
2847
2848 producer
2850 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2851 .unwrap();
2852
2853 let frame = sub.read_frame().await.unwrap().expect("frame");
2854 assert_eq!(frame.payload.len(), 5);
2855 settle().await;
2856
2857 let report = registry.report();
2858 let entry = report
2859 .traffic
2860 .iter()
2861 .find(|e| e.path.as_str() == "demo")
2862 .expect("demo tracked");
2863 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
2864 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
2865 assert_eq!(
2866 entry.publisher.bytes, 5,
2867 "payload counted once, not zero and not doubled"
2868 );
2869 }
2870
2871 #[tokio::test]
2875 async fn test_stats_datagrams_counted_both_sides() {
2876 use crate::Timestamp;
2877 use crate::stats::{Config, Registry, Tier};
2878
2879 tokio::time::pause();
2880
2881 let registry = Registry::new(Config::new());
2882 let ctx = registry.tier(Tier::default()).session("acme");
2883
2884 let origin = Origin::random().produce();
2885 let ingress = origin.clone().with_stats(ctx.clone());
2886 let egress = origin.consume().with_stats(ctx.clone());
2887
2888 let mut announced = egress.announced();
2889 let source = ingress.create_broadcast("demo", announce()).unwrap();
2890 let mut dynamic = source.dynamic();
2891 settle().await;
2892 settle().await;
2893
2894 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2895 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2896 let mut producer = accept_track(&mut dynamic, "video").await;
2897 settle().await;
2898 let mut sub = subscribing.await.unwrap();
2899
2900 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
2901 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
2902 assert_eq!(&datagram.payload[..], b"hello");
2903 settle().await;
2904
2905 let report = registry.report();
2906 let entry = report
2907 .traffic
2908 .iter()
2909 .find(|e| e.path.as_str() == "demo")
2910 .expect("demo tracked");
2911
2912 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
2913 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
2914 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
2915 assert_eq!(traffic.frames, 1, "{side}: one frame");
2916 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
2917 }
2918 }
2919
2920 #[test]
2921 fn origin_rejects_reserved_ids() {
2922 assert!(Origin::new(0).is_err());
2923 assert!(Origin::new(1u64 << 62).is_err());
2924 assert_eq!(Origin::new(1).unwrap().id(), 1);
2925
2926 let mut zero = [0u8].as_slice();
2927 assert_eq!(
2928 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
2929 Origin::UNKNOWN
2930 );
2931 }
2932
2933 #[test]
2934 fn origin_list_push_fails_at_limit() {
2935 let mut list = OriginList::new();
2936 for _ in 0..MAX_HOPS {
2937 list.push(Origin::random()).unwrap();
2938 }
2939 assert_eq!(list.len(), MAX_HOPS);
2940 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
2941 }
2942
2943 #[test]
2944 fn origin_list_replace_first() {
2945 let mut list = OriginList::new();
2946 for _ in 0..3 {
2947 list.push(Origin::UNKNOWN).unwrap();
2948 }
2949
2950 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
2952 assert_eq!(
2953 list.as_slice(),
2954 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
2955 );
2956
2957 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
2959 assert_eq!(list.len(), 3);
2960 }
2961
2962 #[test]
2963 fn origin_list_try_from_vec_enforces_limit() {
2964 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
2965 assert!(OriginList::try_from(under).is_ok());
2966
2967 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
2968 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
2969 }
2970
2971 #[tokio::test]
2972 async fn test_announce() {
2973 tokio::time::pause();
2974
2975 let origin = Origin::random().produce();
2976
2977 let mut consumer1 = origin.consume().announced();
2978 consumer1.assert_next_wait();
2979
2980 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
2982 settle().await;
2983
2984 consumer1.assert_next_some("test1");
2985 consumer1.assert_next_wait();
2986
2987 let mut consumer2 = origin.consume().announced();
2990
2991 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
2993 settle().await;
2994
2995 consumer1.assert_next_some("test2");
2996 consumer1.assert_next_wait();
2997
2998 consumer2.assert_next_some("test1");
2999 consumer2.assert_next_some("test2");
3000 consumer2.assert_next_wait();
3001
3002 broadcast1.finish();
3004 settle().await;
3005
3006 consumer1.assert_next_none("test1");
3008 consumer2.assert_next_none("test1");
3009 consumer1.assert_next_wait();
3010 consumer2.assert_next_wait();
3011
3012 let mut consumer3 = origin.consume().announced();
3014 consumer3.assert_next_some("test2");
3015 consumer3.assert_next_wait();
3016
3017 broadcast2.finish();
3018 settle().await;
3019
3020 consumer1.assert_next_none("test2");
3021 consumer2.assert_next_none("test2");
3022 consumer3.assert_next_none("test2");
3023 }
3024
3025 #[tokio::test]
3029 async fn test_duplicate() {
3030 tokio::time::pause();
3031
3032 let origin = Origin::random().produce();
3033 let consumer = origin.consume();
3034 let mut announced = consumer.announced();
3035
3036 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3037 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3038 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3039 settle().await;
3040 assert!(consumer.get_broadcast("test").is_some());
3041
3042 announced.assert_next_some("test");
3043 announced.assert_next_wait();
3044
3045 broadcast2.finish();
3047 settle().await;
3048 assert!(consumer.get_broadcast("test").is_some());
3049 announced.assert_next_wait();
3050
3051 broadcast1.finish();
3053 settle().await;
3054 assert!(consumer.get_broadcast("test").is_some());
3055 announced.assert_next_wait();
3056
3057 broadcast3.finish();
3059 settle().await;
3060 assert!(consumer.get_broadcast("test").is_none());
3061
3062 announced.assert_next_none("test");
3063 announced.assert_next_wait();
3064 }
3065
3066 #[tokio::test]
3069 async fn test_route_failover() {
3070 tokio::time::pause();
3071
3072 let origin = Origin::random().produce();
3073 let consumer = origin.consume();
3074 let mut announced = consumer.announced();
3075
3076 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3077 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3078
3079 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3081 let mut dynamic_a = source_a.dynamic();
3082 settle().await;
3083 settle().await;
3084 let broadcast = consumer.request_broadcast("test").await.unwrap();
3085 announced.assert_next_some("test");
3086
3087 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3089 let mut dynamic_b = source_b.dynamic();
3090 settle().await;
3091 settle().await;
3092 announced.assert_next_wait();
3093
3094 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3096 let mut producer = accept_track(&mut dynamic_a, "video").await;
3097 settle().await;
3098 dynamic_b.assert_no_request();
3099
3100 let mut sub = subscribing.await.unwrap();
3101 sub.assert_no_group();
3104 assert_eq!(producer.subscription().unwrap().group_start, None);
3105
3106 producer.append_group().unwrap();
3107 producer.append_group().unwrap();
3108 assert_eq!(sub.assert_group().sequence, 0);
3109 assert_eq!(sub.assert_group().sequence, 1);
3110
3111 producer.abort(Error::Dropped).unwrap();
3115 source_a.abort(Error::Dropped).unwrap();
3116 drop(dynamic_a);
3117 settle().await;
3118 announced.assert_next_wait();
3119
3120 let mut producer = accept_track(&mut dynamic_b, "video").await;
3123 settle().await;
3124 sub.assert_no_group();
3125 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3126 producer.create_group(group::Info { sequence: 1 }).unwrap();
3127 producer.create_group(group::Info { sequence: 2 }).unwrap();
3128 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3129 sub.assert_not_closed();
3130 }
3131
3132 #[tokio::test]
3135 async fn test_broadcast_route_watch() {
3136 let mut producer = broadcast::Info::new().produce();
3137 let mut consumer = producer.consume();
3138
3139 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3141
3142 producer.set_route(broadcast::Route::default()).unwrap();
3144 assert!(consumer.route_changed().now_or_never().is_none());
3145
3146 let mut hops = OriginList::new();
3147 hops.push(Origin::new(7).unwrap()).unwrap();
3148 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3149 producer.set_route(route.clone()).unwrap();
3150 assert_eq!(consumer.route_changed().await.unwrap(), route);
3151
3152 let mut fresh = producer.consume();
3154 assert_eq!(fresh.route_changed().await.unwrap(), route);
3155
3156 drop(producer);
3157 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3158 }
3159
3160 #[tokio::test]
3164 async fn test_route_cost_update() {
3165 tokio::time::pause();
3166
3167 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3171 let consumer = origin.consume();
3172 let mut announced = consumer.announced();
3173
3174 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3175 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3176
3177 let mut source_a = origin
3179 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3180 .unwrap();
3181 let mut dynamic_a = source_a.dynamic();
3182 settle().await;
3183 let broadcast = consumer.request_broadcast("test").await.unwrap();
3184 announced.assert_next_some("test");
3185
3186 let mut watch = broadcast.clone();
3187 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3188
3189 let mut source_b = origin
3190 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3191 .unwrap();
3192 let mut dynamic_b = source_b.dynamic();
3193 settle().await;
3194 assert!(
3195 watch.route_changed().now_or_never().is_none(),
3196 "a losing standby must not change the advertised route"
3197 );
3198
3199 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3201 let mut producer = accept_track(&mut dynamic_a, "video").await;
3202 settle().await;
3203 let mut sub = subscribing.await.unwrap();
3204 producer.append_group().unwrap();
3205 assert_eq!(sub.assert_group().sequence, 0);
3206
3207 source_a
3210 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3211 .unwrap();
3212 settle().await;
3213 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3214 announced.assert_next_wait();
3215
3216 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3217 settle().await;
3218 sub.assert_no_group();
3221 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3222 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3223 assert_eq!(sub.assert_group().sequence, 1);
3224 sub.assert_not_closed();
3225
3226 source_b
3228 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3229 .unwrap();
3230 settle().await;
3231 let advertised = watch.route_changed().await.unwrap();
3232 assert_eq!(advertised.hops, hops_b);
3233 assert_eq!(advertised.cost, 5);
3234 announced.assert_next_wait();
3235 }
3236
3237 #[tokio::test]
3240 async fn test_completed_track_survives_route_churn() {
3241 tokio::time::pause();
3242
3243 let origin = Origin::random().produce();
3244 let consumer = origin.consume();
3245
3246 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3247 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3248
3249 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3250 let mut dynamic_a = source_a.dynamic();
3251 settle().await;
3252 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3253 let mut dynamic_b = source_b.dynamic();
3254 settle().await;
3255 settle().await;
3256 let broadcast = consumer.request_broadcast("test").await.unwrap();
3257
3258 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3260 let mut producer = accept_track(&mut dynamic_a, "video").await;
3261 settle().await;
3262 let mut sub = subscribing.await.unwrap();
3263 producer.append_group().unwrap();
3264 assert_eq!(sub.assert_group().sequence, 0);
3265 producer.finish().unwrap();
3266 drop(producer);
3267 settle().await;
3268 sub.assert_closed();
3269
3270 source_a.abort(Error::Dropped).unwrap();
3272 drop(dynamic_a);
3273 settle().await;
3274 dynamic_b.assert_no_request();
3275
3276 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3278 late.assert_closed();
3279 }
3280
3281 #[tokio::test]
3284 async fn test_serve_resets_retry_budget() {
3285 tokio::time::pause();
3286
3287 let origin = Origin::random().produce();
3288 let consumer = origin.consume();
3289
3290 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3291 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3292 let mut dynamic = source.dynamic();
3293 settle().await;
3294 settle().await;
3295 let broadcast = consumer.request_broadcast("test").await.unwrap();
3296
3297 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3299
3300 for _ in 0..2 * MAX_TRACK_RETRIES {
3303 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3304 .await
3305 .expect("timed out waiting for a retry")
3306 .unwrap();
3307 request.reject(Error::NotFound);
3308 let producer = accept_track(&mut dynamic, "video").await;
3309 settle().await;
3310 drop(producer);
3311 }
3312
3313 let _producer = accept_track(&mut dynamic, "video").await;
3314 settle().await;
3315 let mut sub = subscribing.await.unwrap();
3316 sub.assert_not_closed();
3317 }
3318
3319 #[tokio::test]
3323 async fn test_route_handover() {
3324 tokio::time::pause();
3325
3326 let origin = Origin::random().produce();
3327 let consumer = origin.consume();
3328 let mut announced = consumer.announced();
3329
3330 let hops_long = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3331 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3332
3333 let source_a = origin
3334 .create_broadcast("test", announce().with_hops(hops_long))
3335 .unwrap();
3336 let mut dynamic_a = source_a.dynamic();
3337 settle().await;
3338 settle().await;
3339 let broadcast = consumer.request_broadcast("test").await.unwrap();
3340 announced.assert_next_some("test");
3341
3342 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3343 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3344 settle().await;
3345 let mut sub = subscribing.await.unwrap();
3346 producer_a.append_group().unwrap();
3347 producer_a.append_group().unwrap();
3348 assert_eq!(sub.assert_group().sequence, 0);
3349 assert_eq!(sub.assert_group().sequence, 1);
3350
3351 let source_b = origin
3354 .create_broadcast("test", announce().with_hops(hops_short))
3355 .unwrap();
3356 let mut dynamic_b = source_b.dynamic();
3357 settle().await;
3358 settle().await;
3359 announced.assert_next_wait();
3360
3361 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3362 settle().await;
3363
3364 sub.assert_no_group();
3367 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3368 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3369
3370 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3372 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3373 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3374 assert_eq!(sub.assert_group().sequence, 2);
3375 assert_eq!(sub.assert_group().sequence, 3);
3376 sub.assert_no_group();
3377 sub.assert_not_closed();
3378 }
3379
3380 #[tokio::test(start_paused = true)]
3383 async fn test_route_unannounce_immediate() {
3384 let origin = Origin::random().produce();
3385 let consumer = origin.consume();
3386 let mut announced = consumer.announced();
3387
3388 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3389 let mut source = origin
3390 .create_broadcast("test", announce().with_hops(hops.clone()))
3391 .unwrap();
3392 settle().await;
3393 let broadcast = consumer.request_broadcast("test").await.unwrap();
3394 announced.assert_next_some("test");
3395
3396 source.finish();
3399 settle().await;
3400 announced.assert_next_none("test");
3401
3402 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3404 settle().await;
3405 let fresh = consumer.request_broadcast("test").await.unwrap();
3406 announced.assert_next_some("test");
3407 assert!(
3408 !fresh.is_clone(&broadcast),
3409 "re-create must not splice the old broadcast"
3410 );
3411 }
3412
3413 #[tokio::test(start_paused = true)]
3418 async fn test_route_detach_immediate() {
3419 let origin = Origin::random().produce();
3420 let consumer = origin.consume();
3421 let mut announced = consumer.announced();
3422
3423 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3424 let source = origin
3425 .create_broadcast("test", announce().with_hops(hops.clone()))
3426 .unwrap();
3427 let mut dynamic = source.dynamic();
3428 settle().await;
3429 settle().await;
3430 let broadcast = consumer.request_broadcast("test").await.unwrap();
3431 announced.assert_next_some("test");
3432
3433 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3434 let producer = accept_track(&mut dynamic, "video").await;
3435 settle().await;
3436 let mut sub = subscribing.await.unwrap();
3437
3438 drop(producer);
3440 source.abort(Error::Dropped).unwrap();
3441 drop(dynamic);
3442
3443 settle().await;
3444 announced.assert_next_none("test");
3445 sub.assert_error();
3446
3447 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3450 settle().await;
3451 settle().await;
3452 let fresh = consumer.request_broadcast("test").await.unwrap();
3453 announced.assert_next_some("test");
3454 assert!(
3455 !fresh.is_clone(&broadcast),
3456 "re-create must not splice the old broadcast"
3457 );
3458 }
3459
3460 #[tokio::test(start_paused = true)]
3465 async fn test_idle_track_releases_without_respinning() {
3466 let origin = Info::new(Origin::random()).produce();
3467 let consumer = origin.consume();
3468
3469 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3470 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3471 let mut dynamic = source.dynamic();
3472 settle().await;
3473 let broadcast = consumer.request_broadcast("test").await.unwrap();
3474
3475 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3476 let producer = accept_track(&mut dynamic, "video").await;
3477 settle().await;
3478 let sub = subscribing.await.unwrap();
3479
3480 drop(sub);
3483 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
3484 settle().await;
3485 assert!(
3486 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
3487 "the copy must stay spliced inside the linger",
3488 );
3489
3490 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3493 settle().await;
3494 assert!(
3495 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3496 "an idle copy must be released after the linger",
3497 );
3498
3499 for _ in 0..3 {
3503 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3504 settle().await;
3505 assert!(
3506 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3507 "an unread copy must stay released, not be re-spliced",
3508 );
3509 }
3510 assert!(
3511 dynamic.requested_track().now_or_never().is_none(),
3512 "an unread track must not be re-requested",
3513 );
3514 drop(producer);
3515
3516 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3518 let mut producer = accept_track(&mut dynamic, "video").await;
3519 settle().await;
3520 let mut sub = subscribing.await.unwrap();
3521 producer.append_group().unwrap();
3522 assert_eq!(sub.assert_group().sequence, 0);
3523 }
3524
3525 #[tokio::test(start_paused = true)]
3529 async fn test_back_to_back_fetches_reuse_the_track() {
3530 let origin = Info::new(Origin::random()).produce();
3531 let consumer = origin.consume();
3532
3533 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3534 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3535 let mut dynamic = source.dynamic();
3536 settle().await;
3537 let broadcast = consumer.request_broadcast("test").await.unwrap();
3538
3539 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3541 let mut producer = accept_track(&mut dynamic, "video").await;
3542 producer.append_group().unwrap().finish().unwrap();
3543 settle().await;
3544 let first = fetching.await.expect("first fetch");
3545 drop(first);
3546
3547 settle().await;
3549 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3550 settle().await;
3551 assert!(
3552 dynamic.requested_track().now_or_never().is_none(),
3553 "a fetch inside the linger must reuse the track, not re-request it",
3554 );
3555 drop(fetching.await.expect("second fetch"));
3556
3557 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
3559 settle().await;
3560 assert!(
3561 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3562 "the copy must be released once the fetches stop",
3563 );
3564 drop(producer);
3565
3566 settle().await;
3568 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3569 let mut producer = accept_track(&mut dynamic, "video").await;
3570 producer.append_group().unwrap().finish().unwrap();
3571 settle().await;
3572 fetching.await.expect("fetch after the linger");
3573 }
3574
3575 #[tokio::test(start_paused = true)]
3579 async fn test_linger_reconnect_splices() {
3580 let origin = Info::new(Origin::random())
3581 .with_linger(Duration::from_secs(5))
3582 .produce();
3583 let consumer = origin.consume();
3584 let mut announced = consumer.announced();
3585
3586 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3587 let source = origin
3588 .create_broadcast("test", announce().with_hops(hops.clone()))
3589 .unwrap();
3590 let mut dynamic = source.dynamic();
3591 settle().await;
3592 settle().await;
3593 let broadcast = consumer.request_broadcast("test").await.unwrap();
3594 announced.assert_next_some("test");
3595
3596 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3597 let mut producer = accept_track(&mut dynamic, "video").await;
3598 settle().await;
3599 let mut sub = subscribing.await.unwrap();
3600
3601 producer.append_group().unwrap();
3602 producer.append_group().unwrap();
3603 assert_eq!(sub.assert_group().sequence, 0);
3604 assert_eq!(sub.assert_group().sequence, 1);
3605
3606 drop(producer);
3609 source.abort(Error::Dropped).unwrap();
3610 drop(dynamic);
3611 settle().await;
3612
3613 announced.assert_next_wait();
3615 sub.assert_no_group();
3616 sub.assert_not_closed();
3617
3618 let during = consumer.request_broadcast("test").await.unwrap();
3620 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
3621
3622 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3625 let mut dynamic = source.dynamic();
3626 settle().await;
3627 settle().await;
3628 announced.assert_next_wait();
3629 let again = consumer.request_broadcast("test").await.unwrap();
3630 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3631
3632 let mut producer = accept_track(&mut dynamic, "video").await;
3636 settle().await;
3637 sub.assert_no_group();
3638 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3639 producer.create_group(group::Info { sequence: 2 }).unwrap();
3640 assert_eq!(sub.assert_group().sequence, 2);
3641 sub.assert_not_closed();
3642 }
3643
3644 #[tokio::test(start_paused = true)]
3647 async fn test_linger_expiry_closes() {
3648 let origin = Info::new(Origin::random())
3649 .with_linger(Duration::from_secs(5))
3650 .produce();
3651 let consumer = origin.consume();
3652 let mut announced = consumer.announced();
3653
3654 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3655 let source = origin
3656 .create_broadcast("test", announce().with_hops(hops.clone()))
3657 .unwrap();
3658 let mut dynamic = source.dynamic();
3659 settle().await;
3660 settle().await;
3661 let broadcast = consumer.request_broadcast("test").await.unwrap();
3662 announced.assert_next_some("test");
3663
3664 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3665 let producer = accept_track(&mut dynamic, "video").await;
3666 settle().await;
3667 let mut sub = subscribing.await.unwrap();
3668
3669 drop(producer);
3670 source.abort(Error::Dropped).unwrap();
3671 drop(dynamic);
3672 settle().await;
3673 announced.assert_next_wait();
3674
3675 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
3677 settle().await;
3678 announced.assert_next_none("test");
3679 sub.assert_error();
3680
3681 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3683 settle().await;
3684 settle().await;
3685 let fresh = consumer.request_broadcast("test").await.unwrap();
3686 announced.assert_next_some("test");
3687 assert!(
3688 !fresh.is_clone(&broadcast),
3689 "a late re-create must not splice the expired broadcast"
3690 );
3691 }
3692
3693 #[tokio::test(start_paused = true)]
3697 async fn test_linger_forever() {
3698 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
3699 let consumer = origin.consume();
3700 let mut announced = consumer.announced();
3701
3702 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3703 let source = origin
3704 .create_broadcast("test", announce().with_hops(hops.clone()))
3705 .unwrap();
3706 settle().await;
3707 let broadcast = consumer.request_broadcast("test").await.unwrap();
3708 announced.assert_next_some("test");
3709
3710 source.abort(Error::Dropped).unwrap();
3711 settle().await;
3712
3713 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
3715 announced.assert_next_wait();
3716 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3717 settle().await;
3718 settle().await;
3719 let again = consumer.request_broadcast("test").await.unwrap();
3720 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3721 drop(source);
3722 }
3723
3724 #[tokio::test(start_paused = true)]
3727 async fn test_linger_skipped_on_finish() {
3728 let origin = Info::new(Origin::random())
3729 .with_linger(Duration::from_secs(5))
3730 .produce();
3731 let consumer = origin.consume();
3732 let mut announced = consumer.announced();
3733
3734 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3735 let mut source = origin
3736 .create_broadcast("test", announce().with_hops(hops.clone()))
3737 .unwrap();
3738 settle().await;
3739 let broadcast = consumer.request_broadcast("test").await.unwrap();
3740 announced.assert_next_some("test");
3741
3742 source.finish();
3745 settle().await;
3746 announced.assert_next_none("test");
3747
3748 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3750 settle().await;
3751 let fresh = consumer.request_broadcast("test").await.unwrap();
3752 announced.assert_next_some("test");
3753 assert!(
3754 !fresh.is_clone(&broadcast),
3755 "a finish must not leave a lingering broadcast to splice into"
3756 );
3757 }
3758
3759 #[tokio::test]
3762 async fn test_announce_toggle() {
3763 tokio::time::pause();
3764
3765 let origin = Origin::random().produce();
3766 let consumer = origin.consume();
3767 let mut announced = consumer.announced();
3768
3769 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3770 settle().await;
3771
3772 announced.assert_next_wait();
3774 let broadcast = consumer
3775 .get_broadcast("test")
3776 .expect("offline broadcast is still routable");
3777 assert!(!broadcast.route().announce);
3778
3779 let requested = consumer.request_broadcast("test").await.unwrap();
3781 assert!(requested.is_clone(&broadcast));
3782
3783 source.set_route(announce()).unwrap();
3785 settle().await;
3786 let face = announced.assert_next_some("test");
3787 assert!(face.is_clone(&broadcast));
3788
3789 let mut fresh = origin.consume().announced();
3791 fresh.assert_next_some("test");
3792 fresh.assert_next_wait();
3793
3794 source.set_route(broadcast::Route::new()).unwrap();
3796 settle().await;
3797 announced.assert_next_none("test");
3798 assert!(consumer.get_broadcast("test").is_some());
3799 let mut fresh = origin.consume().announced();
3800 fresh.assert_next_wait();
3801
3802 source.finish();
3803 settle().await;
3804 assert!(consumer.get_broadcast("test").is_none());
3805 }
3806
3807 #[tokio::test]
3810 async fn test_announce_beats_offline() {
3811 tokio::time::pause();
3812
3813 let origin = Origin::random().produce();
3814 let consumer = origin.consume();
3815 let mut announced = consumer.announced();
3816
3817 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3819 settle().await;
3820 announced.assert_next_wait();
3821
3822 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
3825 settle().await;
3826 announced.assert_next_some("test");
3827 let face = consumer.get_broadcast("test").unwrap();
3828 assert!(face.route().announce);
3829 assert_eq!(face.route().cost, 10);
3830
3831 announced_source.finish();
3834 settle().await;
3835 announced.assert_next_none("test");
3836 assert!(consumer.get_broadcast("test").is_some());
3837 }
3838
3839 #[tokio::test]
3842 async fn test_better_source_no_churn() {
3843 tokio::time::pause();
3844
3845 let origin = Origin::random().produce();
3846 let mut announced = origin.consume().announced();
3847
3848 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3850 let _a = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3851 settle().await;
3852 let face = announced.assert_next_some("test");
3853
3854 let _b = origin.create_broadcast("test", announce()).unwrap();
3855 settle().await;
3856 announced.assert_next_wait();
3857 let current = origin.consume().get_broadcast("test").unwrap();
3858 assert!(current.is_clone(&face), "the broadcast identity must not change");
3859 assert!(current.route().hops.is_empty());
3861 }
3862
3863 #[tokio::test]
3864 async fn test_duplicate_reverse() {
3865 tokio::time::pause();
3866
3867 let origin = Origin::random().produce();
3868
3869 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3870 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3871 settle().await;
3872 assert!(origin.consume().get_broadcast("test").is_some());
3873
3874 broadcast2.finish();
3876 settle().await;
3877 assert!(origin.consume().get_broadcast("test").is_some());
3878
3879 broadcast1.finish();
3880 settle().await;
3881 assert!(origin.consume().get_broadcast("test").is_none());
3882 }
3883
3884 #[tokio::test]
3885 async fn test_deterministic_tiebreak() {
3886 tokio::time::pause();
3887
3888 fn hops(ids: &[u64]) -> OriginList {
3889 OriginList::try_from(
3890 ids.iter()
3891 .copied()
3892 .map(|id| Origin::new(id).unwrap())
3893 .collect::<Vec<_>>(),
3894 )
3895 .unwrap()
3896 }
3897
3898 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
3901 let origin = Origin::random().produce();
3902 let _a = origin
3903 .create_broadcast("test", announce().with_hops(hops(first)))
3904 .unwrap();
3905 let _b = origin
3906 .create_broadcast("test", announce().with_hops(hops(second)))
3907 .unwrap();
3908 settle().await;
3909 origin.consume().get_broadcast("test").unwrap().route().hops
3910 }
3911
3912 let forward = winner(&[10, 20], &[30, 40]).await;
3915 let reverse = winner(&[30, 40], &[10, 20]).await;
3916 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
3917
3918 assert_eq!(winner(&[10, 20], &[30]).await.len(), 1);
3920 assert_eq!(winner(&[30], &[10, 20]).await.len(), 1);
3921 }
3922
3923 #[tokio::test]
3928 async fn test_many_announces() {
3929 let origin = Origin::random().produce();
3930
3931 let mut consumer = origin.consume().announced();
3932 let mut broadcasts = Vec::new();
3934 for i in 0..256 {
3935 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3936 settle().await;
3937 }
3938
3939 for i in 0..256 {
3940 consumer.assert_next_some(format!("test{i:03}"));
3941 }
3942 consumer.assert_next_wait();
3943 }
3944
3945 #[tokio::test]
3946 async fn test_many_announces_try() {
3947 let origin = Origin::random().produce();
3948
3949 let mut consumer = origin.consume().announced();
3950 let mut broadcasts = Vec::new();
3952 for i in 0..256 {
3953 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3954 settle().await;
3955 }
3956
3957 for i in 0..256 {
3958 consumer.assert_try_next_some(format!("test{i:03}"));
3959 }
3960 }
3961
3962 #[tokio::test]
3963 async fn test_with_root_basic() {
3964 let origin = Origin::random().produce();
3965
3966 let foo_producer = origin.with_root("foo").expect("should create root");
3968 assert_eq!(foo_producer.root().as_str(), "foo");
3969
3970 let mut consumer = origin.consume().announced();
3971
3972 let _broadcast = foo_producer
3974 .create_broadcast("bar/baz", announce())
3975 .expect("publish allowed");
3976 settle().await;
3977 consumer.assert_next_some("foo/bar/baz");
3979
3980 let mut foo_consumer = foo_producer.consume().announced();
3982 foo_consumer.assert_next_some("bar/baz");
3983 }
3984
3985 #[tokio::test]
3986 async fn test_with_root_nested() {
3987 let origin = Origin::random().produce();
3988
3989 let foo_producer = origin.with_root("foo").expect("should create foo root");
3991 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
3992 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
3993
3994 let mut consumer = origin.consume().announced();
3995
3996 let _broadcast = foo_bar_producer
3998 .create_broadcast("baz", announce())
3999 .expect("publish allowed");
4000 settle().await;
4001 consumer.assert_next_some("foo/bar/baz");
4003
4004 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
4006 foo_bar_consumer.assert_next_some("baz");
4007 }
4008
4009 #[tokio::test]
4010 async fn test_publish_scope_allows() {
4011 let origin = Origin::random().produce();
4012
4013 let limited_producer = origin
4015 .scope(&["allowed/path1".into(), "allowed/path2".into()])
4016 .expect("should create limited producer");
4017
4018 let _broadcast = limited_producer
4020 .create_broadcast("allowed/path1", announce())
4021 .expect("publish allowed");
4022 let _keep2 = limited_producer
4023 .create_broadcast("allowed/path1/nested", announce())
4024 .expect("publish allowed");
4025 let _keep3 = limited_producer
4026 .create_broadcast("allowed/path2", announce())
4027 .expect("publish allowed");
4028 settle().await;
4029
4030 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
4032 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
4034 }
4035
4036 #[tokio::test]
4037 async fn test_publish_max_parts() {
4038 let origin = Origin::random().produce();
4039
4040 let at_limit = (0..Path::MAX_PARTS)
4041 .map(|i| i.to_string())
4042 .collect::<Vec<_>>()
4043 .join("/");
4044 let _broadcast = origin
4045 .create_broadcast(at_limit.as_str(), announce())
4046 .expect("publish allowed");
4047 settle().await;
4048
4049 let too_deep = format!("{at_limit}/extra");
4050 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
4051
4052 let rooted = origin.with_root("root").expect("wildcard allows any root");
4054 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
4055 }
4056
4057 #[tokio::test]
4058 async fn test_publish_scope_empty() {
4059 let origin = Origin::random().produce();
4060
4061 assert!(origin.scope(&[]).is_none());
4063 }
4064
4065 #[tokio::test]
4066 async fn test_consume_scope_filters() {
4067 let origin = Origin::random().produce();
4068
4069 let mut consumer = origin.consume().announced();
4070
4071 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
4073 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
4074 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
4075 settle().await;
4076
4077 let mut limited_consumer = origin
4079 .consume()
4080 .scope(&["allowed".into()])
4081 .expect("should create limited consumer")
4082 .announced();
4083
4084 limited_consumer.assert_next_some("allowed");
4086 limited_consumer.assert_next_some("allowed/nested");
4087 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
4091 consumer.assert_next_some("allowed/nested");
4092 consumer.assert_next_some("notallowed");
4093 }
4094
4095 #[tokio::test]
4096 async fn test_consume_scope_multiple_prefixes() {
4097 let origin = Origin::random().produce();
4098
4099 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
4100 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
4101 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
4102 settle().await;
4103
4104 let mut limited_consumer = origin
4106 .consume()
4107 .scope(&["foo".into(), "bar".into()])
4108 .expect("should create limited consumer")
4109 .announced();
4110
4111 limited_consumer.assert_next_some("bar/test");
4113 limited_consumer.assert_next_some("foo/test");
4114 limited_consumer.assert_next_wait(); }
4116
4117 #[tokio::test]
4118 async fn test_with_root_and_publish_scope() {
4119 let origin = Origin::random().produce();
4120
4121 let foo_producer = origin.with_root("foo").expect("should create foo root");
4123
4124 let limited_producer = foo_producer
4126 .scope(&["bar".into(), "goop/pee".into()])
4127 .expect("should create limited producer");
4128
4129 let mut consumer = origin.consume().announced();
4130
4131 let _broadcast = limited_producer
4133 .create_broadcast("bar", announce())
4134 .expect("publish allowed");
4135 let _keep2 = limited_producer
4136 .create_broadcast("bar/nested", announce())
4137 .expect("publish allowed");
4138 let _keep3 = limited_producer
4139 .create_broadcast("goop/pee", announce())
4140 .expect("publish allowed");
4141 let _keep4 = limited_producer
4142 .create_broadcast("goop/pee/nested", announce())
4143 .expect("publish allowed");
4144 settle().await;
4145
4146 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
4148 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
4150
4151 consumer.assert_next_some("foo/bar");
4153 consumer.assert_next_some("foo/bar/nested");
4154 consumer.assert_next_some("foo/goop/pee");
4155 consumer.assert_next_some("foo/goop/pee/nested");
4156 }
4157
4158 #[tokio::test]
4159 async fn test_with_root_and_consume_scope() {
4160 let origin = Origin::random().produce();
4161
4162 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
4164 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
4165 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
4166 settle().await;
4167
4168 let foo_producer = origin.with_root("foo").expect("should create foo root");
4170
4171 let mut limited_consumer = foo_producer
4173 .consume()
4174 .scope(&["bar".into(), "goop/pee".into()])
4175 .expect("should create limited consumer")
4176 .announced();
4177
4178 limited_consumer.assert_next_some("bar/test");
4180 limited_consumer.assert_next_some("goop/pee/test");
4181 limited_consumer.assert_next_wait(); }
4183
4184 #[tokio::test]
4185 async fn test_with_root_unauthorized() {
4186 let origin = Origin::random().produce();
4187
4188 let limited_producer = origin
4190 .scope(&["allowed".into()])
4191 .expect("should create limited producer");
4192
4193 assert!(limited_producer.with_root("notallowed").is_none());
4195
4196 let allowed_root = limited_producer
4198 .with_root("allowed")
4199 .expect("should create allowed root");
4200 assert_eq!(allowed_root.root().as_str(), "allowed");
4201 }
4202
4203 #[tokio::test]
4204 async fn test_wildcard_permission() {
4205 let origin = Origin::random().produce();
4206
4207 let root_producer = origin.clone();
4209
4210 let _broadcast = root_producer
4212 .create_broadcast("any/path", announce())
4213 .expect("publish allowed");
4214 let _keep2 = root_producer
4215 .create_broadcast("other/path", announce())
4216 .expect("publish allowed");
4217 settle().await;
4218
4219 let foo_producer = root_producer.with_root("foo").expect("should create any root");
4221 assert_eq!(foo_producer.root().as_str(), "foo");
4222 }
4223
4224 #[tokio::test]
4225 async fn test_consume_broadcast_with_permissions() {
4226 let origin = Origin::random().produce();
4227
4228 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
4229 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
4230 settle().await;
4231
4232 let limited_consumer = origin
4234 .consume()
4235 .scope(&["allowed".into()])
4236 .expect("should create limited consumer");
4237
4238 let result = limited_consumer.get_broadcast("allowed/test");
4240 assert!(result.is_some());
4241 assert!(
4242 result
4243 .unwrap()
4244 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
4245 );
4246
4247 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
4249
4250 let consumer = origin.consume();
4252 assert!(consumer.get_broadcast("allowed/test").is_some());
4253 assert!(consumer.get_broadcast("notallowed/test").is_some());
4254 }
4255
4256 #[tokio::test]
4257 async fn test_nested_paths_with_permissions() {
4258 let origin = Origin::random().produce();
4259
4260 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
4262
4263 let _broadcast = limited_producer
4265 .create_broadcast("a/b/c", announce())
4266 .expect("publish allowed");
4267 let _keep2 = limited_producer
4268 .create_broadcast("a/b/c/d", announce())
4269 .expect("publish allowed");
4270 let _keep3 = limited_producer
4271 .create_broadcast("a/b/c/d/e", announce())
4272 .expect("publish allowed");
4273 settle().await;
4274
4275 assert!(limited_producer.create_broadcast("a", announce()).is_err());
4277 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
4278 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
4279 }
4280
4281 #[tokio::test]
4282 async fn test_multiple_consumers_with_different_permissions() {
4283 let origin = Origin::random().produce();
4284
4285 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
4287 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
4288 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
4289 settle().await;
4290
4291 let mut foo_consumer = origin
4293 .consume()
4294 .scope(&["foo".into()])
4295 .expect("should create foo consumer")
4296 .announced();
4297
4298 let mut bar_consumer = origin
4299 .consume()
4300 .scope(&["bar".into()])
4301 .expect("should create bar consumer")
4302 .announced();
4303
4304 let mut foobar_consumer = origin
4305 .consume()
4306 .scope(&["foo".into(), "bar".into()])
4307 .expect("should create foobar consumer")
4308 .announced();
4309
4310 foo_consumer.assert_next_some("foo/test");
4312 foo_consumer.assert_next_wait();
4313
4314 bar_consumer.assert_next_some("bar/test");
4315 bar_consumer.assert_next_wait();
4316
4317 foobar_consumer.assert_next_some("bar/test");
4318 foobar_consumer.assert_next_some("foo/test");
4319 foobar_consumer.assert_next_wait();
4320 }
4321
4322 #[tokio::test]
4323 async fn test_select_with_empty_prefix() {
4324 let origin = Origin::random().produce();
4325
4326 let demo_producer = origin.with_root("demo").expect("should create demo root");
4328 let limited_producer = demo_producer
4329 .scope(&["worm-node".into(), "foobar".into()])
4330 .expect("should create limited producer");
4331
4332 let _broadcast1 = limited_producer
4334 .create_broadcast("worm-node/test", announce())
4335 .expect("publish allowed");
4336 let _broadcast2 = limited_producer
4337 .create_broadcast("foobar/test", announce())
4338 .expect("publish allowed");
4339 settle().await;
4340
4341 let mut consumer = limited_producer
4343 .consume()
4344 .scope(&["".into()])
4345 .expect("should create consumer with empty prefix")
4346 .announced();
4347
4348 let a1 = consumer.try_next().expect("expected first announcement");
4350 let a2 = consumer.try_next().expect("expected second announcement");
4351 consumer.assert_next_wait();
4352
4353 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4354 paths.sort();
4355 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
4356 }
4357
4358 #[tokio::test]
4359 async fn test_select_narrowing_scope() {
4360 let origin = Origin::random().produce();
4361
4362 let demo_producer = origin.with_root("demo").expect("should create demo root");
4364 let limited_producer = demo_producer
4365 .scope(&["worm-node".into(), "foobar".into()])
4366 .expect("should create limited producer");
4367
4368 let _broadcast1 = limited_producer
4370 .create_broadcast("worm-node", announce())
4371 .expect("publish allowed");
4372 let _broadcast2 = limited_producer
4373 .create_broadcast("worm-node/foo", announce())
4374 .expect("publish allowed");
4375 let _broadcast3 = limited_producer
4376 .create_broadcast("foobar/bar", announce())
4377 .expect("publish allowed");
4378 settle().await;
4379
4380 let mut worm_consumer = limited_producer
4382 .consume()
4383 .scope(&["worm-node".into()])
4384 .expect("should create worm-node consumer")
4385 .announced();
4386
4387 worm_consumer.assert_next_some("worm-node");
4389 worm_consumer.assert_next_some("worm-node/foo");
4390 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
4394 .consume()
4395 .scope(&["worm-node/foo".into()])
4396 .expect("should create worm-node/foo consumer")
4397 .announced();
4398
4399 foo_consumer.assert_next_some("worm-node/foo");
4400 foo_consumer.assert_next_wait(); }
4402
4403 #[tokio::test]
4404 async fn test_select_multiple_roots_with_empty_prefix() {
4405 let origin = Origin::random().produce();
4406
4407 let limited_producer = origin
4409 .scope(&["app1".into(), "app2".into(), "shared".into()])
4410 .expect("should create limited producer");
4411
4412 let _broadcast1 = limited_producer
4414 .create_broadcast("app1/data", announce())
4415 .expect("publish allowed");
4416 let _broadcast2 = limited_producer
4417 .create_broadcast("app2/config", announce())
4418 .expect("publish allowed");
4419 let _broadcast3 = limited_producer
4420 .create_broadcast("shared/resource", announce())
4421 .expect("publish allowed");
4422 settle().await;
4423
4424 let mut consumer = limited_producer
4426 .consume()
4427 .scope(&["".into()])
4428 .expect("should create consumer with empty prefix")
4429 .announced();
4430
4431 consumer.assert_next_some("app1/data");
4433 consumer.assert_next_some("app2/config");
4434 consumer.assert_next_some("shared/resource");
4435 consumer.assert_next_wait();
4436 }
4437
4438 #[tokio::test]
4439 async fn test_publish_scope_with_empty_prefix() {
4440 let origin = Origin::random().produce();
4441
4442 let limited_producer = origin
4444 .scope(&["services/api".into(), "services/web".into()])
4445 .expect("should create limited producer");
4446
4447 let same_producer = limited_producer
4449 .scope(&["".into()])
4450 .expect("should create producer with empty prefix");
4451
4452 let _broadcast = same_producer
4454 .create_broadcast("services/api", announce())
4455 .expect("publish allowed");
4456 let _keep2 = same_producer
4457 .create_broadcast("services/web", announce())
4458 .expect("publish allowed");
4459 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
4460 assert!(same_producer.create_broadcast("other", announce()).is_err());
4461 }
4462
4463 #[tokio::test]
4464 async fn test_select_narrowing_to_deeper_path() {
4465 let origin = Origin::random().produce();
4466
4467 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
4469
4470 let _broadcast1 = limited_producer
4472 .create_broadcast("org/team1/project1", announce())
4473 .expect("publish allowed");
4474 let _broadcast2 = limited_producer
4475 .create_broadcast("org/team1/project2", announce())
4476 .expect("publish allowed");
4477 let _broadcast3 = limited_producer
4478 .create_broadcast("org/team2/project1", announce())
4479 .expect("publish allowed");
4480 settle().await;
4481
4482 let mut team2_consumer = limited_producer
4484 .consume()
4485 .scope(&["org/team2".into()])
4486 .expect("should create team2 consumer")
4487 .announced();
4488
4489 team2_consumer.assert_next_some("org/team2/project1");
4490 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
4494 .consume()
4495 .scope(&["org/team1/project1".into()])
4496 .expect("should create project1 consumer")
4497 .announced();
4498
4499 project1_consumer.assert_next_some("org/team1/project1");
4501 project1_consumer.assert_next_wait();
4502 }
4503
4504 #[tokio::test]
4505 async fn test_select_with_non_matching_prefix() {
4506 let origin = Origin::random().produce();
4507
4508 let limited_producer = origin
4510 .scope(&["allowed/path".into()])
4511 .expect("should create limited producer");
4512
4513 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
4515
4516 assert!(limited_producer.scope(&["other/path".into()]).is_none());
4518 }
4519
4520 #[tokio::test]
4523 async fn test_with_root_trailing_slash_consumer() {
4524 let origin = Origin::random().produce();
4525
4526 let prefix = "some_prefix/".to_string();
4528 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4529
4530 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4531 settle().await;
4532 consumer.assert_next_some("test");
4533 }
4534
4535 #[tokio::test]
4537 async fn test_with_root_trailing_slash_producer() {
4538 let origin = Origin::random().produce();
4539
4540 let prefix = "some_prefix/".to_string();
4542 let rooted = origin.with_root(prefix).unwrap();
4543
4544 let _b = rooted.create_broadcast("test", announce()).unwrap();
4545 settle().await;
4546
4547 let mut consumer = rooted.consume().announced();
4548 consumer.assert_next_some("test");
4549 }
4550
4551 #[tokio::test]
4553 async fn test_with_root_trailing_slash_unannounce() {
4554 tokio::time::pause();
4555
4556 let origin = Origin::random().produce();
4557
4558 let prefix = "some_prefix/".to_string();
4559 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4560
4561 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4562 settle().await;
4563 consumer.assert_next_some("test");
4564
4565 b.finish();
4567 settle().await;
4568
4569 consumer.assert_next_none("test");
4571 }
4572
4573 #[tokio::test]
4574 async fn test_select_maintains_access_with_wider_prefix() {
4575 let origin = Origin::random().produce();
4576
4577 let demo_producer = origin.with_root("demo").expect("should create demo root");
4579 let user_producer = demo_producer
4580 .scope(&["worm-node".into(), "foobar".into()])
4581 .expect("should create user producer");
4582
4583 let _broadcast1 = user_producer
4585 .create_broadcast("worm-node/data", announce())
4586 .expect("publish allowed");
4587 let _broadcast2 = user_producer
4588 .create_broadcast("foobar", announce())
4589 .expect("publish allowed");
4590 settle().await;
4591
4592 let mut consumer = user_producer
4594 .consume()
4595 .scope(&["".into()])
4596 .expect("scope with empty prefix should not fail when user has specific permissions")
4597 .announced();
4598
4599 let a1 = consumer.try_next().expect("expected first announcement");
4601 let a2 = consumer.try_next().expect("expected second announcement");
4602 consumer.assert_next_wait();
4603
4604 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4605 paths.sort();
4606 assert_eq!(paths, ["foobar", "worm-node/data"]);
4607
4608 let mut narrow_consumer = user_producer
4610 .consume()
4611 .scope(&["worm-node".into()])
4612 .expect("should be able to narrow scope to worm-node")
4613 .announced();
4614
4615 narrow_consumer.assert_next_some("worm-node/data");
4616 narrow_consumer.assert_next_wait(); }
4618
4619 #[tokio::test]
4620 async fn test_duplicate_prefixes_deduped() {
4621 let origin = Origin::random().produce();
4622
4623 let producer = origin
4625 .scope(&["demo".into(), "demo".into()])
4626 .expect("should create producer");
4627
4628 let _broadcast = producer
4629 .create_broadcast("demo/stream", announce())
4630 .expect("publish allowed");
4631 settle().await;
4632
4633 let mut consumer = producer.consume().announced();
4634 consumer.assert_next_some("demo/stream");
4635 consumer.assert_next_wait();
4636 }
4637
4638 #[tokio::test]
4639 async fn test_overlapping_prefixes_deduped() {
4640 let origin = Origin::random().produce();
4641
4642 let producer = origin
4644 .scope(&["demo".into(), "demo/foo".into()])
4645 .expect("should create producer");
4646
4647 let _broadcast = producer
4649 .create_broadcast("demo/bar/stream", announce())
4650 .expect("publish allowed");
4651 settle().await;
4652
4653 let mut consumer = producer.consume().announced();
4654 consumer.assert_next_some("demo/bar/stream");
4655 consumer.assert_next_wait();
4656 }
4657
4658 #[tokio::test]
4659 async fn test_overlapping_prefixes_no_duplicate_announcements() {
4660 let origin = Origin::random().produce();
4661
4662 let producer = origin
4664 .scope(&["demo".into(), "demo/foo".into()])
4665 .expect("should create producer");
4666
4667 let _broadcast = producer
4668 .create_broadcast("demo/foo/stream", announce())
4669 .expect("publish allowed");
4670 settle().await;
4671
4672 let mut consumer = producer.consume().announced();
4673 consumer.assert_next_some("demo/foo/stream");
4675 consumer.assert_next_wait();
4676 }
4677
4678 #[tokio::test]
4679 async fn test_allowed_returns_deduped_prefixes() {
4680 let origin = Origin::random().produce();
4681
4682 let producer = origin
4683 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
4684 .expect("should create producer");
4685
4686 let allowed: Vec<_> = producer.allowed().collect();
4687 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
4688 }
4689
4690 #[tokio::test]
4691 async fn test_announced_broadcast_already_announced() {
4692 let origin = Origin::random().produce();
4693
4694 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4695 settle().await;
4696
4697 let consumer = origin.consume();
4698 let result = consumer.announced_broadcast("test").await.expect("should find it");
4699 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4700 }
4701
4702 #[tokio::test]
4703 async fn test_announced_broadcast_delayed() {
4704 tokio::time::pause();
4705
4706 let origin = Origin::random().produce();
4707
4708 let consumer = origin.consume();
4709
4710 let wait = tokio::spawn({
4712 let consumer = consumer.clone();
4713 async move { consumer.announced_broadcast("test").await }
4714 });
4715
4716 tokio::task::yield_now().await;
4718
4719 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4720 settle().await;
4721
4722 let result = wait.await.unwrap().expect("should find it");
4723 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4724 }
4725
4726 #[tokio::test]
4727 async fn test_announced_broadcast_ignores_unrelated_paths() {
4728 tokio::time::pause();
4729
4730 let origin = Origin::random().produce();
4731
4732 let consumer = origin.consume();
4733
4734 let wait = tokio::spawn({
4735 let consumer = consumer.clone();
4736 async move { consumer.announced_broadcast("target").await }
4737 });
4738
4739 tokio::task::yield_now().await;
4740
4741 let _other = origin.create_broadcast("other", announce()).unwrap();
4743 settle().await;
4744 tokio::task::yield_now().await;
4745 assert!(!wait.is_finished(), "must not resolve on unrelated path");
4746
4747 let _target = origin.create_broadcast("target", announce()).unwrap();
4748 settle().await;
4749 let result = wait.await.unwrap().expect("should find target");
4750 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
4751 }
4752
4753 #[tokio::test]
4754 async fn test_announced_broadcast_skips_nested_paths() {
4755 tokio::time::pause();
4756
4757 let origin = Origin::random().produce();
4758
4759 let consumer = origin.consume();
4760
4761 let wait = tokio::spawn({
4762 let consumer = consumer.clone();
4763 async move { consumer.announced_broadcast("foo").await }
4764 });
4765
4766 tokio::task::yield_now().await;
4767
4768 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
4770 settle().await;
4771 tokio::task::yield_now().await;
4772 assert!(!wait.is_finished(), "must not resolve on a nested path");
4773
4774 let _exact = origin.create_broadcast("foo", announce()).unwrap();
4775 settle().await;
4776 let result = wait.await.unwrap().expect("should find foo exactly");
4777 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
4778 }
4779
4780 #[tokio::test]
4781 async fn test_announced_broadcast_disallowed() {
4782 let origin = Origin::random().produce();
4783 let limited = origin
4784 .consume()
4785 .scope(&["allowed".into()])
4786 .expect("should create limited");
4787
4788 assert!(limited.announced_broadcast("notallowed").await.is_none());
4790 }
4791
4792 #[tokio::test]
4793 async fn test_announced_broadcast_scope_too_narrow() {
4794 let origin = Origin::random().produce();
4797 let limited = origin
4798 .consume()
4799 .scope(&["foo/specific".into()])
4800 .expect("should create limited");
4801
4802 let result = limited
4804 .announced_broadcast("foo")
4805 .now_or_never()
4806 .expect("must not block");
4807 assert!(result.is_none());
4808 }
4809
4810 #[tokio::test]
4814 async fn test_coalesce_announce_then_unannounce() {
4815 tokio::time::pause();
4817
4818 let origin = Origin::random().produce();
4819 let mut announced = origin.consume().announced();
4820
4821 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4822 settle().await;
4823 broadcast.finish();
4824
4825 settle().await;
4826
4827 announced.assert_next_wait();
4828 }
4829
4830 #[tokio::test]
4831 async fn test_coalesce_announce_unannounce_announce() {
4832 tokio::time::pause();
4835
4836 let origin = Origin::random().produce();
4837 let mut announced = origin.consume().announced();
4838
4839 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4840 settle().await;
4841 broadcast1.finish();
4842 settle().await;
4843 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4844 settle().await;
4845
4846 announced.assert_next_some("test");
4847 announced.assert_next_wait();
4848 }
4849
4850 #[tokio::test]
4851 async fn test_coalesce_unannounce_announce_preserved() {
4852 tokio::time::pause();
4855
4856 let origin = Origin::random().produce();
4857 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4858 settle().await;
4859
4860 let mut announced = origin.consume().announced();
4861 announced.assert_next_some("test");
4862
4863 broadcast1.finish();
4865 settle().await;
4866
4867 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4868 settle().await;
4869
4870 announced.assert_next_none("test");
4872 announced.assert_next_some("test");
4873 announced.assert_next_wait();
4874 }
4875
4876 #[tokio::test]
4877 async fn test_coalesce_unannounce_announce_unannounce() {
4878 tokio::time::pause();
4881
4882 let origin = Origin::random().produce();
4883 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4884 settle().await;
4885
4886 let mut announced = origin.consume().announced();
4887 announced.assert_next_some("test");
4888
4889 broadcast1.finish();
4890 settle().await;
4891
4892 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4893 settle().await;
4894 broadcast2.finish();
4895 settle().await;
4896
4897 announced.assert_next_none("test");
4898 announced.assert_next_wait();
4899 }
4900
4901 #[tokio::test]
4902 async fn test_coalesce_churn_bounded() {
4903 tokio::time::pause();
4908
4909 let origin = Origin::random().produce();
4910 let mut announced = origin.consume().announced();
4911
4912 for _ in 0..1000 {
4913 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4914 settle().await;
4915 broadcast.finish();
4916 }
4917 settle().await;
4918
4919 let mut collected = Vec::new();
4920 while let Some(update) = announced.try_next() {
4921 collected.push(update);
4922 }
4923 assert!(
4924 collected.len() <= 1,
4925 "expected at most one pending update, got {}",
4926 collected.len()
4927 );
4928 assert!(
4929 collected.iter().all(|a| a.path == Path::new("test")),
4930 "unexpected path in pending updates",
4931 );
4932 }
4933
4934 #[tokio::test]
4938 async fn test_consumer_clone_is_side_effect_free() {
4939 let origin = Origin::random().produce();
4940
4941 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
4942 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
4943 settle().await;
4944
4945 let consumer = origin.consume();
4946 let mut announced = consumer.announced();
4947
4948 for _ in 0..16 {
4951 let cloned = consumer.clone();
4952 assert!(cloned.get_broadcast("test1").is_some());
4953 assert!(cloned.get_broadcast("test2").is_some());
4954 }
4955
4956 let a1 = announced.try_next().expect("first announcement");
4959 let a2 = announced.try_next().expect("second announcement");
4960 announced.assert_next_wait();
4961
4962 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4963 paths.sort();
4964 assert_eq!(paths, ["test1", "test2"]);
4965
4966 let mut fresh = consumer.announced();
4968 let b1 = fresh.try_next().expect("backlog: first");
4969 let b2 = fresh.try_next().expect("backlog: second");
4970 fresh.assert_next_wait();
4971
4972 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
4973 paths.sort();
4974 assert_eq!(paths, ["test1", "test2"]);
4975 }
4976
4977 #[tokio::test]
4979 async fn dynamic_request_unroutable_without_handler() {
4980 let origin = Origin::random().produce();
4981 let consumer = origin.consume();
4982 assert!(matches!(
4983 consumer.request_broadcast("missing").await,
4984 Err(Error::Unroutable)
4985 ));
4986 }
4987
4988 #[tokio::test(start_paused = true)]
4991 async fn dynamic_request_served_not_announced() {
4992 let origin = Origin::random().produce();
4993 let mut dynamic = origin.dynamic();
4994 let consumer = origin.consume();
4995
4996 let mut announced = origin.consume().announced();
4998 announced.assert_next_wait();
4999
5000 let served = broadcast::Info::new().produce();
5001 let request_fut = consumer.request_broadcast("fallback");
5004
5005 let mut served_dynamic = served.dynamic();
5007
5008 let request = dynamic.requested_broadcast().await.unwrap();
5009 assert_eq!(request.path(), &Path::new("fallback"));
5010 request.accept(&served);
5011
5012 let broadcast = request_fut.await.unwrap();
5013 assert!(broadcast.is_clone(&served.consume()));
5014
5015 let track_fut = broadcast.track("video").unwrap().subscribe(None);
5017 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
5018 let mut track = track_fut.await.unwrap();
5019 producer.append_group().unwrap();
5020 track.assert_group();
5021
5022 announced.assert_next_wait();
5024 }
5025
5026 #[tokio::test(start_paused = true)]
5028 async fn dynamic_request_coalesces() {
5029 let origin = Origin::random().produce();
5030 let mut dynamic = origin.dynamic();
5031 let consumer = origin.consume();
5032
5033 let f1 = consumer.request_broadcast("dup");
5035 let f2 = consumer.request_broadcast("dup");
5036
5037 let request = dynamic.requested_broadcast().await.unwrap();
5039 assert_eq!(request.path(), &Path::new("dup"));
5040 assert!(
5041 dynamic.requested_broadcast().now_or_never().is_none(),
5042 "a coalesced request must not be served twice"
5043 );
5044
5045 let served = broadcast::Info::new().produce();
5047 request.accept(&served);
5048 assert!(f1.await.unwrap().is_clone(&served.consume()));
5049 assert!(f2.await.unwrap().is_clone(&served.consume()));
5050 }
5051
5052 #[tokio::test(start_paused = true)]
5055 async fn dynamic_request_dedups_served() {
5056 let origin = Origin::random().produce();
5057 let mut dynamic = origin.dynamic();
5058 let consumer = origin.consume();
5059
5060 let request_fut = consumer.request_broadcast("fallback");
5061 let request = dynamic.requested_broadcast().await.unwrap();
5062 let served = broadcast::Info::new().produce();
5063 request.accept(&served);
5064 let first = request_fut.await.unwrap();
5065 assert!(first.is_clone(&served.consume()));
5066
5067 let second = consumer.request_broadcast("fallback").await.unwrap();
5069 assert!(second.is_clone(&served.consume()));
5070
5071 assert!(
5073 dynamic.requested_broadcast().now_or_never().is_none(),
5074 "a still-live served broadcast must not be re-requested from the handler"
5075 );
5076 }
5077
5078 #[tokio::test(start_paused = true)]
5080 async fn dynamic_request_reserves_after_close() {
5081 let origin = Origin::random().produce();
5082 let mut dynamic = origin.dynamic();
5083 let consumer = origin.consume();
5084
5085 let request_fut = consumer.request_broadcast("fallback");
5086 let request = dynamic.requested_broadcast().await.unwrap();
5087 let served = broadcast::Info::new().produce();
5088 request.accept(&served);
5089 request_fut.await.unwrap();
5090
5091 drop(served);
5093
5094 let request_fut = consumer.request_broadcast("fallback");
5096 let request = dynamic.requested_broadcast().await.unwrap();
5097 assert_eq!(request.path(), &Path::new("fallback"));
5098 let served = broadcast::Info::new().produce();
5099 request.accept(&served);
5100 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
5101 }
5102
5103 #[tokio::test(start_paused = true)]
5106 async fn dynamic_request_served_cache_bounded() {
5107 let origin = Origin::random().produce();
5108 let mut dynamic = origin.dynamic();
5109 let consumer = origin.consume();
5110
5111 for i in 0..100 {
5112 let path = format!("one-shot/{i}");
5113 let request_fut = consumer.request_broadcast(&path);
5114 let request = dynamic.requested_broadcast().await.unwrap();
5115 let served = broadcast::Info::new().produce();
5116 request.accept(&served);
5117 request_fut.await.unwrap();
5118 drop(served);
5120 }
5121
5122 assert!(
5125 origin.dynamic.read().served.len() <= 4,
5126 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
5127 origin.dynamic.read().served.len()
5128 );
5129 }
5130
5131 #[tokio::test(start_paused = true)]
5134 async fn dynamic_request_coalesces_after_handoff() {
5135 let origin = Origin::random().produce();
5136 let mut dynamic = origin.dynamic();
5137 let consumer = origin.consume();
5138
5139 let f1 = consumer.request_broadcast("fallback");
5140 let request = dynamic.requested_broadcast().await.unwrap();
5142
5143 let f2 = consumer.request_broadcast("fallback");
5145 assert!(
5146 dynamic.requested_broadcast().now_or_never().is_none(),
5147 "a repeat request during hand-off must coalesce, not re-queue"
5148 );
5149
5150 let served = broadcast::Info::new().produce();
5152 request.accept(&served);
5153 assert!(f1.await.unwrap().is_clone(&served.consume()));
5154 assert!(f2.await.unwrap().is_clone(&served.consume()));
5155 }
5156
5157 #[tokio::test(start_paused = true)]
5159 async fn dynamic_request_dropped_after_handoff() {
5160 let origin = Origin::random().produce();
5161 let mut dynamic = origin.dynamic();
5162 let consumer = origin.consume();
5163
5164 let f1 = consumer.request_broadcast("fallback");
5165 let request = dynamic.requested_broadcast().await.unwrap();
5166 let f2 = consumer.request_broadcast("fallback");
5167
5168 drop(request);
5170 assert!(matches!(f1.await, Err(Error::Unroutable)));
5171 assert!(matches!(f2.await, Err(Error::Unroutable)));
5172 }
5173
5174 #[tokio::test(start_paused = true)]
5176 async fn dynamic_request_rejected() {
5177 let origin = Origin::random().produce();
5178 let mut dynamic = origin.dynamic();
5179 let consumer = origin.consume();
5180
5181 let request_fut = consumer.request_broadcast("fallback");
5182
5183 let request = dynamic.requested_broadcast().await.unwrap();
5184 request.reject(Error::Cancel);
5185
5186 assert!(matches!(request_fut.await, Err(Error::Cancel)));
5187 }
5188
5189 #[tokio::test(start_paused = true)]
5193 async fn dynamic_request_rerequest_after_reject() {
5194 let origin = Origin::random().produce();
5195 let mut dynamic = origin.dynamic();
5196 let consumer = origin.consume();
5197
5198 let f1 = consumer.request_broadcast("fallback");
5199 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
5200 assert!(matches!(f1.await, Err(Error::Unroutable)));
5201
5202 let served = broadcast::Info::new().produce();
5203 let f2 = consumer.request_broadcast("fallback");
5205 let request = dynamic.requested_broadcast().await.unwrap();
5206 assert_eq!(request.path(), &Path::new("fallback"));
5207 request.accept(&served);
5208 assert!(f2.await.unwrap().is_clone(&served.consume()));
5209 }
5210
5211 #[tokio::test(start_paused = true)]
5214 async fn dynamic_request_handler_dropped() {
5215 let origin = Origin::random().produce();
5216 let dynamic = origin.dynamic();
5217 let consumer = origin.consume();
5218
5219 let request_fut = consumer.request_broadcast("fallback");
5220 drop(dynamic);
5221 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
5222
5223 assert!(matches!(
5225 consumer.request_broadcast("again").await,
5226 Err(Error::Unroutable)
5227 ));
5228 }
5229
5230 #[tokio::test(start_paused = true)]
5234 async fn dynamic_request_accept_after_handler_dropped() {
5235 let origin = Origin::random().produce();
5236 let mut dynamic = origin.dynamic();
5237 let consumer = origin.consume();
5238
5239 let request_fut = consumer.request_broadcast("fallback");
5240
5241 let request = dynamic.requested_broadcast().await.unwrap();
5243 drop(dynamic);
5244
5245 let served = broadcast::Info::new().produce();
5246 request.accept(&served);
5248 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
5249 }
5250
5251 #[tokio::test(start_paused = true)]
5253 async fn dynamic_request_prefers_announced() {
5254 let origin = Origin::random().produce();
5255 let mut dynamic = origin.dynamic();
5256 let consumer = origin.consume();
5257
5258 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
5259 settle().await;
5260
5261 let got = consumer.request_broadcast("live").await.unwrap();
5262 assert!(
5263 got.is_clone(&consumer.get_broadcast("live").unwrap()),
5264 "should return the published broadcast"
5265 );
5266 assert!(
5267 dynamic.requested_broadcast().now_or_never().is_none(),
5268 "a published path must not queue a fallback request"
5269 );
5270 }
5271
5272 #[tokio::test(start_paused = true)]
5274 async fn dynamic_clone_keeps_alive() {
5275 let origin = Origin::random().produce();
5276 let dynamic = origin.dynamic();
5277 let consumer = origin.consume();
5278
5279 drop(dynamic.clone());
5280
5281 let request_fut = consumer.request_broadcast("fallback");
5284 assert!(
5285 request_fut.now_or_never().is_none(),
5286 "request should stay pending until served"
5287 );
5288 }
5289}