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: Option<web_async::time::Instant> = None;
1417
1418 loop {
1419 let empty = {
1420 let s = state.read();
1421 !s.closed && s.routes.is_empty()
1422 };
1423 deadline = match (empty, 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 let mut sleep = std::pin::pin!(async {
1435 match deadline {
1436 Some(at) => {
1437 web_async::time::sleep(at.saturating_duration_since(web_async::time::Instant::now())).await
1438 }
1439 None => std::future::pending().await,
1440 }
1441 });
1442 let mut fired = false;
1443 kio::wait(|waiter| {
1444 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1445 return Poll::Ready(Step::Serve(name, resume));
1446 }
1447 match state.poll(waiter, |s| {
1450 if s.closed || s.routes.is_empty() != empty {
1451 Poll::Ready(())
1452 } else {
1453 Poll::Pending
1454 }
1455 }) {
1456 Poll::Ready(Ok(guard)) => {
1457 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1458 }
1459 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1460 Poll::Pending => {}
1461 }
1462 if deadline.is_some() && !fired && waiter.poll_future(sleep.as_mut()).is_ready() {
1463 fired = true;
1464 }
1465 match fired {
1466 true => Poll::Ready(Step::Expired),
1467 false => Poll::Pending,
1468 }
1469 })
1470 .await
1471 };
1472
1473 match step {
1474 Step::Serve(name, resume) => {
1475 web_async::spawn(serve_track(state.clone(), name, resume));
1478 }
1479 Step::Changed => {}
1480 Step::Expired => {
1481 let close = {
1485 let Ok(mut s) = state.write() else { break };
1486 if !s.closed && s.routes.is_empty() {
1487 s.closed = true;
1488 true
1489 } else {
1490 false
1491 }
1492 };
1493 if close {
1494 break;
1495 }
1496 }
1497 Step::Closed => break,
1498 }
1499 }
1500
1501 broadcast.abort_spliced(Error::Dropped);
1503
1504 broadcast.finish();
1506
1507 node.lock().remove(&state, &rest);
1510}
1511
1512async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1520 enum Step {
1521 Closed,
1522 Splice(u64, broadcast::Consumer),
1523 Complete,
1524 Failed,
1525 Idle,
1527 Demand,
1529 }
1530
1531 let mut fails = 0u32;
1532 let mut serving: Option<(u64, track::Consumer)> = None;
1534 let mut dead: Option<u64> = None;
1538 let mut idle_since: Option<web_async::time::Instant> = None;
1540
1541 loop {
1542 let serving_id = serving.as_ref().map(|(id, _)| *id);
1543
1544 let used = resume.is_used();
1548 idle_since = match (serving.is_some(), used) {
1549 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
1550 _ => None,
1551 };
1552 let deadline = idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER));
1553
1554 let step = {
1555 let mut sleep = std::pin::pin!(async {
1559 match deadline {
1560 Some(at) => {
1561 web_async::time::sleep(at.saturating_duration_since(web_async::time::Instant::now())).await
1562 }
1563 None => std::future::pending().await,
1564 }
1565 });
1566 let mut fired = false;
1567
1568 kio::wait(|waiter| {
1569 match state.poll(waiter, |s| {
1573 if s.closed
1574 || (used
1575 && matches!(s.active, Some(active) if Some(active) != serving_id && Some(active) != dead))
1576 {
1577 Poll::Ready(())
1578 } else {
1579 Poll::Pending
1580 }
1581 }) {
1582 Poll::Ready(Ok(guard)) => {
1583 if guard.closed {
1584 return Poll::Ready(Step::Closed);
1585 }
1586 let active = guard.active.expect("predicate guaranteed an active source");
1587 let source = guard
1588 .routes
1589 .iter()
1590 .find(|r| r.id == active)
1591 .expect("active source in table")
1592 .source
1593 .clone();
1594 return Poll::Ready(Step::Splice(active, source));
1595 }
1596 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1597 Poll::Pending => {}
1598 }
1599
1600 let edge = match used {
1605 true => resume.poll_unused(waiter),
1606 false => resume.poll_used(waiter),
1607 };
1608 if edge.is_ready() {
1609 return Poll::Ready(Step::Demand);
1610 }
1611
1612 if let Some((_, track)) = &serving
1615 && let Poll::Ready(result) = track.poll_complete(waiter)
1616 {
1617 return Poll::Ready(match result {
1618 Ok(()) => Step::Complete,
1619 Err(_) => Step::Failed,
1620 });
1621 }
1622
1623 if deadline.is_some() && !fired && waiter.poll_future(sleep.as_mut()).is_ready() {
1624 fired = true;
1625 }
1626 if fired {
1627 return Poll::Ready(Step::Idle);
1628 }
1629 Poll::Pending
1630 })
1631 .await
1632 };
1633
1634 match step {
1635 Step::Closed => return,
1637 Step::Complete => {
1638 let _ = resume.finish();
1639 return;
1640 }
1641 Step::Failed => {
1642 serving = None;
1645 }
1646 Step::Demand => {}
1648 Step::Idle => {
1649 if resume.release().is_err() {
1654 return;
1656 }
1657 serving = None;
1658 }
1659 Step::Splice(id, source) => {
1660 let attempt = match source.track(&name) {
1664 Ok(track) => {
1665 let query = track.info().into_inner();
1668 let info = kio::wait(|waiter| {
1669 if let Poll::Ready(result) = query.poll(waiter) {
1670 return Poll::Ready(Some(result));
1671 }
1672 match state.poll(waiter, |s| {
1673 if s.closed || s.active != Some(id) {
1674 Poll::Ready(())
1675 } else {
1676 Poll::Pending
1677 }
1678 }) {
1679 Poll::Ready(_) => Poll::Ready(None),
1680 Poll::Pending => Poll::Pending,
1681 }
1682 })
1683 .await;
1684 match info {
1685 None => continue,
1688 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
1692 Poll::Ready(Err(err)) => Err(err),
1693 _ => Ok(track),
1694 },
1695 Some(Err(err)) => Err(err),
1696 }
1697 }
1698 Err(err) => Err(err),
1699 };
1700
1701 match attempt {
1702 Ok(track) => {
1703 if resume.takeover(&track).is_err() {
1704 return;
1707 }
1708 fails = 0;
1711 dead = None;
1712 serving = Some((id, track));
1713 }
1714 Err(_) if source.is_closing() => {
1718 dead = Some(id);
1719 serving = None;
1720 }
1721 Err(err) => {
1722 fails += 1;
1723 if fails >= MAX_TRACK_RETRIES {
1724 tracing::debug!(name = %name, %err, "aborting unservable track");
1725 let _ = resume.abort(Error::Unroutable);
1726 return;
1727 }
1728 serving = None;
1729 }
1730 }
1731 }
1732 }
1733 }
1734}
1735
1736#[derive(Default)]
1742struct OriginDynamicState {
1743 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
1746
1747 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
1753}
1754
1755#[derive(Default)]
1762struct PendingBroadcast {
1763 resolved: Option<Result<broadcast::Consumer, Error>>,
1764}
1765
1766pub struct Dynamic {
1777 info: Origin,
1778 root: PathOwned,
1779 state: kio::Shared<OriginDynamicState>,
1780}
1781
1782impl Clone for Dynamic {
1783 fn clone(&self) -> Self {
1784 self.state.lock().requests.add_handler();
1788
1789 Self {
1790 info: self.info,
1791 root: self.root.clone(),
1792 state: self.state.clone(),
1793 }
1794 }
1795}
1796
1797impl Dynamic {
1798 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
1799 state.lock().requests.add_handler();
1800
1801 Self { info, root, state }
1802 }
1803
1804 pub fn info(&self) -> &Origin {
1806 &self.info
1807 }
1808
1809 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
1811 let mut state = ready!(self.state.poll(waiter, |state| {
1812 if state.requests.has_queued() {
1813 Poll::Ready(())
1814 } else {
1815 Poll::Pending
1816 }
1817 }));
1818
1819 let path = state.requests.pop().expect("predicate guaranteed a request");
1820 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
1826 Poll::Ready(Ok(Request {
1827 path,
1828 producer,
1829 state: self.state.clone(),
1830 }))
1831 }
1832
1833 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
1836 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
1837 }
1838
1839 pub fn root(&self) -> &Path<'_> {
1841 &self.root
1842 }
1843}
1844
1845impl Drop for Dynamic {
1846 fn drop(&mut self) {
1847 let mut state = self.state.lock();
1850 if state.requests.remove_handler() {
1851 state.requests.drain_queued();
1855 }
1856 }
1857}
1858
1859pub struct Request {
1866 path: PathOwned,
1868
1869 producer: kio::Producer<PendingBroadcast>,
1872
1873 state: kio::Shared<OriginDynamicState>,
1875}
1876
1877impl Request {
1878 pub fn path(&self) -> &Path<'_> {
1880 &self.path
1881 }
1882
1883 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
1889 let broadcast = broadcast.consume();
1890
1891 let resolved = {
1897 let mut state = self.state.lock();
1898 let existing = state.served.insert(self.path.clone(), broadcast.weak());
1899 state
1900 .requests
1901 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1902 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
1903 };
1904
1905 if let Ok(mut pending) = self.producer.write() {
1906 pending.resolved = Some(Ok(resolved));
1907 }
1908 }
1910
1911 pub fn reject(self, err: Error) {
1913 self.state
1914 .lock()
1915 .requests
1916 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1917 if let Ok(mut state) = self.producer.write() {
1918 state.resolved = Some(Err(err));
1919 }
1920 }
1921}
1922
1923impl Drop for Request {
1924 fn drop(&mut self) {
1925 self.state
1933 .lock()
1934 .requests
1935 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1936 }
1937}
1938
1939pub struct Requesting {
1946 inner: RequestState,
1947 stats: stats::Scope,
1950}
1951
1952enum RequestState {
1953 Ready(broadcast::Consumer),
1955 Failed(Error),
1958 Pending(kio::Consumer<PendingBroadcast>),
1960}
1961
1962impl Requesting {
1963 fn ready(broadcast: broadcast::Consumer) -> Self {
1964 Self {
1965 inner: RequestState::Ready(broadcast),
1966 stats: stats::Scope::default(),
1967 }
1968 }
1969
1970 fn failed(error: Error) -> Self {
1971 Self {
1972 inner: RequestState::Failed(error),
1973 stats: stats::Scope::default(),
1974 }
1975 }
1976
1977 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
1978 Self {
1979 inner: RequestState::Pending(consumer),
1980 stats: stats::Scope::default(),
1981 }
1982 }
1983
1984 fn with_stats(mut self, scope: stats::Scope) -> Self {
1985 self.stats = scope;
1986 self
1987 }
1988
1989 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
1991 match &self.inner {
1992 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
1993 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
1994 RequestState::Pending(consumer) => Poll::Ready(
1995 match ready!(consumer.poll(waiter, |state| match &state.resolved {
1996 Some(result) => Poll::Ready(result.clone()),
1997 None => Poll::Pending,
1998 })) {
1999 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2000 Err(_closed) => Err(Error::Unroutable),
2002 },
2003 ),
2004 }
2005 }
2006}
2007
2008impl kio::Pollable for Requesting {
2009 type Output = Result<broadcast::Consumer, Error>;
2010
2011 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2012 self.poll_ok(waiter)
2013 }
2014}
2015
2016pub trait Consume<T> {
2024 fn consume(&self) -> T;
2026}
2027
2028impl<T, U: Consume<T>> Consume<T> for &U {
2029 fn consume(&self) -> T {
2030 (**self).consume()
2031 }
2032}
2033
2034impl Consume<Consumer> for Producer {
2035 fn consume(&self) -> Consumer {
2036 Consumer::new(
2040 self.info,
2041 self.root.clone(),
2042 self.nodes.clone(),
2043 self.dynamic.clone(),
2044 stats::Session::default(),
2045 )
2046 }
2047}
2048
2049impl Consume<Consumer> for Consumer {
2050 fn consume(&self) -> Consumer {
2051 self.clone()
2052 }
2053}
2054
2055impl Consume<broadcast::Consumer> for broadcast::Producer {
2056 fn consume(&self) -> broadcast::Consumer {
2057 self.consume()
2059 }
2060}
2061
2062impl Consume<broadcast::Consumer> for broadcast::Consumer {
2063 fn consume(&self) -> broadcast::Consumer {
2064 self.clone()
2065 }
2066}
2067
2068impl Consume<track::Consumer> for track::Producer {
2069 fn consume(&self) -> track::Consumer {
2070 self.consume()
2071 }
2072}
2073
2074impl Consume<track::Consumer> for track::Consumer {
2075 fn consume(&self) -> track::Consumer {
2076 self.clone()
2077 }
2078}
2079
2080#[derive(Clone)]
2086pub struct Consumer {
2087 info: Origin,
2089 nodes: OriginNodes,
2090
2091 root: PathOwned,
2093
2094 dynamic: kio::Shared<OriginDynamicState>,
2097
2098 stats: stats::Session,
2102}
2103
2104impl std::ops::Deref for Consumer {
2105 type Target = Origin;
2106
2107 fn deref(&self) -> &Self::Target {
2108 &self.info
2109 }
2110}
2111
2112impl Consumer {
2113 fn new(
2114 info: Origin,
2115 root: PathOwned,
2116 nodes: OriginNodes,
2117 dynamic: kio::Shared<OriginDynamicState>,
2118 stats: stats::Session,
2119 ) -> Self {
2120 Self {
2121 info,
2122 nodes,
2123 root,
2124 dynamic,
2125 stats,
2126 }
2127 }
2128
2129 pub fn with_stats(mut self, session: stats::Session) -> Self {
2133 self.stats = session;
2134 self
2135 }
2136
2137 fn untagged(&self) -> Self {
2141 Self {
2142 stats: stats::Session::default(),
2143 ..self.clone()
2144 }
2145 }
2146
2147 pub(crate) fn empty(&self) -> Self {
2152 Self {
2153 info: self.info,
2154 nodes: OriginNodes { nodes: Vec::new() },
2155 root: self.root.clone(),
2156 dynamic: self.dynamic.clone(),
2157 stats: self.stats.clone(),
2158 }
2159 }
2160
2161 pub fn announced(&self) -> AnnounceConsumer {
2168 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2169 }
2170
2171 pub fn consume(&self) -> Self {
2173 self.clone()
2174 }
2175
2176 fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2183 let path = path.as_path();
2184 let (root, rest) = self.nodes.get(&path)?;
2185 let state = root.lock();
2186 state.consume_broadcast(&rest)
2187 }
2188
2189 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2201 let path = path.as_path();
2202
2203 let consumer = self.scope(std::slice::from_ref(&path))?;
2205
2206 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2210 return None;
2211 }
2212
2213 let mut announced = consumer.untagged().announced();
2217 let scope = self.stats.egress(self.root.join(&path).to_owned());
2218 loop {
2219 let OriginAnnounce {
2220 path: announced_path,
2221 broadcast,
2222 } = announced.next().await?;
2223 if announced_path.as_path() == path
2225 && let Some(broadcast) = broadcast
2226 {
2227 return Some(broadcast.with_stats(scope));
2228 }
2229 }
2230 }
2231
2232 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2238 let prefixes = PathPrefixes::new(prefixes);
2239 Some(Consumer::new(
2240 self.info,
2241 self.root.clone(),
2242 self.nodes.select(&prefixes)?,
2243 self.dynamic.clone(),
2244 self.stats.clone(),
2245 ))
2246 }
2247
2248 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2267 let path = path.as_path();
2268
2269 let absolute = self.root.join(&path).to_owned();
2273 let scope = self.stats.egress(&absolute);
2274
2275 if let Some(broadcast) = self.get_broadcast(&path) {
2277 return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope));
2278 }
2279
2280 let mut state = self.dynamic.lock();
2281
2282 if let Some(weak) = state.served.get(&absolute) {
2286 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2287 }
2288
2289 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2292 producer.consume()
2293 } else {
2294 let producer = kio::Producer::<PendingBroadcast>::default();
2295 let consumer = producer.consume();
2296 if state.requests.insert(absolute, producer).is_err() {
2297 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2298 }
2299 consumer
2300 };
2301
2302 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2303 }
2304
2305 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2310 let prefix = prefix.as_path();
2311
2312 Some(Self::new(
2313 self.info,
2314 self.root.join(&prefix).to_owned(),
2315 self.nodes.root(&prefix)?,
2316 self.dynamic.clone(),
2317 self.stats.clone(),
2318 ))
2319 }
2320
2321 pub fn root(&self) -> &Path<'_> {
2323 &self.root
2324 }
2325
2326 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2329 self.nodes.nodes.iter().map(|(root, _)| root)
2330 }
2331
2332 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2334 self.root.join(path)
2335 }
2336}
2337
2338#[derive(Clone)]
2343pub struct AnnounceProducer {
2344 nodes: OriginNodes,
2345 root: PathOwned,
2346}
2347
2348impl AnnounceProducer {
2349 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2350 Self { nodes, root }
2351 }
2352
2353 pub fn consume(&self) -> AnnounceConsumer {
2359 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2362 }
2363
2364 pub fn root(&self) -> &Path<'_> {
2366 &self.root
2367 }
2368}
2369
2370pub struct AnnounceConsumer {
2375 id: ConsumerId,
2376 nodes: OriginNodes,
2377 root: PathOwned,
2378
2379 state: kio::Producer<OriginConsumerState>,
2382
2383 stats: stats::Session,
2386
2387 guards: HashMap<PathOwned, stats::Announce>,
2391}
2392
2393impl AnnounceConsumer {
2394 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2395 let state = kio::Producer::<OriginConsumerState>::default();
2396 let id = ConsumerId::new();
2397
2398 for (_, node) in &nodes.nodes {
2399 let notify = AnnounceConsumerNotify {
2400 root: root.clone(),
2401 state: state.clone(),
2402 };
2403 node.lock().consume(id, notify);
2404 }
2405
2406 Self {
2407 id,
2408 nodes,
2409 root,
2410 state,
2411 stats,
2412 guards: HashMap::new(),
2413 }
2414 }
2415
2416 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2422 let OriginAnnounce { path, broadcast } = update;
2423 let absolute = self.root.join(&path).to_owned();
2424 match broadcast {
2425 Some(broadcast) => {
2426 let scope = self.stats.egress(&absolute);
2427 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2428 OriginAnnounce {
2429 path,
2430 broadcast: Some(broadcast.with_stats(scope)),
2431 }
2432 }
2433 None => {
2434 self.guards.remove(&absolute);
2435 OriginAnnounce { path, broadcast: None }
2436 }
2437 }
2438 }
2439
2440 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2447 kio::wait(|waiter| self.poll_next(waiter)).await
2448 }
2449
2450 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2456 let update = {
2457 let mut state = match ready!(self.state.poll(waiter, |state| {
2458 if state.pending.is_empty() {
2459 Poll::Pending
2460 } else {
2461 Poll::Ready(())
2462 }
2463 })) {
2464 Ok(state) => state,
2465 Err(_) => return Poll::Ready(None),
2467 };
2468 state.take().expect("predicate guaranteed an update")
2469 };
2470 Poll::Ready(Some(self.attribute(update)))
2471 }
2472
2473 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2478 let update = self.state.write().ok()?.take()?;
2479 Some(self.attribute(update))
2480 }
2481
2482 pub fn is_closed(&self) -> bool {
2484 self.state.write().is_err()
2485 }
2486
2487 pub fn root(&self) -> &Path<'_> {
2489 &self.root
2490 }
2491
2492 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2494 self.root.join(path)
2495 }
2496}
2497
2498impl Drop for AnnounceConsumer {
2499 fn drop(&mut self) {
2500 for (_, root) in &self.nodes.nodes {
2501 root.lock().unconsume(self.id);
2502 }
2503 }
2504}
2505
2506#[cfg(test)]
2507use futures::FutureExt;
2508
2509#[cfg(test)]
2510#[allow(missing_docs)] impl AnnounceConsumer {
2512 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2513 let expected = expected.as_path();
2514 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2515 assert_eq!(announce.path, expected, "wrong path");
2516 let announced = announce.broadcast.expect("should be an active announce");
2517 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2518 }
2519
2520 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2524 let expected = expected.as_path();
2525 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2526 assert_eq!(announce.path, expected, "wrong path");
2527 announce.broadcast.expect("should be an active announce")
2528 }
2529
2530 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2531 let expected = expected.as_path();
2532 let announce = self.try_next().expect("no next");
2533 assert_eq!(announce.path, expected, "wrong path");
2534 let announced = announce.broadcast.expect("should be an active announce");
2535 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2536 }
2537
2538 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2540 let expected = expected.as_path();
2541 let announce = self.try_next().expect("no next");
2542 assert_eq!(announce.path, expected, "wrong path");
2543 announce.broadcast.expect("should be an active announce")
2544 }
2545
2546 pub fn assert_next_none(&mut self, expected: impl AsPath) {
2547 let expected = expected.as_path();
2548 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2549 assert_eq!(announce.path, expected, "wrong path");
2550 assert!(announce.broadcast.is_none(), "should be unannounced");
2551 }
2552
2553 pub fn assert_next_wait(&mut self) {
2554 if let Some(res) = self.next().now_or_never() {
2555 panic!("next should block: got {:?}", res.map(|a| a.path));
2556 }
2557 }
2558
2559 }
2568
2569#[cfg(test)]
2570mod tests {
2571 use crate::coding::Decode;
2572 use crate::group;
2573
2574 use super::*;
2575
2576 fn announce() -> broadcast::Route {
2578 broadcast::Route::new().with_announce(true)
2579 }
2580
2581 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
2587 let name = Path::new(name);
2588 let peer_key = fnv_key(&name, [peer]);
2589 (100u64..)
2590 .map(|id| Origin::new(id).unwrap())
2591 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
2592 .unwrap()
2593 }
2594
2595 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
2598 let source = broadcast::Info::new().produce().consume();
2599 FrontState {
2600 path: Path::new("test").to_owned(),
2601 self_origin,
2602 next_route: routes.len() as u64,
2603 routes: routes
2604 .into_iter()
2605 .enumerate()
2606 .map(|(id, route)| FrontRoute {
2607 id: id as u64,
2608 route,
2609 source: source.clone(),
2610 })
2611 .collect(),
2612 active: Some(0),
2613 linger: Duration::ZERO,
2614 closed: false,
2615 }
2616 }
2617
2618 fn sibling_route(peer: Origin) -> broadcast::Route {
2621 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
2622 announce().with_hops(hops)
2623 }
2624
2625 fn upstream_route(cost: u64) -> broadcast::Route {
2627 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
2628 announce().with_hops(hops).with_cost(cost)
2629 }
2630
2631 #[test]
2635 fn test_carrying_gate_keys() {
2636 let peer = Origin::new(3).unwrap();
2637
2638 let mut lost = front_state(
2640 origin_keyed("test", peer, false),
2641 vec![upstream_route(10), sibling_route(peer)],
2642 );
2643 lost.reselect(true);
2644 assert_eq!(
2645 lost.active,
2646 Some(0),
2647 "carrying front re-parented onto a higher-keyed peer"
2648 );
2649 lost.reselect(false);
2650 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
2651
2652 let mut won = front_state(
2654 origin_keyed("test", peer, true),
2655 vec![upstream_route(10), sibling_route(peer)],
2656 );
2657 won.reselect(true);
2658 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
2659 }
2660
2661 #[test]
2666 fn test_carrying_gate_symmetric_race() {
2667 let a = Origin::new(1).unwrap();
2668 let b = Origin::new(2).unwrap();
2669
2670 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
2671 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
2672 a_view.reselect(true);
2673 b_view.reselect(true);
2674
2675 let a_moved = a_view.active == Some(1);
2676 let b_moved = b_view.active == Some(1);
2677 assert!(
2678 a_moved != b_moved,
2679 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
2680 );
2681 }
2682
2683 #[test]
2688 fn test_carrying_switches_to_benign_routes() {
2689 let peer = Origin::new(3).unwrap();
2690 let lost = origin_keyed("test", peer, false);
2691
2692 let mut forwarder = sibling_route(peer).with_cost(4);
2694 forwarder.advertised = 4;
2695 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
2696 state.reselect(true);
2697 assert_eq!(
2698 state.active,
2699 Some(1),
2700 "a cheaper forwarder path must win while carrying"
2701 );
2702
2703 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
2705 let mut state = front_state(lost, vec![upstream_route(10), direct]);
2706 state.reselect(true);
2707 assert_eq!(
2708 state.active,
2709 Some(1),
2710 "a direct publisher route must win while carrying"
2711 );
2712 }
2713
2714 #[test]
2717 fn test_carrying_gate_ignores_unannounced_incumbent() {
2718 let peer = Origin::new(3).unwrap();
2719 let unannounced = upstream_route(10).with_announce(false);
2720 let mut state = front_state(
2721 origin_keyed("test", peer, false),
2722 vec![unannounced, sibling_route(peer)],
2723 );
2724 state.reselect(true);
2725 assert_eq!(
2726 state.active,
2727 Some(1),
2728 "an unannounced incumbent must always be displaced"
2729 );
2730 }
2731
2732 async fn settle() {
2735 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
2736 }
2737
2738 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
2741 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
2742 .await
2743 .expect("timed out waiting for a track request")
2744 .expect("source closed");
2745 assert_eq!(request.name(), name, "unexpected track dispatched");
2746 request.accept(None)
2747 }
2748
2749 #[tokio::test]
2753 async fn test_stats_tagged_end_to_end() {
2754 use crate::Timestamp;
2755 use crate::stats::{Config, Registry, Tier};
2756 use bytes::Bytes;
2757
2758 tokio::time::pause();
2759
2760 let registry = Registry::new(Config::new());
2761 let ctx = registry.tier(Tier::default()).session("acme");
2762
2763 let origin = Origin::random().produce();
2764 let ingress = origin.clone().with_stats(ctx.clone());
2765 let egress = origin.consume().with_stats(ctx.clone());
2766
2767 let mut announced = egress.announced();
2770
2771 let source = ingress.create_broadcast("demo", announce()).unwrap();
2773 let mut dynamic = source.dynamic();
2774 settle().await;
2775 settle().await;
2776
2777 let update = announced.next().await.unwrap();
2779 assert_eq!(update.path.as_str(), "demo");
2780 let broadcast = update.broadcast.unwrap();
2781
2782 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2784 let mut producer = accept_track(&mut dynamic, "video").await;
2785 settle().await;
2786 let mut sub = subscribing.await.unwrap();
2787
2788 let mut group = producer.append_group().unwrap();
2790 group
2791 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2792 .unwrap();
2793 group
2794 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
2795 .unwrap();
2796 group.finish().unwrap();
2797
2798 let mut group_c = sub.recv_group().await.unwrap().unwrap();
2800 let mut frames = 0;
2801 while let Some(frame) = group_c.read_frame().await.unwrap() {
2802 assert_eq!(frame.payload.len(), 5);
2803 frames += 1;
2804 }
2805 assert_eq!(frames, 2);
2806 settle().await;
2807
2808 let report = registry.report();
2809 let entry = report
2810 .traffic
2811 .iter()
2812 .find(|e| e.path.as_str() == "demo")
2813 .expect("demo tracked");
2814 let path_len = "demo".len() as u64;
2815
2816 let egress = &entry.publisher;
2818 assert_eq!(egress.announced, 1, "one egress announce");
2819 assert_eq!(egress.announced_bytes, path_len);
2820 assert_eq!(egress.subscriptions, 1, "one egress subscription");
2821 assert_eq!(egress.broadcasts, 1, "one viewer");
2822 assert_eq!(egress.groups, 1);
2823 assert_eq!(egress.frames, 2);
2824 assert_eq!(egress.bytes, 10);
2825 assert_eq!(egress.fetches, 0);
2826
2827 let ingress = &entry.subscriber;
2829 assert_eq!(ingress.announced, 1, "one ingress announce");
2830 assert_eq!(ingress.announced_bytes, path_len);
2831 assert_eq!(ingress.subscriptions, 1, "one ingress track");
2832 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
2833 assert_eq!(ingress.groups, 1);
2834 assert_eq!(ingress.frames, 2);
2835 assert_eq!(ingress.bytes, 10);
2836
2837 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
2839 let _ = fetched;
2840 settle().await;
2841 let report = registry.report();
2842 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
2843 assert_eq!(entry.publisher.fetches, 1, "one fetch");
2844 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
2845 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
2846 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
2850 }
2851
2852 #[tokio::test]
2857 async fn test_stats_read_frame_counts_once() {
2858 use crate::Timestamp;
2859 use crate::stats::{Config, Registry, Tier};
2860 use bytes::Bytes;
2861
2862 tokio::time::pause();
2863
2864 let registry = Registry::new(Config::new());
2865 let ctx = registry.tier(Tier::default()).session("acme");
2866
2867 let origin = Origin::random().produce();
2868 let ingress = origin.clone().with_stats(ctx.clone());
2869 let egress = origin.consume().with_stats(ctx.clone());
2870
2871 let mut announced = egress.announced();
2872 let source = ingress.create_broadcast("demo", announce()).unwrap();
2873 let mut dynamic = source.dynamic();
2874 settle().await;
2875 settle().await;
2876
2877 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2878 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2879 let mut producer = accept_track(&mut dynamic, "video").await;
2880 settle().await;
2881 let mut sub = subscribing.await.unwrap();
2882
2883 producer
2885 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2886 .unwrap();
2887
2888 let frame = sub.read_frame().await.unwrap().expect("frame");
2889 assert_eq!(frame.payload.len(), 5);
2890 settle().await;
2891
2892 let report = registry.report();
2893 let entry = report
2894 .traffic
2895 .iter()
2896 .find(|e| e.path.as_str() == "demo")
2897 .expect("demo tracked");
2898 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
2899 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
2900 assert_eq!(
2901 entry.publisher.bytes, 5,
2902 "payload counted once, not zero and not doubled"
2903 );
2904 }
2905
2906 #[tokio::test]
2910 async fn test_stats_datagrams_counted_both_sides() {
2911 use crate::Timestamp;
2912 use crate::stats::{Config, Registry, Tier};
2913
2914 tokio::time::pause();
2915
2916 let registry = Registry::new(Config::new());
2917 let ctx = registry.tier(Tier::default()).session("acme");
2918
2919 let origin = Origin::random().produce();
2920 let ingress = origin.clone().with_stats(ctx.clone());
2921 let egress = origin.consume().with_stats(ctx.clone());
2922
2923 let mut announced = egress.announced();
2924 let source = ingress.create_broadcast("demo", announce()).unwrap();
2925 let mut dynamic = source.dynamic();
2926 settle().await;
2927 settle().await;
2928
2929 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2930 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2931 let mut producer = accept_track(&mut dynamic, "video").await;
2932 settle().await;
2933 let mut sub = subscribing.await.unwrap();
2934
2935 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
2936 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
2937 assert_eq!(&datagram.payload[..], b"hello");
2938 settle().await;
2939
2940 let report = registry.report();
2941 let entry = report
2942 .traffic
2943 .iter()
2944 .find(|e| e.path.as_str() == "demo")
2945 .expect("demo tracked");
2946
2947 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
2948 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
2949 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
2950 assert_eq!(traffic.frames, 1, "{side}: one frame");
2951 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
2952 }
2953 }
2954
2955 #[test]
2956 fn origin_rejects_reserved_ids() {
2957 assert!(Origin::new(0).is_err());
2958 assert!(Origin::new(1u64 << 62).is_err());
2959 assert_eq!(Origin::new(1).unwrap().id(), 1);
2960
2961 let mut zero = [0u8].as_slice();
2962 assert_eq!(
2963 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
2964 Origin::UNKNOWN
2965 );
2966 }
2967
2968 #[test]
2969 fn origin_list_push_fails_at_limit() {
2970 let mut list = OriginList::new();
2971 for _ in 0..MAX_HOPS {
2972 list.push(Origin::random()).unwrap();
2973 }
2974 assert_eq!(list.len(), MAX_HOPS);
2975 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
2976 }
2977
2978 #[test]
2979 fn origin_list_replace_first() {
2980 let mut list = OriginList::new();
2981 for _ in 0..3 {
2982 list.push(Origin::UNKNOWN).unwrap();
2983 }
2984
2985 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
2987 assert_eq!(
2988 list.as_slice(),
2989 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
2990 );
2991
2992 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
2994 assert_eq!(list.len(), 3);
2995 }
2996
2997 #[test]
2998 fn origin_list_try_from_vec_enforces_limit() {
2999 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3000 assert!(OriginList::try_from(under).is_ok());
3001
3002 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3003 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3004 }
3005
3006 #[tokio::test]
3007 async fn test_announce() {
3008 tokio::time::pause();
3009
3010 let origin = Origin::random().produce();
3011
3012 let mut consumer1 = origin.consume().announced();
3013 consumer1.assert_next_wait();
3014
3015 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3017 settle().await;
3018
3019 consumer1.assert_next_some("test1");
3020 consumer1.assert_next_wait();
3021
3022 let mut consumer2 = origin.consume().announced();
3025
3026 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3028 settle().await;
3029
3030 consumer1.assert_next_some("test2");
3031 consumer1.assert_next_wait();
3032
3033 consumer2.assert_next_some("test1");
3034 consumer2.assert_next_some("test2");
3035 consumer2.assert_next_wait();
3036
3037 broadcast1.finish();
3039 settle().await;
3040
3041 consumer1.assert_next_none("test1");
3043 consumer2.assert_next_none("test1");
3044 consumer1.assert_next_wait();
3045 consumer2.assert_next_wait();
3046
3047 let mut consumer3 = origin.consume().announced();
3049 consumer3.assert_next_some("test2");
3050 consumer3.assert_next_wait();
3051
3052 broadcast2.finish();
3053 settle().await;
3054
3055 consumer1.assert_next_none("test2");
3056 consumer2.assert_next_none("test2");
3057 consumer3.assert_next_none("test2");
3058 }
3059
3060 #[tokio::test]
3064 async fn test_duplicate() {
3065 tokio::time::pause();
3066
3067 let origin = Origin::random().produce();
3068 let consumer = origin.consume();
3069 let mut announced = consumer.announced();
3070
3071 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3072 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3073 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3074 settle().await;
3075 assert!(consumer.get_broadcast("test").is_some());
3076
3077 announced.assert_next_some("test");
3078 announced.assert_next_wait();
3079
3080 broadcast2.finish();
3082 settle().await;
3083 assert!(consumer.get_broadcast("test").is_some());
3084 announced.assert_next_wait();
3085
3086 broadcast1.finish();
3088 settle().await;
3089 assert!(consumer.get_broadcast("test").is_some());
3090 announced.assert_next_wait();
3091
3092 broadcast3.finish();
3094 settle().await;
3095 assert!(consumer.get_broadcast("test").is_none());
3096
3097 announced.assert_next_none("test");
3098 announced.assert_next_wait();
3099 }
3100
3101 #[tokio::test]
3104 async fn test_route_failover() {
3105 tokio::time::pause();
3106
3107 let origin = Origin::random().produce();
3108 let consumer = origin.consume();
3109 let mut announced = consumer.announced();
3110
3111 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3112 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3113
3114 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3116 let mut dynamic_a = source_a.dynamic();
3117 settle().await;
3118 settle().await;
3119 let broadcast = consumer.request_broadcast("test").await.unwrap();
3120 announced.assert_next_some("test");
3121
3122 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3124 let mut dynamic_b = source_b.dynamic();
3125 settle().await;
3126 settle().await;
3127 announced.assert_next_wait();
3128
3129 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3131 let mut producer = accept_track(&mut dynamic_a, "video").await;
3132 settle().await;
3133 dynamic_b.assert_no_request();
3134
3135 let mut sub = subscribing.await.unwrap();
3136 sub.assert_no_group();
3139 assert_eq!(producer.subscription().unwrap().group_start, None);
3140
3141 producer.append_group().unwrap();
3142 producer.append_group().unwrap();
3143 assert_eq!(sub.assert_group().sequence, 0);
3144 assert_eq!(sub.assert_group().sequence, 1);
3145
3146 producer.abort(Error::Dropped).unwrap();
3150 source_a.abort(Error::Dropped).unwrap();
3151 drop(dynamic_a);
3152 settle().await;
3153 announced.assert_next_wait();
3154
3155 let mut producer = accept_track(&mut dynamic_b, "video").await;
3158 settle().await;
3159 sub.assert_no_group();
3160 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3161 producer.create_group(group::Info { sequence: 1 }).unwrap();
3162 producer.create_group(group::Info { sequence: 2 }).unwrap();
3163 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3164 sub.assert_not_closed();
3165 }
3166
3167 #[tokio::test]
3170 async fn test_broadcast_route_watch() {
3171 let mut producer = broadcast::Info::new().produce();
3172 let mut consumer = producer.consume();
3173
3174 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3176
3177 producer.set_route(broadcast::Route::default()).unwrap();
3179 assert!(consumer.route_changed().now_or_never().is_none());
3180
3181 let mut hops = OriginList::new();
3182 hops.push(Origin::new(7).unwrap()).unwrap();
3183 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3184 producer.set_route(route.clone()).unwrap();
3185 assert_eq!(consumer.route_changed().await.unwrap(), route);
3186
3187 let mut fresh = producer.consume();
3189 assert_eq!(fresh.route_changed().await.unwrap(), route);
3190
3191 drop(producer);
3192 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3193 }
3194
3195 #[tokio::test]
3199 async fn test_route_cost_update() {
3200 tokio::time::pause();
3201
3202 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3206 let consumer = origin.consume();
3207 let mut announced = consumer.announced();
3208
3209 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3210 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3211
3212 let mut source_a = origin
3214 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3215 .unwrap();
3216 let mut dynamic_a = source_a.dynamic();
3217 settle().await;
3218 let broadcast = consumer.request_broadcast("test").await.unwrap();
3219 announced.assert_next_some("test");
3220
3221 let mut watch = broadcast.clone();
3222 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3223
3224 let mut source_b = origin
3225 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3226 .unwrap();
3227 let mut dynamic_b = source_b.dynamic();
3228 settle().await;
3229 assert!(
3230 watch.route_changed().now_or_never().is_none(),
3231 "a losing standby must not change the advertised route"
3232 );
3233
3234 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3236 let mut producer = accept_track(&mut dynamic_a, "video").await;
3237 settle().await;
3238 let mut sub = subscribing.await.unwrap();
3239 producer.append_group().unwrap();
3240 assert_eq!(sub.assert_group().sequence, 0);
3241
3242 source_a
3245 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3246 .unwrap();
3247 settle().await;
3248 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3249 announced.assert_next_wait();
3250
3251 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3252 settle().await;
3253 sub.assert_no_group();
3256 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3257 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3258 assert_eq!(sub.assert_group().sequence, 1);
3259 sub.assert_not_closed();
3260
3261 source_b
3263 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3264 .unwrap();
3265 settle().await;
3266 let advertised = watch.route_changed().await.unwrap();
3267 assert_eq!(advertised.hops, hops_b);
3268 assert_eq!(advertised.cost, 5);
3269 announced.assert_next_wait();
3270 }
3271
3272 #[tokio::test]
3275 async fn test_completed_track_survives_route_churn() {
3276 tokio::time::pause();
3277
3278 let origin = Origin::random().produce();
3279 let consumer = origin.consume();
3280
3281 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3282 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3283
3284 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3285 let mut dynamic_a = source_a.dynamic();
3286 settle().await;
3287 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3288 let mut dynamic_b = source_b.dynamic();
3289 settle().await;
3290 settle().await;
3291 let broadcast = consumer.request_broadcast("test").await.unwrap();
3292
3293 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3295 let mut producer = accept_track(&mut dynamic_a, "video").await;
3296 settle().await;
3297 let mut sub = subscribing.await.unwrap();
3298 producer.append_group().unwrap();
3299 assert_eq!(sub.assert_group().sequence, 0);
3300 producer.finish().unwrap();
3301 drop(producer);
3302 settle().await;
3303 sub.assert_closed();
3304
3305 source_a.abort(Error::Dropped).unwrap();
3307 drop(dynamic_a);
3308 settle().await;
3309 dynamic_b.assert_no_request();
3310
3311 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3313 late.assert_closed();
3314 }
3315
3316 #[tokio::test]
3319 async fn test_serve_resets_retry_budget() {
3320 tokio::time::pause();
3321
3322 let origin = Origin::random().produce();
3323 let consumer = origin.consume();
3324
3325 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3326 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3327 let mut dynamic = source.dynamic();
3328 settle().await;
3329 settle().await;
3330 let broadcast = consumer.request_broadcast("test").await.unwrap();
3331
3332 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3334
3335 for _ in 0..2 * MAX_TRACK_RETRIES {
3338 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3339 .await
3340 .expect("timed out waiting for a retry")
3341 .unwrap();
3342 request.reject(Error::NotFound);
3343 let producer = accept_track(&mut dynamic, "video").await;
3344 settle().await;
3345 drop(producer);
3346 }
3347
3348 let _producer = accept_track(&mut dynamic, "video").await;
3349 settle().await;
3350 let mut sub = subscribing.await.unwrap();
3351 sub.assert_not_closed();
3352 }
3353
3354 #[tokio::test]
3358 async fn test_route_handover() {
3359 tokio::time::pause();
3360
3361 let origin = Origin::random().produce();
3362 let consumer = origin.consume();
3363 let mut announced = consumer.announced();
3364
3365 let hops_long = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3366 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3367
3368 let source_a = origin
3369 .create_broadcast("test", announce().with_hops(hops_long))
3370 .unwrap();
3371 let mut dynamic_a = source_a.dynamic();
3372 settle().await;
3373 settle().await;
3374 let broadcast = consumer.request_broadcast("test").await.unwrap();
3375 announced.assert_next_some("test");
3376
3377 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3378 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3379 settle().await;
3380 let mut sub = subscribing.await.unwrap();
3381 producer_a.append_group().unwrap();
3382 producer_a.append_group().unwrap();
3383 assert_eq!(sub.assert_group().sequence, 0);
3384 assert_eq!(sub.assert_group().sequence, 1);
3385
3386 let source_b = origin
3389 .create_broadcast("test", announce().with_hops(hops_short))
3390 .unwrap();
3391 let mut dynamic_b = source_b.dynamic();
3392 settle().await;
3393 settle().await;
3394 announced.assert_next_wait();
3395
3396 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3397 settle().await;
3398
3399 sub.assert_no_group();
3402 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3403 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3404
3405 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3407 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3408 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3409 assert_eq!(sub.assert_group().sequence, 2);
3410 assert_eq!(sub.assert_group().sequence, 3);
3411 sub.assert_no_group();
3412 sub.assert_not_closed();
3413 }
3414
3415 #[tokio::test(start_paused = true)]
3418 async fn test_route_unannounce_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 mut source = origin
3425 .create_broadcast("test", announce().with_hops(hops.clone()))
3426 .unwrap();
3427 settle().await;
3428 let broadcast = consumer.request_broadcast("test").await.unwrap();
3429 announced.assert_next_some("test");
3430
3431 source.finish();
3434 settle().await;
3435 announced.assert_next_none("test");
3436
3437 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3439 settle().await;
3440 let fresh = consumer.request_broadcast("test").await.unwrap();
3441 announced.assert_next_some("test");
3442 assert!(
3443 !fresh.is_clone(&broadcast),
3444 "re-create must not splice the old broadcast"
3445 );
3446 }
3447
3448 #[tokio::test(start_paused = true)]
3453 async fn test_route_detach_immediate() {
3454 let origin = Origin::random().produce();
3455 let consumer = origin.consume();
3456 let mut announced = consumer.announced();
3457
3458 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3459 let source = origin
3460 .create_broadcast("test", announce().with_hops(hops.clone()))
3461 .unwrap();
3462 let mut dynamic = source.dynamic();
3463 settle().await;
3464 settle().await;
3465 let broadcast = consumer.request_broadcast("test").await.unwrap();
3466 announced.assert_next_some("test");
3467
3468 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3469 let producer = accept_track(&mut dynamic, "video").await;
3470 settle().await;
3471 let mut sub = subscribing.await.unwrap();
3472
3473 drop(producer);
3475 source.abort(Error::Dropped).unwrap();
3476 drop(dynamic);
3477
3478 settle().await;
3479 announced.assert_next_none("test");
3480 sub.assert_error();
3481
3482 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3485 settle().await;
3486 settle().await;
3487 let fresh = consumer.request_broadcast("test").await.unwrap();
3488 announced.assert_next_some("test");
3489 assert!(
3490 !fresh.is_clone(&broadcast),
3491 "re-create must not splice the old broadcast"
3492 );
3493 }
3494
3495 #[tokio::test(start_paused = true)]
3500 async fn test_idle_track_releases_without_respinning() {
3501 let origin = Info::new(Origin::random()).produce();
3502 let consumer = origin.consume();
3503
3504 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3505 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3506 let mut dynamic = source.dynamic();
3507 settle().await;
3508 let broadcast = consumer.request_broadcast("test").await.unwrap();
3509
3510 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3511 let producer = accept_track(&mut dynamic, "video").await;
3512 settle().await;
3513 let sub = subscribing.await.unwrap();
3514
3515 drop(sub);
3518 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
3519 settle().await;
3520 assert!(
3521 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
3522 "the copy must stay spliced inside the linger",
3523 );
3524
3525 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3528 settle().await;
3529 assert!(
3530 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3531 "an idle copy must be released after the linger",
3532 );
3533
3534 for _ in 0..3 {
3538 tokio::time::sleep(TRACK_IDLE_LINGER).await;
3539 settle().await;
3540 assert!(
3541 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3542 "an unread copy must stay released, not be re-spliced",
3543 );
3544 }
3545 assert!(
3546 dynamic.requested_track().now_or_never().is_none(),
3547 "an unread track must not be re-requested",
3548 );
3549 drop(producer);
3550
3551 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3553 let mut producer = accept_track(&mut dynamic, "video").await;
3554 settle().await;
3555 let mut sub = subscribing.await.unwrap();
3556 producer.append_group().unwrap();
3557 assert_eq!(sub.assert_group().sequence, 0);
3558 }
3559
3560 #[tokio::test(start_paused = true)]
3564 async fn test_back_to_back_fetches_reuse_the_track() {
3565 let origin = Info::new(Origin::random()).produce();
3566 let consumer = origin.consume();
3567
3568 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3569 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3570 let mut dynamic = source.dynamic();
3571 settle().await;
3572 let broadcast = consumer.request_broadcast("test").await.unwrap();
3573
3574 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3576 let mut producer = accept_track(&mut dynamic, "video").await;
3577 producer.append_group().unwrap().finish().unwrap();
3578 settle().await;
3579 let first = fetching.await.expect("first fetch");
3580 drop(first);
3581
3582 settle().await;
3584 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3585 settle().await;
3586 assert!(
3587 dynamic.requested_track().now_or_never().is_none(),
3588 "a fetch inside the linger must reuse the track, not re-request it",
3589 );
3590 drop(fetching.await.expect("second fetch"));
3591
3592 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
3594 settle().await;
3595 assert!(
3596 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
3597 "the copy must be released once the fetches stop",
3598 );
3599 drop(producer);
3600
3601 settle().await;
3603 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
3604 let mut producer = accept_track(&mut dynamic, "video").await;
3605 producer.append_group().unwrap().finish().unwrap();
3606 settle().await;
3607 fetching.await.expect("fetch after the linger");
3608 }
3609
3610 #[tokio::test(start_paused = true)]
3614 async fn test_linger_reconnect_splices() {
3615 let origin = Info::new(Origin::random())
3616 .with_linger(Duration::from_secs(5))
3617 .produce();
3618 let consumer = origin.consume();
3619 let mut announced = consumer.announced();
3620
3621 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3622 let source = origin
3623 .create_broadcast("test", announce().with_hops(hops.clone()))
3624 .unwrap();
3625 let mut dynamic = source.dynamic();
3626 settle().await;
3627 settle().await;
3628 let broadcast = consumer.request_broadcast("test").await.unwrap();
3629 announced.assert_next_some("test");
3630
3631 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3632 let mut producer = accept_track(&mut dynamic, "video").await;
3633 settle().await;
3634 let mut sub = subscribing.await.unwrap();
3635
3636 producer.append_group().unwrap();
3637 producer.append_group().unwrap();
3638 assert_eq!(sub.assert_group().sequence, 0);
3639 assert_eq!(sub.assert_group().sequence, 1);
3640
3641 drop(producer);
3644 source.abort(Error::Dropped).unwrap();
3645 drop(dynamic);
3646 settle().await;
3647
3648 announced.assert_next_wait();
3650 sub.assert_no_group();
3651 sub.assert_not_closed();
3652
3653 let during = consumer.request_broadcast("test").await.unwrap();
3655 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
3656
3657 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3660 let mut dynamic = source.dynamic();
3661 settle().await;
3662 settle().await;
3663 announced.assert_next_wait();
3664 let again = consumer.request_broadcast("test").await.unwrap();
3665 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3666
3667 let mut producer = accept_track(&mut dynamic, "video").await;
3671 settle().await;
3672 sub.assert_no_group();
3673 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3674 producer.create_group(group::Info { sequence: 2 }).unwrap();
3675 assert_eq!(sub.assert_group().sequence, 2);
3676 sub.assert_not_closed();
3677 }
3678
3679 #[tokio::test(start_paused = true)]
3682 async fn test_linger_expiry_closes() {
3683 let origin = Info::new(Origin::random())
3684 .with_linger(Duration::from_secs(5))
3685 .produce();
3686 let consumer = origin.consume();
3687 let mut announced = consumer.announced();
3688
3689 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3690 let source = origin
3691 .create_broadcast("test", announce().with_hops(hops.clone()))
3692 .unwrap();
3693 let mut dynamic = source.dynamic();
3694 settle().await;
3695 settle().await;
3696 let broadcast = consumer.request_broadcast("test").await.unwrap();
3697 announced.assert_next_some("test");
3698
3699 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3700 let producer = accept_track(&mut dynamic, "video").await;
3701 settle().await;
3702 let mut sub = subscribing.await.unwrap();
3703
3704 drop(producer);
3705 source.abort(Error::Dropped).unwrap();
3706 drop(dynamic);
3707 settle().await;
3708 announced.assert_next_wait();
3709
3710 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
3712 settle().await;
3713 announced.assert_next_none("test");
3714 sub.assert_error();
3715
3716 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3718 settle().await;
3719 settle().await;
3720 let fresh = consumer.request_broadcast("test").await.unwrap();
3721 announced.assert_next_some("test");
3722 assert!(
3723 !fresh.is_clone(&broadcast),
3724 "a late re-create must not splice the expired broadcast"
3725 );
3726 }
3727
3728 #[tokio::test(start_paused = true)]
3732 async fn test_linger_forever() {
3733 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
3734 let consumer = origin.consume();
3735 let mut announced = consumer.announced();
3736
3737 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3738 let source = origin
3739 .create_broadcast("test", announce().with_hops(hops.clone()))
3740 .unwrap();
3741 settle().await;
3742 let broadcast = consumer.request_broadcast("test").await.unwrap();
3743 announced.assert_next_some("test");
3744
3745 source.abort(Error::Dropped).unwrap();
3746 settle().await;
3747
3748 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
3750 announced.assert_next_wait();
3751 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3752 settle().await;
3753 settle().await;
3754 let again = consumer.request_broadcast("test").await.unwrap();
3755 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3756 drop(source);
3757 }
3758
3759 #[tokio::test(start_paused = true)]
3762 async fn test_linger_skipped_on_finish() {
3763 let origin = Info::new(Origin::random())
3764 .with_linger(Duration::from_secs(5))
3765 .produce();
3766 let consumer = origin.consume();
3767 let mut announced = consumer.announced();
3768
3769 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3770 let mut source = origin
3771 .create_broadcast("test", announce().with_hops(hops.clone()))
3772 .unwrap();
3773 settle().await;
3774 let broadcast = consumer.request_broadcast("test").await.unwrap();
3775 announced.assert_next_some("test");
3776
3777 source.finish();
3780 settle().await;
3781 announced.assert_next_none("test");
3782
3783 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3785 settle().await;
3786 let fresh = consumer.request_broadcast("test").await.unwrap();
3787 announced.assert_next_some("test");
3788 assert!(
3789 !fresh.is_clone(&broadcast),
3790 "a finish must not leave a lingering broadcast to splice into"
3791 );
3792 }
3793
3794 #[tokio::test]
3797 async fn test_announce_toggle() {
3798 tokio::time::pause();
3799
3800 let origin = Origin::random().produce();
3801 let consumer = origin.consume();
3802 let mut announced = consumer.announced();
3803
3804 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3805 settle().await;
3806
3807 announced.assert_next_wait();
3809 let broadcast = consumer
3810 .get_broadcast("test")
3811 .expect("offline broadcast is still routable");
3812 assert!(!broadcast.route().announce);
3813
3814 let requested = consumer.request_broadcast("test").await.unwrap();
3816 assert!(requested.is_clone(&broadcast));
3817
3818 source.set_route(announce()).unwrap();
3820 settle().await;
3821 let face = announced.assert_next_some("test");
3822 assert!(face.is_clone(&broadcast));
3823
3824 let mut fresh = origin.consume().announced();
3826 fresh.assert_next_some("test");
3827 fresh.assert_next_wait();
3828
3829 source.set_route(broadcast::Route::new()).unwrap();
3831 settle().await;
3832 announced.assert_next_none("test");
3833 assert!(consumer.get_broadcast("test").is_some());
3834 let mut fresh = origin.consume().announced();
3835 fresh.assert_next_wait();
3836
3837 source.finish();
3838 settle().await;
3839 assert!(consumer.get_broadcast("test").is_none());
3840 }
3841
3842 #[tokio::test]
3845 async fn test_announce_beats_offline() {
3846 tokio::time::pause();
3847
3848 let origin = Origin::random().produce();
3849 let consumer = origin.consume();
3850 let mut announced = consumer.announced();
3851
3852 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3854 settle().await;
3855 announced.assert_next_wait();
3856
3857 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
3860 settle().await;
3861 announced.assert_next_some("test");
3862 let face = consumer.get_broadcast("test").unwrap();
3863 assert!(face.route().announce);
3864 assert_eq!(face.route().cost, 10);
3865
3866 announced_source.finish();
3869 settle().await;
3870 announced.assert_next_none("test");
3871 assert!(consumer.get_broadcast("test").is_some());
3872 }
3873
3874 #[tokio::test]
3877 async fn test_better_source_no_churn() {
3878 tokio::time::pause();
3879
3880 let origin = Origin::random().produce();
3881 let mut announced = origin.consume().announced();
3882
3883 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3885 let _a = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3886 settle().await;
3887 let face = announced.assert_next_some("test");
3888
3889 let _b = origin.create_broadcast("test", announce()).unwrap();
3890 settle().await;
3891 announced.assert_next_wait();
3892 let current = origin.consume().get_broadcast("test").unwrap();
3893 assert!(current.is_clone(&face), "the broadcast identity must not change");
3894 assert!(current.route().hops.is_empty());
3896 }
3897
3898 #[tokio::test]
3899 async fn test_duplicate_reverse() {
3900 tokio::time::pause();
3901
3902 let origin = Origin::random().produce();
3903
3904 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3905 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3906 settle().await;
3907 assert!(origin.consume().get_broadcast("test").is_some());
3908
3909 broadcast2.finish();
3911 settle().await;
3912 assert!(origin.consume().get_broadcast("test").is_some());
3913
3914 broadcast1.finish();
3915 settle().await;
3916 assert!(origin.consume().get_broadcast("test").is_none());
3917 }
3918
3919 #[tokio::test]
3920 async fn test_deterministic_tiebreak() {
3921 tokio::time::pause();
3922
3923 fn hops(ids: &[u64]) -> OriginList {
3924 OriginList::try_from(
3925 ids.iter()
3926 .copied()
3927 .map(|id| Origin::new(id).unwrap())
3928 .collect::<Vec<_>>(),
3929 )
3930 .unwrap()
3931 }
3932
3933 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
3936 let origin = Origin::random().produce();
3937 let _a = origin
3938 .create_broadcast("test", announce().with_hops(hops(first)))
3939 .unwrap();
3940 let _b = origin
3941 .create_broadcast("test", announce().with_hops(hops(second)))
3942 .unwrap();
3943 settle().await;
3944 origin.consume().get_broadcast("test").unwrap().route().hops
3945 }
3946
3947 let forward = winner(&[10, 20], &[30, 40]).await;
3950 let reverse = winner(&[30, 40], &[10, 20]).await;
3951 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
3952
3953 assert_eq!(winner(&[10, 20], &[30]).await.len(), 1);
3955 assert_eq!(winner(&[30], &[10, 20]).await.len(), 1);
3956 }
3957
3958 #[tokio::test]
3963 async fn test_many_announces() {
3964 let origin = Origin::random().produce();
3965
3966 let mut consumer = origin.consume().announced();
3967 let mut broadcasts = Vec::new();
3969 for i in 0..256 {
3970 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3971 settle().await;
3972 }
3973
3974 for i in 0..256 {
3975 consumer.assert_next_some(format!("test{i:03}"));
3976 }
3977 consumer.assert_next_wait();
3978 }
3979
3980 #[tokio::test]
3981 async fn test_many_announces_try() {
3982 let origin = Origin::random().produce();
3983
3984 let mut consumer = origin.consume().announced();
3985 let mut broadcasts = Vec::new();
3987 for i in 0..256 {
3988 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3989 settle().await;
3990 }
3991
3992 for i in 0..256 {
3993 consumer.assert_try_next_some(format!("test{i:03}"));
3994 }
3995 }
3996
3997 #[tokio::test]
3998 async fn test_with_root_basic() {
3999 let origin = Origin::random().produce();
4000
4001 let foo_producer = origin.with_root("foo").expect("should create root");
4003 assert_eq!(foo_producer.root().as_str(), "foo");
4004
4005 let mut consumer = origin.consume().announced();
4006
4007 let _broadcast = foo_producer
4009 .create_broadcast("bar/baz", announce())
4010 .expect("publish allowed");
4011 settle().await;
4012 consumer.assert_next_some("foo/bar/baz");
4014
4015 let mut foo_consumer = foo_producer.consume().announced();
4017 foo_consumer.assert_next_some("bar/baz");
4018 }
4019
4020 #[tokio::test]
4021 async fn test_with_root_nested() {
4022 let origin = Origin::random().produce();
4023
4024 let foo_producer = origin.with_root("foo").expect("should create foo root");
4026 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
4027 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
4028
4029 let mut consumer = origin.consume().announced();
4030
4031 let _broadcast = foo_bar_producer
4033 .create_broadcast("baz", announce())
4034 .expect("publish allowed");
4035 settle().await;
4036 consumer.assert_next_some("foo/bar/baz");
4038
4039 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
4041 foo_bar_consumer.assert_next_some("baz");
4042 }
4043
4044 #[tokio::test]
4045 async fn test_publish_scope_allows() {
4046 let origin = Origin::random().produce();
4047
4048 let limited_producer = origin
4050 .scope(&["allowed/path1".into(), "allowed/path2".into()])
4051 .expect("should create limited producer");
4052
4053 let _broadcast = limited_producer
4055 .create_broadcast("allowed/path1", announce())
4056 .expect("publish allowed");
4057 let _keep2 = limited_producer
4058 .create_broadcast("allowed/path1/nested", announce())
4059 .expect("publish allowed");
4060 let _keep3 = limited_producer
4061 .create_broadcast("allowed/path2", announce())
4062 .expect("publish allowed");
4063 settle().await;
4064
4065 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
4067 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
4069 }
4070
4071 #[tokio::test]
4072 async fn test_publish_max_parts() {
4073 let origin = Origin::random().produce();
4074
4075 let at_limit = (0..Path::MAX_PARTS)
4076 .map(|i| i.to_string())
4077 .collect::<Vec<_>>()
4078 .join("/");
4079 let _broadcast = origin
4080 .create_broadcast(at_limit.as_str(), announce())
4081 .expect("publish allowed");
4082 settle().await;
4083
4084 let too_deep = format!("{at_limit}/extra");
4085 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
4086
4087 let rooted = origin.with_root("root").expect("wildcard allows any root");
4089 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
4090 }
4091
4092 #[tokio::test]
4093 async fn test_publish_scope_empty() {
4094 let origin = Origin::random().produce();
4095
4096 assert!(origin.scope(&[]).is_none());
4098 }
4099
4100 #[tokio::test]
4101 async fn test_consume_scope_filters() {
4102 let origin = Origin::random().produce();
4103
4104 let mut consumer = origin.consume().announced();
4105
4106 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
4108 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
4109 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
4110 settle().await;
4111
4112 let mut limited_consumer = origin
4114 .consume()
4115 .scope(&["allowed".into()])
4116 .expect("should create limited consumer")
4117 .announced();
4118
4119 limited_consumer.assert_next_some("allowed");
4121 limited_consumer.assert_next_some("allowed/nested");
4122 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
4126 consumer.assert_next_some("allowed/nested");
4127 consumer.assert_next_some("notallowed");
4128 }
4129
4130 #[tokio::test]
4131 async fn test_consume_scope_multiple_prefixes() {
4132 let origin = Origin::random().produce();
4133
4134 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
4135 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
4136 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
4137 settle().await;
4138
4139 let mut limited_consumer = origin
4141 .consume()
4142 .scope(&["foo".into(), "bar".into()])
4143 .expect("should create limited consumer")
4144 .announced();
4145
4146 limited_consumer.assert_next_some("bar/test");
4148 limited_consumer.assert_next_some("foo/test");
4149 limited_consumer.assert_next_wait(); }
4151
4152 #[tokio::test]
4153 async fn test_with_root_and_publish_scope() {
4154 let origin = Origin::random().produce();
4155
4156 let foo_producer = origin.with_root("foo").expect("should create foo root");
4158
4159 let limited_producer = foo_producer
4161 .scope(&["bar".into(), "goop/pee".into()])
4162 .expect("should create limited producer");
4163
4164 let mut consumer = origin.consume().announced();
4165
4166 let _broadcast = limited_producer
4168 .create_broadcast("bar", announce())
4169 .expect("publish allowed");
4170 let _keep2 = limited_producer
4171 .create_broadcast("bar/nested", announce())
4172 .expect("publish allowed");
4173 let _keep3 = limited_producer
4174 .create_broadcast("goop/pee", announce())
4175 .expect("publish allowed");
4176 let _keep4 = limited_producer
4177 .create_broadcast("goop/pee/nested", announce())
4178 .expect("publish allowed");
4179 settle().await;
4180
4181 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
4183 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
4185
4186 consumer.assert_next_some("foo/bar");
4188 consumer.assert_next_some("foo/bar/nested");
4189 consumer.assert_next_some("foo/goop/pee");
4190 consumer.assert_next_some("foo/goop/pee/nested");
4191 }
4192
4193 #[tokio::test]
4194 async fn test_with_root_and_consume_scope() {
4195 let origin = Origin::random().produce();
4196
4197 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
4199 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
4200 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
4201 settle().await;
4202
4203 let foo_producer = origin.with_root("foo").expect("should create foo root");
4205
4206 let mut limited_consumer = foo_producer
4208 .consume()
4209 .scope(&["bar".into(), "goop/pee".into()])
4210 .expect("should create limited consumer")
4211 .announced();
4212
4213 limited_consumer.assert_next_some("bar/test");
4215 limited_consumer.assert_next_some("goop/pee/test");
4216 limited_consumer.assert_next_wait(); }
4218
4219 #[tokio::test]
4220 async fn test_with_root_unauthorized() {
4221 let origin = Origin::random().produce();
4222
4223 let limited_producer = origin
4225 .scope(&["allowed".into()])
4226 .expect("should create limited producer");
4227
4228 assert!(limited_producer.with_root("notallowed").is_none());
4230
4231 let allowed_root = limited_producer
4233 .with_root("allowed")
4234 .expect("should create allowed root");
4235 assert_eq!(allowed_root.root().as_str(), "allowed");
4236 }
4237
4238 #[tokio::test]
4239 async fn test_wildcard_permission() {
4240 let origin = Origin::random().produce();
4241
4242 let root_producer = origin.clone();
4244
4245 let _broadcast = root_producer
4247 .create_broadcast("any/path", announce())
4248 .expect("publish allowed");
4249 let _keep2 = root_producer
4250 .create_broadcast("other/path", announce())
4251 .expect("publish allowed");
4252 settle().await;
4253
4254 let foo_producer = root_producer.with_root("foo").expect("should create any root");
4256 assert_eq!(foo_producer.root().as_str(), "foo");
4257 }
4258
4259 #[tokio::test]
4260 async fn test_consume_broadcast_with_permissions() {
4261 let origin = Origin::random().produce();
4262
4263 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
4264 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
4265 settle().await;
4266
4267 let limited_consumer = origin
4269 .consume()
4270 .scope(&["allowed".into()])
4271 .expect("should create limited consumer");
4272
4273 let result = limited_consumer.get_broadcast("allowed/test");
4275 assert!(result.is_some());
4276 assert!(
4277 result
4278 .unwrap()
4279 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
4280 );
4281
4282 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
4284
4285 let consumer = origin.consume();
4287 assert!(consumer.get_broadcast("allowed/test").is_some());
4288 assert!(consumer.get_broadcast("notallowed/test").is_some());
4289 }
4290
4291 #[tokio::test]
4292 async fn test_nested_paths_with_permissions() {
4293 let origin = Origin::random().produce();
4294
4295 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
4297
4298 let _broadcast = limited_producer
4300 .create_broadcast("a/b/c", announce())
4301 .expect("publish allowed");
4302 let _keep2 = limited_producer
4303 .create_broadcast("a/b/c/d", announce())
4304 .expect("publish allowed");
4305 let _keep3 = limited_producer
4306 .create_broadcast("a/b/c/d/e", announce())
4307 .expect("publish allowed");
4308 settle().await;
4309
4310 assert!(limited_producer.create_broadcast("a", announce()).is_err());
4312 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
4313 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
4314 }
4315
4316 #[tokio::test]
4317 async fn test_multiple_consumers_with_different_permissions() {
4318 let origin = Origin::random().produce();
4319
4320 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
4322 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
4323 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
4324 settle().await;
4325
4326 let mut foo_consumer = origin
4328 .consume()
4329 .scope(&["foo".into()])
4330 .expect("should create foo consumer")
4331 .announced();
4332
4333 let mut bar_consumer = origin
4334 .consume()
4335 .scope(&["bar".into()])
4336 .expect("should create bar consumer")
4337 .announced();
4338
4339 let mut foobar_consumer = origin
4340 .consume()
4341 .scope(&["foo".into(), "bar".into()])
4342 .expect("should create foobar consumer")
4343 .announced();
4344
4345 foo_consumer.assert_next_some("foo/test");
4347 foo_consumer.assert_next_wait();
4348
4349 bar_consumer.assert_next_some("bar/test");
4350 bar_consumer.assert_next_wait();
4351
4352 foobar_consumer.assert_next_some("bar/test");
4353 foobar_consumer.assert_next_some("foo/test");
4354 foobar_consumer.assert_next_wait();
4355 }
4356
4357 #[tokio::test]
4358 async fn test_select_with_empty_prefix() {
4359 let origin = Origin::random().produce();
4360
4361 let demo_producer = origin.with_root("demo").expect("should create demo root");
4363 let limited_producer = demo_producer
4364 .scope(&["worm-node".into(), "foobar".into()])
4365 .expect("should create limited producer");
4366
4367 let _broadcast1 = limited_producer
4369 .create_broadcast("worm-node/test", announce())
4370 .expect("publish allowed");
4371 let _broadcast2 = limited_producer
4372 .create_broadcast("foobar/test", announce())
4373 .expect("publish allowed");
4374 settle().await;
4375
4376 let mut consumer = limited_producer
4378 .consume()
4379 .scope(&["".into()])
4380 .expect("should create consumer with empty prefix")
4381 .announced();
4382
4383 let a1 = consumer.try_next().expect("expected first announcement");
4385 let a2 = consumer.try_next().expect("expected second announcement");
4386 consumer.assert_next_wait();
4387
4388 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4389 paths.sort();
4390 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
4391 }
4392
4393 #[tokio::test]
4394 async fn test_select_narrowing_scope() {
4395 let origin = Origin::random().produce();
4396
4397 let demo_producer = origin.with_root("demo").expect("should create demo root");
4399 let limited_producer = demo_producer
4400 .scope(&["worm-node".into(), "foobar".into()])
4401 .expect("should create limited producer");
4402
4403 let _broadcast1 = limited_producer
4405 .create_broadcast("worm-node", announce())
4406 .expect("publish allowed");
4407 let _broadcast2 = limited_producer
4408 .create_broadcast("worm-node/foo", announce())
4409 .expect("publish allowed");
4410 let _broadcast3 = limited_producer
4411 .create_broadcast("foobar/bar", announce())
4412 .expect("publish allowed");
4413 settle().await;
4414
4415 let mut worm_consumer = limited_producer
4417 .consume()
4418 .scope(&["worm-node".into()])
4419 .expect("should create worm-node consumer")
4420 .announced();
4421
4422 worm_consumer.assert_next_some("worm-node");
4424 worm_consumer.assert_next_some("worm-node/foo");
4425 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
4429 .consume()
4430 .scope(&["worm-node/foo".into()])
4431 .expect("should create worm-node/foo consumer")
4432 .announced();
4433
4434 foo_consumer.assert_next_some("worm-node/foo");
4435 foo_consumer.assert_next_wait(); }
4437
4438 #[tokio::test]
4439 async fn test_select_multiple_roots_with_empty_prefix() {
4440 let origin = Origin::random().produce();
4441
4442 let limited_producer = origin
4444 .scope(&["app1".into(), "app2".into(), "shared".into()])
4445 .expect("should create limited producer");
4446
4447 let _broadcast1 = limited_producer
4449 .create_broadcast("app1/data", announce())
4450 .expect("publish allowed");
4451 let _broadcast2 = limited_producer
4452 .create_broadcast("app2/config", announce())
4453 .expect("publish allowed");
4454 let _broadcast3 = limited_producer
4455 .create_broadcast("shared/resource", announce())
4456 .expect("publish allowed");
4457 settle().await;
4458
4459 let mut consumer = limited_producer
4461 .consume()
4462 .scope(&["".into()])
4463 .expect("should create consumer with empty prefix")
4464 .announced();
4465
4466 consumer.assert_next_some("app1/data");
4468 consumer.assert_next_some("app2/config");
4469 consumer.assert_next_some("shared/resource");
4470 consumer.assert_next_wait();
4471 }
4472
4473 #[tokio::test]
4474 async fn test_publish_scope_with_empty_prefix() {
4475 let origin = Origin::random().produce();
4476
4477 let limited_producer = origin
4479 .scope(&["services/api".into(), "services/web".into()])
4480 .expect("should create limited producer");
4481
4482 let same_producer = limited_producer
4484 .scope(&["".into()])
4485 .expect("should create producer with empty prefix");
4486
4487 let _broadcast = same_producer
4489 .create_broadcast("services/api", announce())
4490 .expect("publish allowed");
4491 let _keep2 = same_producer
4492 .create_broadcast("services/web", announce())
4493 .expect("publish allowed");
4494 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
4495 assert!(same_producer.create_broadcast("other", announce()).is_err());
4496 }
4497
4498 #[tokio::test]
4499 async fn test_select_narrowing_to_deeper_path() {
4500 let origin = Origin::random().produce();
4501
4502 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
4504
4505 let _broadcast1 = limited_producer
4507 .create_broadcast("org/team1/project1", announce())
4508 .expect("publish allowed");
4509 let _broadcast2 = limited_producer
4510 .create_broadcast("org/team1/project2", announce())
4511 .expect("publish allowed");
4512 let _broadcast3 = limited_producer
4513 .create_broadcast("org/team2/project1", announce())
4514 .expect("publish allowed");
4515 settle().await;
4516
4517 let mut team2_consumer = limited_producer
4519 .consume()
4520 .scope(&["org/team2".into()])
4521 .expect("should create team2 consumer")
4522 .announced();
4523
4524 team2_consumer.assert_next_some("org/team2/project1");
4525 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
4529 .consume()
4530 .scope(&["org/team1/project1".into()])
4531 .expect("should create project1 consumer")
4532 .announced();
4533
4534 project1_consumer.assert_next_some("org/team1/project1");
4536 project1_consumer.assert_next_wait();
4537 }
4538
4539 #[tokio::test]
4540 async fn test_select_with_non_matching_prefix() {
4541 let origin = Origin::random().produce();
4542
4543 let limited_producer = origin
4545 .scope(&["allowed/path".into()])
4546 .expect("should create limited producer");
4547
4548 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
4550
4551 assert!(limited_producer.scope(&["other/path".into()]).is_none());
4553 }
4554
4555 #[tokio::test]
4558 async fn test_with_root_trailing_slash_consumer() {
4559 let origin = Origin::random().produce();
4560
4561 let prefix = "some_prefix/".to_string();
4563 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4564
4565 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4566 settle().await;
4567 consumer.assert_next_some("test");
4568 }
4569
4570 #[tokio::test]
4572 async fn test_with_root_trailing_slash_producer() {
4573 let origin = Origin::random().produce();
4574
4575 let prefix = "some_prefix/".to_string();
4577 let rooted = origin.with_root(prefix).unwrap();
4578
4579 let _b = rooted.create_broadcast("test", announce()).unwrap();
4580 settle().await;
4581
4582 let mut consumer = rooted.consume().announced();
4583 consumer.assert_next_some("test");
4584 }
4585
4586 #[tokio::test]
4588 async fn test_with_root_trailing_slash_unannounce() {
4589 tokio::time::pause();
4590
4591 let origin = Origin::random().produce();
4592
4593 let prefix = "some_prefix/".to_string();
4594 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4595
4596 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4597 settle().await;
4598 consumer.assert_next_some("test");
4599
4600 b.finish();
4602 settle().await;
4603
4604 consumer.assert_next_none("test");
4606 }
4607
4608 #[tokio::test]
4609 async fn test_select_maintains_access_with_wider_prefix() {
4610 let origin = Origin::random().produce();
4611
4612 let demo_producer = origin.with_root("demo").expect("should create demo root");
4614 let user_producer = demo_producer
4615 .scope(&["worm-node".into(), "foobar".into()])
4616 .expect("should create user producer");
4617
4618 let _broadcast1 = user_producer
4620 .create_broadcast("worm-node/data", announce())
4621 .expect("publish allowed");
4622 let _broadcast2 = user_producer
4623 .create_broadcast("foobar", announce())
4624 .expect("publish allowed");
4625 settle().await;
4626
4627 let mut consumer = user_producer
4629 .consume()
4630 .scope(&["".into()])
4631 .expect("scope with empty prefix should not fail when user has specific permissions")
4632 .announced();
4633
4634 let a1 = consumer.try_next().expect("expected first announcement");
4636 let a2 = consumer.try_next().expect("expected second announcement");
4637 consumer.assert_next_wait();
4638
4639 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4640 paths.sort();
4641 assert_eq!(paths, ["foobar", "worm-node/data"]);
4642
4643 let mut narrow_consumer = user_producer
4645 .consume()
4646 .scope(&["worm-node".into()])
4647 .expect("should be able to narrow scope to worm-node")
4648 .announced();
4649
4650 narrow_consumer.assert_next_some("worm-node/data");
4651 narrow_consumer.assert_next_wait(); }
4653
4654 #[tokio::test]
4655 async fn test_duplicate_prefixes_deduped() {
4656 let origin = Origin::random().produce();
4657
4658 let producer = origin
4660 .scope(&["demo".into(), "demo".into()])
4661 .expect("should create producer");
4662
4663 let _broadcast = producer
4664 .create_broadcast("demo/stream", announce())
4665 .expect("publish allowed");
4666 settle().await;
4667
4668 let mut consumer = producer.consume().announced();
4669 consumer.assert_next_some("demo/stream");
4670 consumer.assert_next_wait();
4671 }
4672
4673 #[tokio::test]
4674 async fn test_overlapping_prefixes_deduped() {
4675 let origin = Origin::random().produce();
4676
4677 let producer = origin
4679 .scope(&["demo".into(), "demo/foo".into()])
4680 .expect("should create producer");
4681
4682 let _broadcast = producer
4684 .create_broadcast("demo/bar/stream", announce())
4685 .expect("publish allowed");
4686 settle().await;
4687
4688 let mut consumer = producer.consume().announced();
4689 consumer.assert_next_some("demo/bar/stream");
4690 consumer.assert_next_wait();
4691 }
4692
4693 #[tokio::test]
4694 async fn test_overlapping_prefixes_no_duplicate_announcements() {
4695 let origin = Origin::random().produce();
4696
4697 let producer = origin
4699 .scope(&["demo".into(), "demo/foo".into()])
4700 .expect("should create producer");
4701
4702 let _broadcast = producer
4703 .create_broadcast("demo/foo/stream", announce())
4704 .expect("publish allowed");
4705 settle().await;
4706
4707 let mut consumer = producer.consume().announced();
4708 consumer.assert_next_some("demo/foo/stream");
4710 consumer.assert_next_wait();
4711 }
4712
4713 #[tokio::test]
4714 async fn test_allowed_returns_deduped_prefixes() {
4715 let origin = Origin::random().produce();
4716
4717 let producer = origin
4718 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
4719 .expect("should create producer");
4720
4721 let allowed: Vec<_> = producer.allowed().collect();
4722 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
4723 }
4724
4725 #[tokio::test]
4726 async fn test_announced_broadcast_already_announced() {
4727 let origin = Origin::random().produce();
4728
4729 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4730 settle().await;
4731
4732 let consumer = origin.consume();
4733 let result = consumer.announced_broadcast("test").await.expect("should find it");
4734 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4735 }
4736
4737 #[tokio::test]
4738 async fn test_announced_broadcast_delayed() {
4739 tokio::time::pause();
4740
4741 let origin = Origin::random().produce();
4742
4743 let consumer = origin.consume();
4744
4745 let wait = tokio::spawn({
4747 let consumer = consumer.clone();
4748 async move { consumer.announced_broadcast("test").await }
4749 });
4750
4751 tokio::task::yield_now().await;
4753
4754 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4755 settle().await;
4756
4757 let result = wait.await.unwrap().expect("should find it");
4758 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4759 }
4760
4761 #[tokio::test]
4762 async fn test_announced_broadcast_ignores_unrelated_paths() {
4763 tokio::time::pause();
4764
4765 let origin = Origin::random().produce();
4766
4767 let consumer = origin.consume();
4768
4769 let wait = tokio::spawn({
4770 let consumer = consumer.clone();
4771 async move { consumer.announced_broadcast("target").await }
4772 });
4773
4774 tokio::task::yield_now().await;
4775
4776 let _other = origin.create_broadcast("other", announce()).unwrap();
4778 settle().await;
4779 tokio::task::yield_now().await;
4780 assert!(!wait.is_finished(), "must not resolve on unrelated path");
4781
4782 let _target = origin.create_broadcast("target", announce()).unwrap();
4783 settle().await;
4784 let result = wait.await.unwrap().expect("should find target");
4785 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
4786 }
4787
4788 #[tokio::test]
4789 async fn test_announced_broadcast_skips_nested_paths() {
4790 tokio::time::pause();
4791
4792 let origin = Origin::random().produce();
4793
4794 let consumer = origin.consume();
4795
4796 let wait = tokio::spawn({
4797 let consumer = consumer.clone();
4798 async move { consumer.announced_broadcast("foo").await }
4799 });
4800
4801 tokio::task::yield_now().await;
4802
4803 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
4805 settle().await;
4806 tokio::task::yield_now().await;
4807 assert!(!wait.is_finished(), "must not resolve on a nested path");
4808
4809 let _exact = origin.create_broadcast("foo", announce()).unwrap();
4810 settle().await;
4811 let result = wait.await.unwrap().expect("should find foo exactly");
4812 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
4813 }
4814
4815 #[tokio::test]
4816 async fn test_announced_broadcast_disallowed() {
4817 let origin = Origin::random().produce();
4818 let limited = origin
4819 .consume()
4820 .scope(&["allowed".into()])
4821 .expect("should create limited");
4822
4823 assert!(limited.announced_broadcast("notallowed").await.is_none());
4825 }
4826
4827 #[tokio::test]
4828 async fn test_announced_broadcast_scope_too_narrow() {
4829 let origin = Origin::random().produce();
4832 let limited = origin
4833 .consume()
4834 .scope(&["foo/specific".into()])
4835 .expect("should create limited");
4836
4837 let result = limited
4839 .announced_broadcast("foo")
4840 .now_or_never()
4841 .expect("must not block");
4842 assert!(result.is_none());
4843 }
4844
4845 #[tokio::test]
4849 async fn test_coalesce_announce_then_unannounce() {
4850 tokio::time::pause();
4852
4853 let origin = Origin::random().produce();
4854 let mut announced = origin.consume().announced();
4855
4856 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4857 settle().await;
4858 broadcast.finish();
4859
4860 settle().await;
4861
4862 announced.assert_next_wait();
4863 }
4864
4865 #[tokio::test]
4866 async fn test_coalesce_announce_unannounce_announce() {
4867 tokio::time::pause();
4870
4871 let origin = Origin::random().produce();
4872 let mut announced = origin.consume().announced();
4873
4874 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4875 settle().await;
4876 broadcast1.finish();
4877 settle().await;
4878 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4879 settle().await;
4880
4881 announced.assert_next_some("test");
4882 announced.assert_next_wait();
4883 }
4884
4885 #[tokio::test]
4886 async fn test_coalesce_unannounce_announce_preserved() {
4887 tokio::time::pause();
4890
4891 let origin = Origin::random().produce();
4892 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4893 settle().await;
4894
4895 let mut announced = origin.consume().announced();
4896 announced.assert_next_some("test");
4897
4898 broadcast1.finish();
4900 settle().await;
4901
4902 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4903 settle().await;
4904
4905 announced.assert_next_none("test");
4907 announced.assert_next_some("test");
4908 announced.assert_next_wait();
4909 }
4910
4911 #[tokio::test]
4912 async fn test_coalesce_unannounce_announce_unannounce() {
4913 tokio::time::pause();
4916
4917 let origin = Origin::random().produce();
4918 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4919 settle().await;
4920
4921 let mut announced = origin.consume().announced();
4922 announced.assert_next_some("test");
4923
4924 broadcast1.finish();
4925 settle().await;
4926
4927 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4928 settle().await;
4929 broadcast2.finish();
4930 settle().await;
4931
4932 announced.assert_next_none("test");
4933 announced.assert_next_wait();
4934 }
4935
4936 #[tokio::test]
4937 async fn test_coalesce_churn_bounded() {
4938 tokio::time::pause();
4943
4944 let origin = Origin::random().produce();
4945 let mut announced = origin.consume().announced();
4946
4947 for _ in 0..1000 {
4948 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4949 settle().await;
4950 broadcast.finish();
4951 }
4952 settle().await;
4953
4954 let mut collected = Vec::new();
4955 while let Some(update) = announced.try_next() {
4956 collected.push(update);
4957 }
4958 assert!(
4959 collected.len() <= 1,
4960 "expected at most one pending update, got {}",
4961 collected.len()
4962 );
4963 assert!(
4964 collected.iter().all(|a| a.path == Path::new("test")),
4965 "unexpected path in pending updates",
4966 );
4967 }
4968
4969 #[tokio::test]
4973 async fn test_consumer_clone_is_side_effect_free() {
4974 let origin = Origin::random().produce();
4975
4976 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
4977 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
4978 settle().await;
4979
4980 let consumer = origin.consume();
4981 let mut announced = consumer.announced();
4982
4983 for _ in 0..16 {
4986 let cloned = consumer.clone();
4987 assert!(cloned.get_broadcast("test1").is_some());
4988 assert!(cloned.get_broadcast("test2").is_some());
4989 }
4990
4991 let a1 = announced.try_next().expect("first announcement");
4994 let a2 = announced.try_next().expect("second announcement");
4995 announced.assert_next_wait();
4996
4997 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4998 paths.sort();
4999 assert_eq!(paths, ["test1", "test2"]);
5000
5001 let mut fresh = consumer.announced();
5003 let b1 = fresh.try_next().expect("backlog: first");
5004 let b2 = fresh.try_next().expect("backlog: second");
5005 fresh.assert_next_wait();
5006
5007 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
5008 paths.sort();
5009 assert_eq!(paths, ["test1", "test2"]);
5010 }
5011
5012 #[tokio::test]
5014 async fn dynamic_request_unroutable_without_handler() {
5015 let origin = Origin::random().produce();
5016 let consumer = origin.consume();
5017 assert!(matches!(
5018 consumer.request_broadcast("missing").await,
5019 Err(Error::Unroutable)
5020 ));
5021 }
5022
5023 #[tokio::test(start_paused = true)]
5026 async fn dynamic_request_served_not_announced() {
5027 let origin = Origin::random().produce();
5028 let mut dynamic = origin.dynamic();
5029 let consumer = origin.consume();
5030
5031 let mut announced = origin.consume().announced();
5033 announced.assert_next_wait();
5034
5035 let served = broadcast::Info::new().produce();
5036 let request_fut = consumer.request_broadcast("fallback");
5039
5040 let mut served_dynamic = served.dynamic();
5042
5043 let request = dynamic.requested_broadcast().await.unwrap();
5044 assert_eq!(request.path(), &Path::new("fallback"));
5045 request.accept(&served);
5046
5047 let broadcast = request_fut.await.unwrap();
5048 assert!(broadcast.is_clone(&served.consume()));
5049
5050 let track_fut = broadcast.track("video").unwrap().subscribe(None);
5052 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
5053 let mut track = track_fut.await.unwrap();
5054 producer.append_group().unwrap();
5055 track.assert_group();
5056
5057 announced.assert_next_wait();
5059 }
5060
5061 #[tokio::test(start_paused = true)]
5063 async fn dynamic_request_coalesces() {
5064 let origin = Origin::random().produce();
5065 let mut dynamic = origin.dynamic();
5066 let consumer = origin.consume();
5067
5068 let f1 = consumer.request_broadcast("dup");
5070 let f2 = consumer.request_broadcast("dup");
5071
5072 let request = dynamic.requested_broadcast().await.unwrap();
5074 assert_eq!(request.path(), &Path::new("dup"));
5075 assert!(
5076 dynamic.requested_broadcast().now_or_never().is_none(),
5077 "a coalesced request must not be served twice"
5078 );
5079
5080 let served = broadcast::Info::new().produce();
5082 request.accept(&served);
5083 assert!(f1.await.unwrap().is_clone(&served.consume()));
5084 assert!(f2.await.unwrap().is_clone(&served.consume()));
5085 }
5086
5087 #[tokio::test(start_paused = true)]
5090 async fn dynamic_request_dedups_served() {
5091 let origin = Origin::random().produce();
5092 let mut dynamic = origin.dynamic();
5093 let consumer = origin.consume();
5094
5095 let request_fut = consumer.request_broadcast("fallback");
5096 let request = dynamic.requested_broadcast().await.unwrap();
5097 let served = broadcast::Info::new().produce();
5098 request.accept(&served);
5099 let first = request_fut.await.unwrap();
5100 assert!(first.is_clone(&served.consume()));
5101
5102 let second = consumer.request_broadcast("fallback").await.unwrap();
5104 assert!(second.is_clone(&served.consume()));
5105
5106 assert!(
5108 dynamic.requested_broadcast().now_or_never().is_none(),
5109 "a still-live served broadcast must not be re-requested from the handler"
5110 );
5111 }
5112
5113 #[tokio::test(start_paused = true)]
5115 async fn dynamic_request_reserves_after_close() {
5116 let origin = Origin::random().produce();
5117 let mut dynamic = origin.dynamic();
5118 let consumer = origin.consume();
5119
5120 let request_fut = consumer.request_broadcast("fallback");
5121 let request = dynamic.requested_broadcast().await.unwrap();
5122 let served = broadcast::Info::new().produce();
5123 request.accept(&served);
5124 request_fut.await.unwrap();
5125
5126 drop(served);
5128
5129 let request_fut = consumer.request_broadcast("fallback");
5131 let request = dynamic.requested_broadcast().await.unwrap();
5132 assert_eq!(request.path(), &Path::new("fallback"));
5133 let served = broadcast::Info::new().produce();
5134 request.accept(&served);
5135 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
5136 }
5137
5138 #[tokio::test(start_paused = true)]
5141 async fn dynamic_request_served_cache_bounded() {
5142 let origin = Origin::random().produce();
5143 let mut dynamic = origin.dynamic();
5144 let consumer = origin.consume();
5145
5146 for i in 0..100 {
5147 let path = format!("one-shot/{i}");
5148 let request_fut = consumer.request_broadcast(&path);
5149 let request = dynamic.requested_broadcast().await.unwrap();
5150 let served = broadcast::Info::new().produce();
5151 request.accept(&served);
5152 request_fut.await.unwrap();
5153 drop(served);
5155 }
5156
5157 assert!(
5160 origin.dynamic.read().served.len() <= 4,
5161 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
5162 origin.dynamic.read().served.len()
5163 );
5164 }
5165
5166 #[tokio::test(start_paused = true)]
5169 async fn dynamic_request_coalesces_after_handoff() {
5170 let origin = Origin::random().produce();
5171 let mut dynamic = origin.dynamic();
5172 let consumer = origin.consume();
5173
5174 let f1 = consumer.request_broadcast("fallback");
5175 let request = dynamic.requested_broadcast().await.unwrap();
5177
5178 let f2 = consumer.request_broadcast("fallback");
5180 assert!(
5181 dynamic.requested_broadcast().now_or_never().is_none(),
5182 "a repeat request during hand-off must coalesce, not re-queue"
5183 );
5184
5185 let served = broadcast::Info::new().produce();
5187 request.accept(&served);
5188 assert!(f1.await.unwrap().is_clone(&served.consume()));
5189 assert!(f2.await.unwrap().is_clone(&served.consume()));
5190 }
5191
5192 #[tokio::test(start_paused = true)]
5194 async fn dynamic_request_dropped_after_handoff() {
5195 let origin = Origin::random().produce();
5196 let mut dynamic = origin.dynamic();
5197 let consumer = origin.consume();
5198
5199 let f1 = consumer.request_broadcast("fallback");
5200 let request = dynamic.requested_broadcast().await.unwrap();
5201 let f2 = consumer.request_broadcast("fallback");
5202
5203 drop(request);
5205 assert!(matches!(f1.await, Err(Error::Unroutable)));
5206 assert!(matches!(f2.await, Err(Error::Unroutable)));
5207 }
5208
5209 #[tokio::test(start_paused = true)]
5211 async fn dynamic_request_rejected() {
5212 let origin = Origin::random().produce();
5213 let mut dynamic = origin.dynamic();
5214 let consumer = origin.consume();
5215
5216 let request_fut = consumer.request_broadcast("fallback");
5217
5218 let request = dynamic.requested_broadcast().await.unwrap();
5219 request.reject(Error::Cancel);
5220
5221 assert!(matches!(request_fut.await, Err(Error::Cancel)));
5222 }
5223
5224 #[tokio::test(start_paused = true)]
5228 async fn dynamic_request_rerequest_after_reject() {
5229 let origin = Origin::random().produce();
5230 let mut dynamic = origin.dynamic();
5231 let consumer = origin.consume();
5232
5233 let f1 = consumer.request_broadcast("fallback");
5234 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
5235 assert!(matches!(f1.await, Err(Error::Unroutable)));
5236
5237 let served = broadcast::Info::new().produce();
5238 let f2 = consumer.request_broadcast("fallback");
5240 let request = dynamic.requested_broadcast().await.unwrap();
5241 assert_eq!(request.path(), &Path::new("fallback"));
5242 request.accept(&served);
5243 assert!(f2.await.unwrap().is_clone(&served.consume()));
5244 }
5245
5246 #[tokio::test(start_paused = true)]
5249 async fn dynamic_request_handler_dropped() {
5250 let origin = Origin::random().produce();
5251 let dynamic = origin.dynamic();
5252 let consumer = origin.consume();
5253
5254 let request_fut = consumer.request_broadcast("fallback");
5255 drop(dynamic);
5256 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
5257
5258 assert!(matches!(
5260 consumer.request_broadcast("again").await,
5261 Err(Error::Unroutable)
5262 ));
5263 }
5264
5265 #[tokio::test(start_paused = true)]
5269 async fn dynamic_request_accept_after_handler_dropped() {
5270 let origin = Origin::random().produce();
5271 let mut dynamic = origin.dynamic();
5272 let consumer = origin.consume();
5273
5274 let request_fut = consumer.request_broadcast("fallback");
5275
5276 let request = dynamic.requested_broadcast().await.unwrap();
5278 drop(dynamic);
5279
5280 let served = broadcast::Info::new().produce();
5281 request.accept(&served);
5283 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
5284 }
5285
5286 #[tokio::test(start_paused = true)]
5288 async fn dynamic_request_prefers_announced() {
5289 let origin = Origin::random().produce();
5290 let mut dynamic = origin.dynamic();
5291 let consumer = origin.consume();
5292
5293 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
5294 settle().await;
5295
5296 let got = consumer.request_broadcast("live").await.unwrap();
5297 assert!(
5298 got.is_clone(&consumer.get_broadcast("live").unwrap()),
5299 "should return the published broadcast"
5300 );
5301 assert!(
5302 dynamic.requested_broadcast().now_or_never().is_none(),
5303 "a published path must not queue a fallback request"
5304 );
5305 }
5306
5307 #[tokio::test(start_paused = true)]
5309 async fn dynamic_clone_keeps_alive() {
5310 let origin = Origin::random().produce();
5311 let dynamic = origin.dynamic();
5312 let consumer = origin.consume();
5313
5314 drop(dynamic.clone());
5315
5316 let request_fut = consumer.request_broadcast("fallback");
5319 assert!(
5320 request_fut.now_or_never().is_none(),
5321 "request should stay pending until served"
5322 );
5323 }
5324}