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 linger: Duration,
120}
121
122impl Default for Info {
123 fn default() -> Self {
126 Self {
127 id: Origin::UNKNOWN,
128 pool: cache::Pool::default(),
129 linger: Duration::ZERO,
130 }
131 }
132}
133
134impl Info {
135 pub fn new(id: Origin) -> Self {
137 Self { id, ..Self::default() }
138 }
139
140 pub fn with_pool(mut self, pool: cache::Pool) -> Self {
142 self.pool = pool;
143 self
144 }
145
146 pub fn with_linger(mut self, linger: Duration) -> Self {
149 self.linger = linger;
150 self
151 }
152
153 pub fn produce(self) -> Producer {
155 Producer::new(self)
156 }
157}
158
159impl TryFrom<u64> for Origin {
160 type Error = InvalidOrigin;
161
162 fn try_from(id: u64) -> Result<Self, Self::Error> {
163 Self::new(id)
164 }
165}
166
167impl fmt::Display for Origin {
168 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
169 self.id.fmt(f)
170 }
171}
172
173impl<V: Copy> Encode<V> for Origin
174where
175 u64: Encode<V>,
176{
177 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
178 self.id.encode(w, version)
179 }
180}
181
182impl<V: Copy> Decode<V> for Origin
183where
184 u64: Decode<V>,
185{
186 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
187 let id = u64::decode(r, version)?;
188 if id >= 1u64 << 62 {
189 return Err(DecodeError::InvalidValue);
190 }
191 Ok(Self { id })
192 }
193}
194
195pub(crate) const MAX_HOPS: usize = 32;
201
202#[derive(Debug, Clone, Default, PartialEq, Eq)]
207pub struct OriginList(Vec<Origin>);
208
209#[derive(Debug, Clone, Copy, PartialEq, Eq)]
211#[non_exhaustive]
212pub struct TooManyOrigins;
213
214impl fmt::Display for TooManyOrigins {
215 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
216 write!(f, "too many origins (max {MAX_HOPS})")
217 }
218}
219
220impl std::error::Error for TooManyOrigins {}
221
222impl From<TooManyOrigins> for DecodeError {
223 fn from(_: TooManyOrigins) -> Self {
224 DecodeError::BoundsExceeded
225 }
226}
227
228impl OriginList {
229 pub fn new() -> Self {
231 Self(Vec::new())
232 }
233
234 pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
236 if self.0.len() >= MAX_HOPS {
237 return Err(TooManyOrigins);
238 }
239 self.0.push(origin);
240 Ok(())
241 }
242
243 pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
246 for entry in &mut self.0 {
247 if *entry == target {
248 *entry = replacement;
249 return true;
250 }
251 }
252 false
253 }
254
255 pub fn contains(&self, origin: &Origin) -> bool {
257 self.0.contains(origin)
258 }
259
260 pub fn len(&self) -> usize {
262 self.0.len()
263 }
264
265 pub fn is_empty(&self) -> bool {
267 self.0.is_empty()
268 }
269
270 pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
272 self.0.iter()
273 }
274
275 pub fn as_slice(&self) -> &[Origin] {
277 &self.0
278 }
279}
280
281impl TryFrom<Vec<Origin>> for OriginList {
282 type Error = TooManyOrigins;
283
284 fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
285 if v.len() > MAX_HOPS {
286 return Err(TooManyOrigins);
287 }
288 Ok(Self(v))
289 }
290}
291
292impl<'a> IntoIterator for &'a OriginList {
293 type Item = &'a Origin;
294 type IntoIter = std::slice::Iter<'a, Origin>;
295
296 fn into_iter(self) -> Self::IntoIter {
297 self.iter()
298 }
299}
300
301impl<V: Copy> Encode<V> for OriginList
302where
303 u64: Encode<V>,
304 Origin: Encode<V>,
305{
306 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
307 (self.0.len() as u64).encode(w, version)?;
308 for origin in &self.0 {
309 origin.encode(w, version)?;
310 }
311 Ok(())
312 }
313}
314
315impl<V: Copy> Decode<V> for OriginList
316where
317 u64: Decode<V>,
318 Origin: Decode<V>,
319{
320 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
321 let count = u64::decode(r, version)? as usize;
322 if count > MAX_HOPS {
323 return Err(DecodeError::BoundsExceeded);
324 }
325 let mut list = Vec::with_capacity(count);
326 for _ in 0..count {
327 list.push(Origin::decode(r, version)?);
328 }
329 Ok(Self(list))
330 }
331}
332
333static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
334
335#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
336struct ConsumerId(u64);
337
338impl ConsumerId {
339 fn new() -> Self {
340 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
341 }
342}
343
344struct OriginBroadcast {
350 path: PathOwned,
351 broadcast: broadcast::Producer,
353 state: kio::Producer<FrontState>,
356 announced: bool,
357}
358
359fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
368 (hops.len(), fnv_key(name, hops.iter().copied()))
369}
370
371fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
384 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
386
387 let mut hash = SEED;
388 for &byte in name.as_str().as_bytes() {
389 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
390 }
391 for origin in origins {
392 for &byte in &origin.id().to_le_bytes() {
393 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
394 }
395 }
396
397 hash
398}
399
400fn route_order(name: &Path, route: &broadcast::Route) -> (bool, u64, usize, u64) {
409 let (len, hash) = route_key(name, &route.hops);
410 (!route.announce, route.cost, len, hash)
411}
412
413enum PendingUpdate {
422 Announce(broadcast::Consumer),
423 Unannounce,
424 UnannounceAnnounce(broadcast::Consumer),
425}
426
427#[derive(Default)]
432struct OriginConsumerState {
433 pending: BTreeMap<PathOwned, PendingUpdate>,
434}
435
436impl OriginConsumerState {
437 fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
438 let new = match self.pending.remove(&path) {
439 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
441 Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
443 PendingUpdate::UnannounceAnnounce(broadcast)
444 }
445 };
446 self.pending.insert(path, new);
447 }
448
449 fn apply_unannounce(&mut self, path: PathOwned) {
450 match self.pending.remove(&path) {
451 Some(PendingUpdate::Announce(_)) => {}
453 None | Some(PendingUpdate::Unannounce) => {
454 self.pending.insert(path, PendingUpdate::Unannounce);
455 }
456 Some(PendingUpdate::UnannounceAnnounce(_)) => {
459 self.pending.insert(path, PendingUpdate::Unannounce);
460 }
461 }
462 }
463
464 fn take(&mut self) -> Option<OriginAnnounce> {
466 let path = self.pending.keys().next()?.clone();
467 let broadcast = match self.pending.remove(&path).unwrap() {
468 PendingUpdate::Announce(broadcast) => Some(broadcast),
469 PendingUpdate::Unannounce => None,
470 PendingUpdate::UnannounceAnnounce(broadcast) => {
471 self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
474 None
475 }
476 };
477 Some(OriginAnnounce { path, broadcast })
478 }
479}
480
481#[derive(Clone)]
482struct AnnounceConsumerNotify {
483 root: PathOwned,
484 state: kio::Producer<OriginConsumerState>,
485}
486
487impl AnnounceConsumerNotify {
488 fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer) {
489 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
490 self.state
491 .write()
492 .ok()
493 .expect("consumer closed")
494 .apply_announce(path, broadcast);
495 }
496
497 fn unannounce(&self, path: impl AsPath) {
498 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
499 self.state.write().ok().expect("consumer closed").apply_unannounce(path);
500 }
501}
502
503struct NotifyNode {
504 parent: Option<Lock<NotifyNode>>,
505
506 consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
509}
510
511impl NotifyNode {
512 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
513 Self {
514 parent,
515 consumers: HashMap::new(),
516 }
517 }
518
519 fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer) {
520 for consumer in self.consumers.values() {
521 consumer.announce(path.as_path(), broadcast.clone());
522 }
523
524 if let Some(parent) = &self.parent {
525 parent.lock().announce(path, broadcast);
526 }
527 }
528
529 fn unannounce(&mut self, path: impl AsPath) {
530 for consumer in self.consumers.values() {
531 consumer.unannounce(path.as_path());
532 }
533
534 if let Some(parent) = &self.parent {
535 parent.lock().unannounce(path);
536 }
537 }
538}
539
540struct OriginNode {
541 broadcast: Option<OriginBroadcast>,
544
545 nested: HashMap<String, Lock<OriginNode>>,
547
548 notify: Lock<NotifyNode>,
550}
551
552impl OriginNode {
553 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
554 Self {
555 broadcast: None,
556 nested: HashMap::new(),
557 notify: Lock::new(NotifyNode::new(parent)),
558 }
559 }
560
561 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
562 let (dir, rest) = path.next_part().expect("leaf called with empty path");
563
564 let next = self.entry(dir);
565 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
566 }
567
568 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
569 match self.nested.get(dir) {
570 Some(next) => next.clone(),
571 None => {
572 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
573 self.nested.insert(dir.to_string(), next.clone());
574 next
575 }
576 }
577 }
578
579 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
583 let Some(existing) = &mut self.broadcast else { return };
584 if !existing.state.same_channel(expect) || existing.announced == announce {
585 return;
586 }
587 existing.announced = announce;
588 let path = existing.path.clone();
589 let consumer = existing.broadcast.consume();
590 let mut notify = self.notify.lock();
591 if announce {
592 notify.announce(&path, &consumer);
593 } else {
594 notify.unannounce(&path);
595 }
596 }
597
598 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
599 self.consume_initial(&mut notify);
600 self.notify.lock().consumers.insert(id, notify);
601 }
602
603 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
604 if let Some(broadcast) = &self.broadcast
607 && broadcast.announced
608 {
609 notify.announce(&broadcast.path, broadcast.broadcast.consume());
610 }
611
612 for nested in self.nested.values() {
614 nested.lock().consume_initial(notify);
615 }
616 }
617
618 fn consume_broadcast(&self, rest: impl AsPath) -> Option<broadcast::Consumer> {
619 let rest = rest.as_path();
620
621 if let Some((dir, rest)) = rest.next_part() {
622 let node = self.nested.get(dir)?.lock();
623 node.consume_broadcast(&rest)
624 } else {
625 self.broadcast.as_ref().map(|b| b.broadcast.consume())
626 }
627 }
628
629 fn unconsume(&mut self, id: ConsumerId) {
630 self.notify.lock().consumers.remove(&id).expect("consumer not found");
631 if self.is_empty() {
632 }
635 }
636
637 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
641 let relative = relative.as_path();
642
643 if let Some((dir, relative)) = relative.next_part() {
644 let Some(nested) = self.nested.get(dir) else { return };
645 let nested = nested.clone();
646 let mut locked = nested.lock();
647 locked.remove(expect, &relative);
648
649 if locked.is_empty() {
650 drop(locked);
651 self.nested.remove(dir);
652 }
653 } else if let Some(existing) = &self.broadcast
654 && existing.state.same_channel(expect)
655 {
656 let existing = self.broadcast.take().expect("checked above");
657 if existing.announced {
658 self.notify.lock().unannounce(&existing.path);
659 }
660 }
661 }
662
663 fn is_empty(&self) -> bool {
664 self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
665 }
666}
667
668#[derive(Clone)]
669struct OriginNodes {
670 nodes: Vec<(PathOwned, Lock<OriginNode>)>,
671}
672
673impl OriginNodes {
674 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
677 let mut roots = Vec::new();
678
679 for (root, state) in &self.nodes {
680 for prefix in prefixes {
681 if root.has_prefix(prefix) {
682 roots.push((root.to_owned(), state.clone()));
684 continue;
685 }
686
687 if let Some(suffix) = prefix.strip_prefix(root) {
688 let nested = state.lock().leaf(&suffix);
690 roots.push((prefix.to_owned(), nested));
691 }
692 }
693 }
694
695 if roots.is_empty() {
696 None
697 } else {
698 Some(Self { nodes: roots })
699 }
700 }
701
702 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
703 let new_root = new_root.as_path();
704 let mut roots = Vec::new();
705
706 if new_root.is_empty() {
707 return Some(self.clone());
708 }
709
710 for (root, state) in &self.nodes {
711 if let Some(suffix) = root.strip_prefix(&new_root) {
712 roots.push((suffix.to_owned(), state.clone()));
714 } else if let Some(suffix) = new_root.strip_prefix(root) {
715 let nested = state.lock().leaf(&suffix);
718 roots.push(("".into(), nested));
719 }
720 }
721
722 if roots.is_empty() {
723 None
724 } else {
725 Some(Self { nodes: roots })
726 }
727 }
728
729 pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
731 let path = path.as_path();
732
733 for (root, state) in &self.nodes {
734 if let Some(suffix) = path.strip_prefix(root) {
735 return Some((state.clone(), suffix.to_owned()));
736 }
737 }
738
739 None
740 }
741}
742
743impl Default for OriginNodes {
744 fn default() -> Self {
745 Self {
746 nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
747 }
748 }
749}
750
751#[derive(Clone)]
753pub struct OriginAnnounce {
754 pub path: PathOwned,
756 pub broadcast: Option<broadcast::Consumer>,
762}
763
764#[derive(Clone)]
766pub struct Producer {
767 info: Origin,
771
772 nodes: OriginNodes,
775
776 root: PathOwned,
778
779 dynamic: kio::Shared<OriginDynamicState>,
783
784 pool: cache::Pool,
787
788 linger: Duration,
791
792 stats: stats::Session,
796}
797
798impl std::ops::Deref for Producer {
799 type Target = Origin;
800
801 fn deref(&self) -> &Self::Target {
802 &self.info
803 }
804}
805
806impl Producer {
807 pub fn new(info: Info) -> Self {
811 Self {
812 info: info.id,
813 nodes: OriginNodes::default(),
814 root: PathOwned::default(),
815 dynamic: kio::Shared::default(),
816 pool: info.pool,
817 linger: info.linger,
818 stats: stats::Session::default(),
819 }
820 }
821
822 pub fn with_stats(mut self, session: stats::Session) -> Self {
826 self.stats = session;
827 self
828 }
829
830 pub fn with_linger(mut self, linger: Duration) -> Self {
839 self.linger = linger;
840 self
841 }
842
843 pub fn info(&self) -> Info {
846 Info {
847 id: self.info,
848 pool: self.pool.clone(),
849 linger: self.linger,
850 }
851 }
852
853 pub(crate) fn empty(info: Origin) -> Self {
858 Self {
859 info,
860 nodes: OriginNodes { nodes: Vec::new() },
861 root: PathOwned::default(),
862 dynamic: kio::Shared::default(),
863 pool: cache::Pool::default(),
864 linger: Duration::ZERO,
865 stats: stats::Session::default(),
866 }
867 }
868
869 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
908 let path = path.as_path();
909
910 debug_assert!(
911 !route.hops.contains(&self.info),
912 "create_broadcast called with a looping hop chain",
913 );
914
915 let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
916 let full = self.root.join(&path).to_owned();
917
918 if full.parts().count() > Path::MAX_PARTS {
922 return Err(BoundsExceeded.into());
923 }
924
925 let ingress = self.stats.ingress(&full);
929
930 let mut source = broadcast::Info { origin: self.info() }
931 .produce()
932 .with_stats(ingress.clone());
933 source.set_route(route).expect("fresh producer");
934
935 web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
936
937 Ok(source)
938 }
939
940 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
946 let prefixes = PathPrefixes::new(prefixes);
947 Some(Producer {
948 info: self.info,
949 nodes: self.nodes.select(&prefixes)?,
950 root: self.root.clone(),
951 dynamic: self.dynamic.clone(),
952 pool: self.pool.clone(),
953 linger: self.linger,
954 stats: self.stats.clone(),
955 })
956 }
957
958 pub fn dynamic(&self) -> Dynamic {
967 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
968 }
969
970 pub fn consume(&self) -> Consumer {
975 Consumer::new(
978 self.info,
979 self.root.clone(),
980 self.nodes.clone(),
981 self.dynamic.clone(),
982 stats::Session::default(),
983 )
984 }
985
986 pub fn announces(&self) -> AnnounceProducer {
992 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
993 }
994
995 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1000 let prefix = prefix.as_path();
1001
1002 Some(Self {
1003 info: self.info,
1004 root: self.root.join(&prefix).to_owned(),
1005 nodes: self.nodes.root(&prefix)?,
1006 dynamic: self.dynamic.clone(),
1007 pool: self.pool.clone(),
1008 linger: self.linger,
1009 stats: self.stats.clone(),
1010 })
1011 }
1012
1013 pub fn root(&self) -> &Path<'_> {
1015 &self.root
1016 }
1017
1018 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1021 self.nodes.nodes.iter().map(|(root, _)| root)
1022 }
1023
1024 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1026 self.root.join(path)
1027 }
1028}
1029
1030const MAX_TRACK_RETRIES: u32 = 3;
1035
1036struct FrontRoute {
1038 id: u64,
1039 route: broadcast::Route,
1042 source: broadcast::Consumer,
1044}
1045
1046struct FrontState {
1048 path: PathOwned,
1050 self_origin: Origin,
1052 next_route: u64,
1053 routes: Vec<FrontRoute>,
1054 active: Option<u64>,
1056 linger: Duration,
1059 closed: bool,
1063}
1064
1065impl FrontState {
1066 fn best_route(&self) -> Option<u64> {
1069 self.routes
1070 .iter()
1071 .min_by_key(|r| route_order(&self.path.as_path(), &r.route))
1072 .map(|r| r.id)
1073 }
1074
1075 fn reselect(&mut self, carrying: bool) {
1093 let best = self.best_route();
1094 if carrying
1095 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1096 && best_id != cur_id
1097 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1098 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1099 && incumbent.route.announce
1100 && candidate.route.cost < incumbent.route.cost
1101 && candidate.route.advertised == 0
1102 && candidate.route.hops.len() >= 2
1103 && !self.handover_allowed(&candidate.route)
1104 {
1105 return;
1107 }
1108 self.active = best;
1109 }
1110
1111 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1123 let name = self.path.as_path();
1124 match route.hops.iter().last() {
1125 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1126 None => true,
1127 }
1128 }
1129
1130 fn active_route(&self) -> Option<broadcast::Route> {
1133 let id = self.active?;
1134 self.routes.iter().find(|r| r.id == id).map(|r| r.route.clone())
1135 }
1136}
1137
1138fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1148 let mut leaf_guard = leaf.lock();
1153 let advert = state.read().active_route();
1154 if let Some(advert) = advert {
1155 let announce = advert.announce;
1156 let _ = broadcast.clone().set_route(advert);
1157 leaf_guard.set_announced(state, announce);
1158 }
1159}
1160
1161fn detach_source(
1174 state: &kio::Producer<FrontState>,
1175 broadcast: &broadcast::Producer,
1176 leaf: &Lock<OriginNode>,
1177 id: u64,
1178 graceful: bool,
1179) {
1180 let close = {
1181 let carrying = broadcast.demand().is_used();
1182 let Ok(mut s) = state.write() else { return };
1183 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1184 return;
1185 };
1186 s.routes.remove(pos);
1187 s.reselect(carrying);
1188 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1189 s.closed = true;
1192 true
1193 } else {
1194 false
1195 }
1196 };
1197 if close {
1198 broadcast.abort_spliced(Error::Dropped);
1199 }
1200 sync_front(state, broadcast, leaf);
1201}
1202
1203async fn run_source(
1207 origin: Info,
1208 node: Lock<OriginNode>,
1209 full: PathOwned,
1210 rest: PathOwned,
1211 mut source: broadcast::Consumer,
1212 ingress: stats::Scope,
1213) {
1214 let Ok(route) = source.route_changed().await else {
1218 return;
1220 };
1221
1222 let mut announce = route.announce.then(|| ingress.announce());
1227
1228 let leaf = if rest.is_empty() {
1229 node.clone()
1230 } else {
1231 node.lock().leaf(&rest)
1232 };
1233
1234 let (state, broadcast, id) = attach_source(&origin, &node, &leaf, &full, &rest, &source, route);
1235
1236 loop {
1237 match source.route_changed().await {
1238 Ok(route) => {
1239 let announced = route.announce;
1240 {
1241 let carrying = broadcast.demand().is_used();
1242 let Ok(mut s) = state.write() else { return };
1243 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1244 return;
1245 };
1246 if entry.route == route {
1247 continue;
1248 }
1249 entry.route = route;
1250 s.reselect(carrying);
1251 }
1252 match (announced, announce.is_some()) {
1254 (true, false) => announce = Some(ingress.announce()),
1255 (false, true) => announce = None,
1256 _ => {}
1257 }
1258 sync_front(&state, &broadcast, &leaf);
1259 }
1260 Err(_) => {
1261 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1264 return;
1265 }
1266 }
1267 }
1268}
1269
1270fn attach_source(
1275 origin: &Info,
1276 node: &Lock<OriginNode>,
1277 leaf: &Lock<OriginNode>,
1278 full: &PathOwned,
1279 rest: &PathOwned,
1280 source: &broadcast::Consumer,
1281 route: broadcast::Route,
1282) -> (kio::Producer<FrontState>, broadcast::Producer, u64) {
1283 let mut leaf_guard = leaf.lock();
1284
1285 if let Some(existing) = &leaf_guard.broadcast {
1288 let mut joined = None;
1289 let carrying = existing.broadcast.demand().is_used();
1290 if let Ok(mut s) = existing.state.write()
1291 && !s.closed
1292 {
1293 let id = s.next_route;
1294 s.next_route += 1;
1295 s.routes.push(FrontRoute {
1296 id,
1297 route: route.clone(),
1298 source: source.clone(),
1299 });
1300 s.reselect(carrying);
1301 joined = Some(id);
1302 }
1303 if let Some(id) = joined {
1304 let state = existing.state.clone();
1305 let broadcast = existing.broadcast.clone();
1306 drop(leaf_guard);
1307 sync_front(&state, &broadcast, leaf);
1308 return (state, broadcast, id);
1309 }
1310 }
1311
1312 let announce = route.announce;
1314 let broadcast = broadcast::Producer::new_spliced(broadcast::Info { origin: origin.clone() });
1315 let _ = broadcast.clone().set_route(route.clone());
1316 let state = kio::Producer::new(FrontState {
1317 path: full.clone(),
1318 self_origin: origin.id,
1319 next_route: 1,
1320 routes: vec![FrontRoute {
1321 id: 0,
1322 route,
1323 source: source.clone(),
1324 }],
1325 active: Some(0),
1326 linger: origin.linger,
1327 closed: false,
1328 });
1329
1330 if let Some(stale) = leaf_guard.broadcast.take()
1334 && stale.announced
1335 {
1336 leaf_guard.notify.lock().unannounce(&stale.path);
1337 }
1338 let entry = OriginBroadcast {
1339 path: full.clone(),
1340 broadcast: broadcast.clone(),
1341 state: state.clone(),
1342 announced: announce,
1343 };
1344 if entry.announced {
1345 leaf_guard.notify.lock().announce(full, &broadcast.consume());
1346 }
1347 leaf_guard.broadcast = Some(entry);
1348 drop(leaf_guard);
1349
1350 web_async::spawn(run_front(state.clone(), broadcast.clone(), node.clone(), rest.clone()));
1351
1352 (state, broadcast, 0)
1353}
1354
1355async fn run_front(
1358 state: kio::Producer<FrontState>,
1359 mut broadcast: broadcast::Producer,
1360 node: Lock<OriginNode>,
1361 rest: PathOwned,
1362) {
1363 enum Step {
1364 Serve(Arc<str>, super::resume::Producer),
1365 Changed,
1367 Expired,
1369 Closed,
1370 }
1371
1372 let linger = state.read().linger;
1373 let mut deadline: Option<web_async::time::Instant> = None;
1378
1379 loop {
1380 let empty = {
1381 let s = state.read();
1382 !s.closed && s.routes.is_empty()
1383 };
1384 deadline = match (empty, deadline) {
1385 (true, None) => web_async::time::Instant::now().checked_add(linger),
1388 (true, at) => at,
1389 (false, _) => None,
1390 };
1391
1392 let step = {
1393 let mut sleep = std::pin::pin!(async {
1396 match deadline {
1397 Some(at) => {
1398 web_async::time::sleep(at.saturating_duration_since(web_async::time::Instant::now())).await
1399 }
1400 None => std::future::pending().await,
1401 }
1402 });
1403 let mut fired = false;
1404 kio::wait(|waiter| {
1405 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1406 return Poll::Ready(Step::Serve(name, resume));
1407 }
1408 match state.poll(waiter, |s| {
1411 if s.closed || s.routes.is_empty() != empty {
1412 Poll::Ready(())
1413 } else {
1414 Poll::Pending
1415 }
1416 }) {
1417 Poll::Ready(Ok(guard)) => {
1418 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1419 }
1420 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1421 Poll::Pending => {}
1422 }
1423 if deadline.is_some() && !fired && waiter.poll_future(sleep.as_mut()).is_ready() {
1424 fired = true;
1425 }
1426 match fired {
1427 true => Poll::Ready(Step::Expired),
1428 false => Poll::Pending,
1429 }
1430 })
1431 .await
1432 };
1433
1434 match step {
1435 Step::Serve(name, resume) => {
1436 web_async::spawn(serve_track(state.clone(), name, resume));
1439 }
1440 Step::Changed => {}
1441 Step::Expired => {
1442 let close = {
1446 let Ok(mut s) = state.write() else { break };
1447 if !s.closed && s.routes.is_empty() {
1448 s.closed = true;
1449 true
1450 } else {
1451 false
1452 }
1453 };
1454 if close {
1455 break;
1456 }
1457 }
1458 Step::Closed => break,
1459 }
1460 }
1461
1462 broadcast.abort_spliced(Error::Dropped);
1464
1465 broadcast.finish();
1467
1468 node.lock().remove(&state, &rest);
1471}
1472
1473async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1481 enum Step {
1482 Closed,
1483 Splice(u64, broadcast::Consumer),
1484 Complete,
1485 Failed,
1486 }
1487
1488 let mut fails = 0u32;
1489 let mut serving: Option<(u64, track::Consumer)> = None;
1491 let mut dead: Option<u64> = None;
1495
1496 loop {
1497 let serving_id = serving.as_ref().map(|(id, _)| *id);
1498 let step = kio::wait(|waiter| {
1499 match state.poll(waiter, |s| {
1503 if s.closed || matches!(s.active, Some(active) if Some(active) != serving_id && Some(active) != dead) {
1504 Poll::Ready(())
1505 } else {
1506 Poll::Pending
1507 }
1508 }) {
1509 Poll::Ready(Ok(guard)) => {
1510 if guard.closed {
1511 return Poll::Ready(Step::Closed);
1512 }
1513 let active = guard.active.expect("predicate guaranteed an active source");
1514 let source = guard
1515 .routes
1516 .iter()
1517 .find(|r| r.id == active)
1518 .expect("active source in table")
1519 .source
1520 .clone();
1521 return Poll::Ready(Step::Splice(active, source));
1522 }
1523 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1524 Poll::Pending => {}
1525 }
1526
1527 if let Some((_, track)) = &serving
1530 && let Poll::Ready(result) = track.poll_complete(waiter)
1531 {
1532 return Poll::Ready(match result {
1533 Ok(()) => Step::Complete,
1534 Err(_) => Step::Failed,
1535 });
1536 }
1537 Poll::Pending
1538 })
1539 .await;
1540
1541 match step {
1542 Step::Closed => return,
1544 Step::Complete => {
1545 let _ = resume.finish();
1546 return;
1547 }
1548 Step::Failed => {
1549 serving = None;
1552 }
1553 Step::Splice(id, source) => {
1554 let attempt = match source.track(&name) {
1558 Ok(track) => {
1559 let query = track.info().into_inner();
1562 let info = kio::wait(|waiter| {
1563 if let Poll::Ready(result) = query.poll(waiter) {
1564 return Poll::Ready(Some(result));
1565 }
1566 match state.poll(waiter, |s| {
1567 if s.closed || s.active != Some(id) {
1568 Poll::Ready(())
1569 } else {
1570 Poll::Pending
1571 }
1572 }) {
1573 Poll::Ready(_) => Poll::Ready(None),
1574 Poll::Pending => Poll::Pending,
1575 }
1576 })
1577 .await;
1578 match info {
1579 None => continue,
1582 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
1586 Poll::Ready(Err(err)) => Err(err),
1587 _ => Ok(track),
1588 },
1589 Some(Err(err)) => Err(err),
1590 }
1591 }
1592 Err(err) => Err(err),
1593 };
1594
1595 match attempt {
1596 Ok(track) => {
1597 if resume.takeover(&track).is_err() {
1598 return;
1601 }
1602 fails = 0;
1605 dead = None;
1606 serving = Some((id, track));
1607 }
1608 Err(_) if source.is_closing() => {
1612 dead = Some(id);
1613 serving = None;
1614 }
1615 Err(err) => {
1616 fails += 1;
1617 if fails >= MAX_TRACK_RETRIES {
1618 tracing::debug!(name = %name, %err, "aborting unservable track");
1619 let _ = resume.abort(Error::Unroutable);
1620 return;
1621 }
1622 serving = None;
1623 }
1624 }
1625 }
1626 }
1627 }
1628}
1629
1630#[derive(Default)]
1636struct OriginDynamicState {
1637 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
1640
1641 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
1647}
1648
1649#[derive(Default)]
1656struct PendingBroadcast {
1657 resolved: Option<Result<broadcast::Consumer, Error>>,
1658}
1659
1660pub struct Dynamic {
1671 info: Origin,
1672 root: PathOwned,
1673 state: kio::Shared<OriginDynamicState>,
1674}
1675
1676impl Clone for Dynamic {
1677 fn clone(&self) -> Self {
1678 self.state.lock().requests.add_handler();
1682
1683 Self {
1684 info: self.info,
1685 root: self.root.clone(),
1686 state: self.state.clone(),
1687 }
1688 }
1689}
1690
1691impl Dynamic {
1692 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
1693 state.lock().requests.add_handler();
1694
1695 Self { info, root, state }
1696 }
1697
1698 pub fn info(&self) -> &Origin {
1700 &self.info
1701 }
1702
1703 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
1705 let mut state = ready!(self.state.poll(waiter, |state| {
1706 if state.requests.has_queued() {
1707 Poll::Ready(())
1708 } else {
1709 Poll::Pending
1710 }
1711 }));
1712
1713 let path = state.requests.pop().expect("predicate guaranteed a request");
1714 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
1720 Poll::Ready(Ok(Request {
1721 path,
1722 producer,
1723 state: self.state.clone(),
1724 }))
1725 }
1726
1727 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
1730 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
1731 }
1732
1733 pub fn root(&self) -> &Path<'_> {
1735 &self.root
1736 }
1737}
1738
1739impl Drop for Dynamic {
1740 fn drop(&mut self) {
1741 let mut state = self.state.lock();
1744 if state.requests.remove_handler() {
1745 state.requests.drain_queued();
1749 }
1750 }
1751}
1752
1753pub struct Request {
1760 path: PathOwned,
1762
1763 producer: kio::Producer<PendingBroadcast>,
1766
1767 state: kio::Shared<OriginDynamicState>,
1769}
1770
1771impl Request {
1772 pub fn path(&self) -> &Path<'_> {
1774 &self.path
1775 }
1776
1777 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
1783 let broadcast = broadcast.consume();
1784
1785 let resolved = {
1791 let mut state = self.state.lock();
1792 let existing = state.served.insert(self.path.clone(), broadcast.weak());
1793 state
1794 .requests
1795 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1796 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
1797 };
1798
1799 if let Ok(mut pending) = self.producer.write() {
1800 pending.resolved = Some(Ok(resolved));
1801 }
1802 }
1804
1805 pub fn reject(self, err: Error) {
1807 self.state
1808 .lock()
1809 .requests
1810 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1811 if let Ok(mut state) = self.producer.write() {
1812 state.resolved = Some(Err(err));
1813 }
1814 }
1815}
1816
1817impl Drop for Request {
1818 fn drop(&mut self) {
1819 self.state
1827 .lock()
1828 .requests
1829 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
1830 }
1831}
1832
1833pub struct Requesting {
1840 inner: RequestState,
1841 stats: stats::Scope,
1844}
1845
1846enum RequestState {
1847 Ready(broadcast::Consumer),
1849 Failed(Error),
1852 Pending(kio::Consumer<PendingBroadcast>),
1854}
1855
1856impl Requesting {
1857 fn ready(broadcast: broadcast::Consumer) -> Self {
1858 Self {
1859 inner: RequestState::Ready(broadcast),
1860 stats: stats::Scope::default(),
1861 }
1862 }
1863
1864 fn failed(error: Error) -> Self {
1865 Self {
1866 inner: RequestState::Failed(error),
1867 stats: stats::Scope::default(),
1868 }
1869 }
1870
1871 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
1872 Self {
1873 inner: RequestState::Pending(consumer),
1874 stats: stats::Scope::default(),
1875 }
1876 }
1877
1878 fn with_stats(mut self, scope: stats::Scope) -> Self {
1879 self.stats = scope;
1880 self
1881 }
1882
1883 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
1885 match &self.inner {
1886 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
1887 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
1888 RequestState::Pending(consumer) => Poll::Ready(
1889 match ready!(consumer.poll(waiter, |state| match &state.resolved {
1890 Some(result) => Poll::Ready(result.clone()),
1891 None => Poll::Pending,
1892 })) {
1893 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
1894 Err(_closed) => Err(Error::Unroutable),
1896 },
1897 ),
1898 }
1899 }
1900}
1901
1902impl kio::Pollable for Requesting {
1903 type Output = Result<broadcast::Consumer, Error>;
1904
1905 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
1906 self.poll_ok(waiter)
1907 }
1908}
1909
1910pub trait Consume<T> {
1918 fn consume(&self) -> T;
1920}
1921
1922impl<T, U: Consume<T>> Consume<T> for &U {
1923 fn consume(&self) -> T {
1924 (**self).consume()
1925 }
1926}
1927
1928impl Consume<Consumer> for Producer {
1929 fn consume(&self) -> Consumer {
1930 Consumer::new(
1934 self.info,
1935 self.root.clone(),
1936 self.nodes.clone(),
1937 self.dynamic.clone(),
1938 stats::Session::default(),
1939 )
1940 }
1941}
1942
1943impl Consume<Consumer> for Consumer {
1944 fn consume(&self) -> Consumer {
1945 self.clone()
1946 }
1947}
1948
1949impl Consume<broadcast::Consumer> for broadcast::Producer {
1950 fn consume(&self) -> broadcast::Consumer {
1951 self.consume()
1953 }
1954}
1955
1956impl Consume<broadcast::Consumer> for broadcast::Consumer {
1957 fn consume(&self) -> broadcast::Consumer {
1958 self.clone()
1959 }
1960}
1961
1962impl Consume<track::Consumer> for track::Producer {
1963 fn consume(&self) -> track::Consumer {
1964 self.consume()
1965 }
1966}
1967
1968impl Consume<track::Consumer> for track::Consumer {
1969 fn consume(&self) -> track::Consumer {
1970 self.clone()
1971 }
1972}
1973
1974#[derive(Clone)]
1980pub struct Consumer {
1981 info: Origin,
1983 nodes: OriginNodes,
1984
1985 root: PathOwned,
1987
1988 dynamic: kio::Shared<OriginDynamicState>,
1991
1992 stats: stats::Session,
1996}
1997
1998impl std::ops::Deref for Consumer {
1999 type Target = Origin;
2000
2001 fn deref(&self) -> &Self::Target {
2002 &self.info
2003 }
2004}
2005
2006impl Consumer {
2007 fn new(
2008 info: Origin,
2009 root: PathOwned,
2010 nodes: OriginNodes,
2011 dynamic: kio::Shared<OriginDynamicState>,
2012 stats: stats::Session,
2013 ) -> Self {
2014 Self {
2015 info,
2016 nodes,
2017 root,
2018 dynamic,
2019 stats,
2020 }
2021 }
2022
2023 pub fn with_stats(mut self, session: stats::Session) -> Self {
2027 self.stats = session;
2028 self
2029 }
2030
2031 fn untagged(&self) -> Self {
2035 Self {
2036 stats: stats::Session::default(),
2037 ..self.clone()
2038 }
2039 }
2040
2041 pub(crate) fn empty(&self) -> Self {
2046 Self {
2047 info: self.info,
2048 nodes: OriginNodes { nodes: Vec::new() },
2049 root: self.root.clone(),
2050 dynamic: self.dynamic.clone(),
2051 stats: self.stats.clone(),
2052 }
2053 }
2054
2055 pub fn announced(&self) -> AnnounceConsumer {
2062 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2063 }
2064
2065 pub fn consume(&self) -> Self {
2067 self.clone()
2068 }
2069
2070 fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2077 let path = path.as_path();
2078 let (root, rest) = self.nodes.get(&path)?;
2079 let state = root.lock();
2080 state.consume_broadcast(&rest)
2081 }
2082
2083 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2095 let path = path.as_path();
2096
2097 let consumer = self.scope(std::slice::from_ref(&path))?;
2099
2100 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2104 return None;
2105 }
2106
2107 let mut announced = consumer.untagged().announced();
2111 let scope = self.stats.egress(self.root.join(&path).to_owned());
2112 loop {
2113 let OriginAnnounce {
2114 path: announced_path,
2115 broadcast,
2116 } = announced.next().await?;
2117 if announced_path.as_path() == path
2119 && let Some(broadcast) = broadcast
2120 {
2121 return Some(broadcast.with_stats(scope));
2122 }
2123 }
2124 }
2125
2126 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2132 let prefixes = PathPrefixes::new(prefixes);
2133 Some(Consumer::new(
2134 self.info,
2135 self.root.clone(),
2136 self.nodes.select(&prefixes)?,
2137 self.dynamic.clone(),
2138 self.stats.clone(),
2139 ))
2140 }
2141
2142 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2161 let path = path.as_path();
2162
2163 let absolute = self.root.join(&path).to_owned();
2167 let scope = self.stats.egress(&absolute);
2168
2169 if let Some(broadcast) = self.get_broadcast(&path) {
2171 return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope));
2172 }
2173
2174 let mut state = self.dynamic.lock();
2175
2176 if let Some(weak) = state.served.get(&absolute) {
2180 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2181 }
2182
2183 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2186 producer.consume()
2187 } else {
2188 let producer = kio::Producer::<PendingBroadcast>::default();
2189 let consumer = producer.consume();
2190 if state.requests.insert(absolute, producer).is_err() {
2191 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2192 }
2193 consumer
2194 };
2195
2196 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2197 }
2198
2199 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2204 let prefix = prefix.as_path();
2205
2206 Some(Self::new(
2207 self.info,
2208 self.root.join(&prefix).to_owned(),
2209 self.nodes.root(&prefix)?,
2210 self.dynamic.clone(),
2211 self.stats.clone(),
2212 ))
2213 }
2214
2215 pub fn root(&self) -> &Path<'_> {
2217 &self.root
2218 }
2219
2220 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2223 self.nodes.nodes.iter().map(|(root, _)| root)
2224 }
2225
2226 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2228 self.root.join(path)
2229 }
2230}
2231
2232#[derive(Clone)]
2237pub struct AnnounceProducer {
2238 nodes: OriginNodes,
2239 root: PathOwned,
2240}
2241
2242impl AnnounceProducer {
2243 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2244 Self { nodes, root }
2245 }
2246
2247 pub fn consume(&self) -> AnnounceConsumer {
2253 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2256 }
2257
2258 pub fn root(&self) -> &Path<'_> {
2260 &self.root
2261 }
2262}
2263
2264pub struct AnnounceConsumer {
2269 id: ConsumerId,
2270 nodes: OriginNodes,
2271 root: PathOwned,
2272
2273 state: kio::Producer<OriginConsumerState>,
2276
2277 stats: stats::Session,
2280
2281 guards: HashMap<PathOwned, stats::Announce>,
2285}
2286
2287impl AnnounceConsumer {
2288 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2289 let state = kio::Producer::<OriginConsumerState>::default();
2290 let id = ConsumerId::new();
2291
2292 for (_, node) in &nodes.nodes {
2293 let notify = AnnounceConsumerNotify {
2294 root: root.clone(),
2295 state: state.clone(),
2296 };
2297 node.lock().consume(id, notify);
2298 }
2299
2300 Self {
2301 id,
2302 nodes,
2303 root,
2304 state,
2305 stats,
2306 guards: HashMap::new(),
2307 }
2308 }
2309
2310 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2316 let OriginAnnounce { path, broadcast } = update;
2317 let absolute = self.root.join(&path).to_owned();
2318 match broadcast {
2319 Some(broadcast) => {
2320 let scope = self.stats.egress(&absolute);
2321 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2322 OriginAnnounce {
2323 path,
2324 broadcast: Some(broadcast.with_stats(scope)),
2325 }
2326 }
2327 None => {
2328 self.guards.remove(&absolute);
2329 OriginAnnounce { path, broadcast: None }
2330 }
2331 }
2332 }
2333
2334 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2341 kio::wait(|waiter| self.poll_next(waiter)).await
2342 }
2343
2344 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2350 let update = {
2351 let mut state = match ready!(self.state.poll(waiter, |state| {
2352 if state.pending.is_empty() {
2353 Poll::Pending
2354 } else {
2355 Poll::Ready(())
2356 }
2357 })) {
2358 Ok(state) => state,
2359 Err(_) => return Poll::Ready(None),
2361 };
2362 state.take().expect("predicate guaranteed an update")
2363 };
2364 Poll::Ready(Some(self.attribute(update)))
2365 }
2366
2367 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2372 let update = self.state.write().ok()?.take()?;
2373 Some(self.attribute(update))
2374 }
2375
2376 pub fn is_closed(&self) -> bool {
2378 self.state.write().is_err()
2379 }
2380
2381 pub fn root(&self) -> &Path<'_> {
2383 &self.root
2384 }
2385
2386 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2388 self.root.join(path)
2389 }
2390}
2391
2392impl Drop for AnnounceConsumer {
2393 fn drop(&mut self) {
2394 for (_, root) in &self.nodes.nodes {
2395 root.lock().unconsume(self.id);
2396 }
2397 }
2398}
2399
2400#[cfg(test)]
2401use futures::FutureExt;
2402
2403#[cfg(test)]
2404#[allow(missing_docs)] impl AnnounceConsumer {
2406 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2407 let expected = expected.as_path();
2408 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2409 assert_eq!(announce.path, expected, "wrong path");
2410 let announced = announce.broadcast.expect("should be an active announce");
2411 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2412 }
2413
2414 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2418 let expected = expected.as_path();
2419 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2420 assert_eq!(announce.path, expected, "wrong path");
2421 announce.broadcast.expect("should be an active announce")
2422 }
2423
2424 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2425 let expected = expected.as_path();
2426 let announce = self.try_next().expect("no next");
2427 assert_eq!(announce.path, expected, "wrong path");
2428 let announced = announce.broadcast.expect("should be an active announce");
2429 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2430 }
2431
2432 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2434 let expected = expected.as_path();
2435 let announce = self.try_next().expect("no next");
2436 assert_eq!(announce.path, expected, "wrong path");
2437 announce.broadcast.expect("should be an active announce")
2438 }
2439
2440 pub fn assert_next_none(&mut self, expected: impl AsPath) {
2441 let expected = expected.as_path();
2442 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2443 assert_eq!(announce.path, expected, "wrong path");
2444 assert!(announce.broadcast.is_none(), "should be unannounced");
2445 }
2446
2447 pub fn assert_next_wait(&mut self) {
2448 if let Some(res) = self.next().now_or_never() {
2449 panic!("next should block: got {:?}", res.map(|a| a.path));
2450 }
2451 }
2452
2453 }
2462
2463#[cfg(test)]
2464mod tests {
2465 use crate::coding::Decode;
2466 use crate::group;
2467
2468 use super::*;
2469
2470 fn announce() -> broadcast::Route {
2472 broadcast::Route::new().with_announce(true)
2473 }
2474
2475 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
2481 let name = Path::new(name);
2482 let peer_key = fnv_key(&name, [peer]);
2483 (100u64..)
2484 .map(|id| Origin::new(id).unwrap())
2485 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
2486 .unwrap()
2487 }
2488
2489 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
2492 let source = broadcast::Info::new().produce().consume();
2493 FrontState {
2494 path: Path::new("test").to_owned(),
2495 self_origin,
2496 next_route: routes.len() as u64,
2497 routes: routes
2498 .into_iter()
2499 .enumerate()
2500 .map(|(id, route)| FrontRoute {
2501 id: id as u64,
2502 route,
2503 source: source.clone(),
2504 })
2505 .collect(),
2506 active: Some(0),
2507 linger: Duration::ZERO,
2508 closed: false,
2509 }
2510 }
2511
2512 fn sibling_route(peer: Origin) -> broadcast::Route {
2515 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
2516 announce().with_hops(hops)
2517 }
2518
2519 fn upstream_route(cost: u64) -> broadcast::Route {
2521 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
2522 announce().with_hops(hops).with_cost(cost)
2523 }
2524
2525 #[test]
2529 fn test_carrying_gate_keys() {
2530 let peer = Origin::new(3).unwrap();
2531
2532 let mut lost = front_state(
2534 origin_keyed("test", peer, false),
2535 vec![upstream_route(10), sibling_route(peer)],
2536 );
2537 lost.reselect(true);
2538 assert_eq!(
2539 lost.active,
2540 Some(0),
2541 "carrying front re-parented onto a higher-keyed peer"
2542 );
2543 lost.reselect(false);
2544 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
2545
2546 let mut won = front_state(
2548 origin_keyed("test", peer, true),
2549 vec![upstream_route(10), sibling_route(peer)],
2550 );
2551 won.reselect(true);
2552 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
2553 }
2554
2555 #[test]
2560 fn test_carrying_gate_symmetric_race() {
2561 let a = Origin::new(1).unwrap();
2562 let b = Origin::new(2).unwrap();
2563
2564 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
2565 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
2566 a_view.reselect(true);
2567 b_view.reselect(true);
2568
2569 let a_moved = a_view.active == Some(1);
2570 let b_moved = b_view.active == Some(1);
2571 assert!(
2572 a_moved != b_moved,
2573 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
2574 );
2575 }
2576
2577 #[test]
2582 fn test_carrying_switches_to_benign_routes() {
2583 let peer = Origin::new(3).unwrap();
2584 let lost = origin_keyed("test", peer, false);
2585
2586 let mut forwarder = sibling_route(peer).with_cost(4);
2588 forwarder.advertised = 4;
2589 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
2590 state.reselect(true);
2591 assert_eq!(
2592 state.active,
2593 Some(1),
2594 "a cheaper forwarder path must win while carrying"
2595 );
2596
2597 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
2599 let mut state = front_state(lost, vec![upstream_route(10), direct]);
2600 state.reselect(true);
2601 assert_eq!(
2602 state.active,
2603 Some(1),
2604 "a direct publisher route must win while carrying"
2605 );
2606 }
2607
2608 #[test]
2611 fn test_carrying_gate_ignores_unannounced_incumbent() {
2612 let peer = Origin::new(3).unwrap();
2613 let unannounced = upstream_route(10).with_announce(false);
2614 let mut state = front_state(
2615 origin_keyed("test", peer, false),
2616 vec![unannounced, sibling_route(peer)],
2617 );
2618 state.reselect(true);
2619 assert_eq!(
2620 state.active,
2621 Some(1),
2622 "an unannounced incumbent must always be displaced"
2623 );
2624 }
2625
2626 async fn settle() {
2629 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
2630 }
2631
2632 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
2635 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
2636 .await
2637 .expect("timed out waiting for a track request")
2638 .expect("source closed");
2639 assert_eq!(request.name(), name, "unexpected track dispatched");
2640 request.accept(None)
2641 }
2642
2643 #[tokio::test]
2647 async fn test_stats_tagged_end_to_end() {
2648 use crate::Timestamp;
2649 use crate::stats::{Config, Registry, Tier};
2650 use bytes::Bytes;
2651
2652 tokio::time::pause();
2653
2654 let registry = Registry::new(Config::new());
2655 let ctx = registry.tier(Tier::default()).session("acme");
2656
2657 let origin = Origin::random().produce();
2658 let ingress = origin.clone().with_stats(ctx.clone());
2659 let egress = origin.consume().with_stats(ctx.clone());
2660
2661 let mut announced = egress.announced();
2664
2665 let source = ingress.create_broadcast("demo", announce()).unwrap();
2667 let mut dynamic = source.dynamic();
2668 settle().await;
2669 settle().await;
2670
2671 let update = announced.next().await.unwrap();
2673 assert_eq!(update.path.as_str(), "demo");
2674 let broadcast = update.broadcast.unwrap();
2675
2676 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2678 let mut producer = accept_track(&mut dynamic, "video").await;
2679 settle().await;
2680 let mut sub = subscribing.await.unwrap();
2681
2682 let mut group = producer.append_group().unwrap();
2684 group
2685 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2686 .unwrap();
2687 group
2688 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
2689 .unwrap();
2690 group.finish().unwrap();
2691
2692 let mut group_c = sub.recv_group().await.unwrap().unwrap();
2694 let mut frames = 0;
2695 while let Some(frame) = group_c.read_frame().await.unwrap() {
2696 assert_eq!(frame.payload.len(), 5);
2697 frames += 1;
2698 }
2699 assert_eq!(frames, 2);
2700 settle().await;
2701
2702 let report = registry.report();
2703 let entry = report
2704 .traffic
2705 .iter()
2706 .find(|e| e.path.as_str() == "demo")
2707 .expect("demo tracked");
2708 let path_len = "demo".len() as u64;
2709
2710 let egress = &entry.publisher;
2712 assert_eq!(egress.announced, 1, "one egress announce");
2713 assert_eq!(egress.announced_bytes, path_len);
2714 assert_eq!(egress.subscriptions, 1, "one egress subscription");
2715 assert_eq!(egress.broadcasts, 1, "one viewer");
2716 assert_eq!(egress.groups, 1);
2717 assert_eq!(egress.frames, 2);
2718 assert_eq!(egress.bytes, 10);
2719 assert_eq!(egress.fetches, 0);
2720
2721 let ingress = &entry.subscriber;
2723 assert_eq!(ingress.announced, 1, "one ingress announce");
2724 assert_eq!(ingress.announced_bytes, path_len);
2725 assert_eq!(ingress.subscriptions, 1, "one ingress track");
2726 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
2727 assert_eq!(ingress.groups, 1);
2728 assert_eq!(ingress.frames, 2);
2729 assert_eq!(ingress.bytes, 10);
2730
2731 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
2733 let _ = fetched;
2734 settle().await;
2735 let report = registry.report();
2736 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
2737 assert_eq!(entry.publisher.fetches, 1, "one fetch");
2738 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
2739 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
2740 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
2744 }
2745
2746 #[tokio::test]
2751 async fn test_stats_read_frame_counts_once() {
2752 use crate::Timestamp;
2753 use crate::stats::{Config, Registry, Tier};
2754 use bytes::Bytes;
2755
2756 tokio::time::pause();
2757
2758 let registry = Registry::new(Config::new());
2759 let ctx = registry.tier(Tier::default()).session("acme");
2760
2761 let origin = Origin::random().produce();
2762 let ingress = origin.clone().with_stats(ctx.clone());
2763 let egress = origin.consume().with_stats(ctx.clone());
2764
2765 let mut announced = egress.announced();
2766 let source = ingress.create_broadcast("demo", announce()).unwrap();
2767 let mut dynamic = source.dynamic();
2768 settle().await;
2769 settle().await;
2770
2771 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2772 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2773 let mut producer = accept_track(&mut dynamic, "video").await;
2774 settle().await;
2775 let mut sub = subscribing.await.unwrap();
2776
2777 producer
2779 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2780 .unwrap();
2781
2782 let frame = sub.read_frame().await.unwrap().expect("frame");
2783 assert_eq!(frame.payload.len(), 5);
2784 settle().await;
2785
2786 let report = registry.report();
2787 let entry = report
2788 .traffic
2789 .iter()
2790 .find(|e| e.path.as_str() == "demo")
2791 .expect("demo tracked");
2792 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
2793 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
2794 assert_eq!(
2795 entry.publisher.bytes, 5,
2796 "payload counted once, not zero and not doubled"
2797 );
2798 }
2799
2800 #[tokio::test]
2804 async fn test_stats_datagrams_counted_both_sides() {
2805 use crate::Timestamp;
2806 use crate::stats::{Config, Registry, Tier};
2807
2808 tokio::time::pause();
2809
2810 let registry = Registry::new(Config::new());
2811 let ctx = registry.tier(Tier::default()).session("acme");
2812
2813 let origin = Origin::random().produce();
2814 let ingress = origin.clone().with_stats(ctx.clone());
2815 let egress = origin.consume().with_stats(ctx.clone());
2816
2817 let mut announced = egress.announced();
2818 let source = ingress.create_broadcast("demo", announce()).unwrap();
2819 let mut dynamic = source.dynamic();
2820 settle().await;
2821 settle().await;
2822
2823 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
2824 let subscribing = broadcast.track("video").unwrap().subscribe(None);
2825 let mut producer = accept_track(&mut dynamic, "video").await;
2826 settle().await;
2827 let mut sub = subscribing.await.unwrap();
2828
2829 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
2830 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
2831 assert_eq!(&datagram.payload[..], b"hello");
2832 settle().await;
2833
2834 let report = registry.report();
2835 let entry = report
2836 .traffic
2837 .iter()
2838 .find(|e| e.path.as_str() == "demo")
2839 .expect("demo tracked");
2840
2841 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
2842 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
2843 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
2844 assert_eq!(traffic.frames, 1, "{side}: one frame");
2845 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
2846 }
2847 }
2848
2849 #[test]
2850 fn origin_rejects_reserved_ids() {
2851 assert!(Origin::new(0).is_err());
2852 assert!(Origin::new(1u64 << 62).is_err());
2853 assert_eq!(Origin::new(1).unwrap().id(), 1);
2854
2855 let mut zero = [0u8].as_slice();
2856 assert_eq!(
2857 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
2858 Origin::UNKNOWN
2859 );
2860 }
2861
2862 #[test]
2863 fn origin_list_push_fails_at_limit() {
2864 let mut list = OriginList::new();
2865 for _ in 0..MAX_HOPS {
2866 list.push(Origin::random()).unwrap();
2867 }
2868 assert_eq!(list.len(), MAX_HOPS);
2869 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
2870 }
2871
2872 #[test]
2873 fn origin_list_replace_first() {
2874 let mut list = OriginList::new();
2875 for _ in 0..3 {
2876 list.push(Origin::UNKNOWN).unwrap();
2877 }
2878
2879 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
2881 assert_eq!(
2882 list.as_slice(),
2883 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
2884 );
2885
2886 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
2888 assert_eq!(list.len(), 3);
2889 }
2890
2891 #[test]
2892 fn origin_list_try_from_vec_enforces_limit() {
2893 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
2894 assert!(OriginList::try_from(under).is_ok());
2895
2896 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
2897 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
2898 }
2899
2900 #[tokio::test]
2901 async fn test_announce() {
2902 tokio::time::pause();
2903
2904 let origin = Origin::random().produce();
2905
2906 let mut consumer1 = origin.consume().announced();
2907 consumer1.assert_next_wait();
2908
2909 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
2911 settle().await;
2912
2913 consumer1.assert_next_some("test1");
2914 consumer1.assert_next_wait();
2915
2916 let mut consumer2 = origin.consume().announced();
2919
2920 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
2922 settle().await;
2923
2924 consumer1.assert_next_some("test2");
2925 consumer1.assert_next_wait();
2926
2927 consumer2.assert_next_some("test1");
2928 consumer2.assert_next_some("test2");
2929 consumer2.assert_next_wait();
2930
2931 broadcast1.finish();
2933 settle().await;
2934
2935 consumer1.assert_next_none("test1");
2937 consumer2.assert_next_none("test1");
2938 consumer1.assert_next_wait();
2939 consumer2.assert_next_wait();
2940
2941 let mut consumer3 = origin.consume().announced();
2943 consumer3.assert_next_some("test2");
2944 consumer3.assert_next_wait();
2945
2946 broadcast2.finish();
2947 settle().await;
2948
2949 consumer1.assert_next_none("test2");
2950 consumer2.assert_next_none("test2");
2951 consumer3.assert_next_none("test2");
2952 }
2953
2954 #[tokio::test]
2958 async fn test_duplicate() {
2959 tokio::time::pause();
2960
2961 let origin = Origin::random().produce();
2962 let consumer = origin.consume();
2963 let mut announced = consumer.announced();
2964
2965 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
2966 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
2967 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
2968 settle().await;
2969 assert!(consumer.get_broadcast("test").is_some());
2970
2971 announced.assert_next_some("test");
2972 announced.assert_next_wait();
2973
2974 broadcast2.finish();
2976 settle().await;
2977 assert!(consumer.get_broadcast("test").is_some());
2978 announced.assert_next_wait();
2979
2980 broadcast1.finish();
2982 settle().await;
2983 assert!(consumer.get_broadcast("test").is_some());
2984 announced.assert_next_wait();
2985
2986 broadcast3.finish();
2988 settle().await;
2989 assert!(consumer.get_broadcast("test").is_none());
2990
2991 announced.assert_next_none("test");
2992 announced.assert_next_wait();
2993 }
2994
2995 #[tokio::test]
2998 async fn test_route_failover() {
2999 tokio::time::pause();
3000
3001 let origin = Origin::random().produce();
3002 let consumer = origin.consume();
3003 let mut announced = consumer.announced();
3004
3005 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3006 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3007
3008 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3010 let mut dynamic_a = source_a.dynamic();
3011 settle().await;
3012 settle().await;
3013 let broadcast = consumer.request_broadcast("test").await.unwrap();
3014 announced.assert_next_some("test");
3015
3016 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3018 let mut dynamic_b = source_b.dynamic();
3019 settle().await;
3020 settle().await;
3021 announced.assert_next_wait();
3022
3023 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3025 let mut producer = accept_track(&mut dynamic_a, "video").await;
3026 settle().await;
3027 dynamic_b.assert_no_request();
3028
3029 let mut sub = subscribing.await.unwrap();
3030 sub.assert_no_group();
3033 assert_eq!(producer.subscription().unwrap().group_start, None);
3034
3035 producer.append_group().unwrap();
3036 producer.append_group().unwrap();
3037 assert_eq!(sub.assert_group().sequence, 0);
3038 assert_eq!(sub.assert_group().sequence, 1);
3039
3040 producer.abort(Error::Dropped).unwrap();
3044 source_a.abort(Error::Dropped).unwrap();
3045 drop(dynamic_a);
3046 settle().await;
3047 announced.assert_next_wait();
3048
3049 let mut producer = accept_track(&mut dynamic_b, "video").await;
3052 settle().await;
3053 sub.assert_no_group();
3054 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3055 producer.create_group(group::Info { sequence: 1 }).unwrap();
3056 producer.create_group(group::Info { sequence: 2 }).unwrap();
3057 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3058 sub.assert_not_closed();
3059 }
3060
3061 #[tokio::test]
3064 async fn test_broadcast_route_watch() {
3065 let mut producer = broadcast::Info::new().produce();
3066 let mut consumer = producer.consume();
3067
3068 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3070
3071 producer.set_route(broadcast::Route::default()).unwrap();
3073 assert!(consumer.route_changed().now_or_never().is_none());
3074
3075 let mut hops = OriginList::new();
3076 hops.push(Origin::new(7).unwrap()).unwrap();
3077 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3078 producer.set_route(route.clone()).unwrap();
3079 assert_eq!(consumer.route_changed().await.unwrap(), route);
3080
3081 let mut fresh = producer.consume();
3083 assert_eq!(fresh.route_changed().await.unwrap(), route);
3084
3085 drop(producer);
3086 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3087 }
3088
3089 #[tokio::test]
3093 async fn test_route_cost_update() {
3094 tokio::time::pause();
3095
3096 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3100 let consumer = origin.consume();
3101 let mut announced = consumer.announced();
3102
3103 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3104 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3105
3106 let mut source_a = origin
3108 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3109 .unwrap();
3110 let mut dynamic_a = source_a.dynamic();
3111 settle().await;
3112 let broadcast = consumer.request_broadcast("test").await.unwrap();
3113 announced.assert_next_some("test");
3114
3115 let mut watch = broadcast.clone();
3116 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3117
3118 let mut source_b = origin
3119 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3120 .unwrap();
3121 let mut dynamic_b = source_b.dynamic();
3122 settle().await;
3123 assert!(
3124 watch.route_changed().now_or_never().is_none(),
3125 "a losing standby must not change the advertised route"
3126 );
3127
3128 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3130 let mut producer = accept_track(&mut dynamic_a, "video").await;
3131 settle().await;
3132 let mut sub = subscribing.await.unwrap();
3133 producer.append_group().unwrap();
3134 assert_eq!(sub.assert_group().sequence, 0);
3135
3136 source_a
3139 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3140 .unwrap();
3141 settle().await;
3142 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3143 announced.assert_next_wait();
3144
3145 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3146 settle().await;
3147 sub.assert_no_group();
3150 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3151 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3152 assert_eq!(sub.assert_group().sequence, 1);
3153 sub.assert_not_closed();
3154
3155 source_b
3157 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3158 .unwrap();
3159 settle().await;
3160 let advertised = watch.route_changed().await.unwrap();
3161 assert_eq!(advertised.hops, hops_b);
3162 assert_eq!(advertised.cost, 5);
3163 announced.assert_next_wait();
3164 }
3165
3166 #[tokio::test]
3169 async fn test_completed_track_survives_route_churn() {
3170 tokio::time::pause();
3171
3172 let origin = Origin::random().produce();
3173 let consumer = origin.consume();
3174
3175 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3176 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3177
3178 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3179 let mut dynamic_a = source_a.dynamic();
3180 settle().await;
3181 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3182 let mut dynamic_b = source_b.dynamic();
3183 settle().await;
3184 settle().await;
3185 let broadcast = consumer.request_broadcast("test").await.unwrap();
3186
3187 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3189 let mut producer = accept_track(&mut dynamic_a, "video").await;
3190 settle().await;
3191 let mut sub = subscribing.await.unwrap();
3192 producer.append_group().unwrap();
3193 assert_eq!(sub.assert_group().sequence, 0);
3194 producer.finish().unwrap();
3195 drop(producer);
3196 settle().await;
3197 sub.assert_closed();
3198
3199 source_a.abort(Error::Dropped).unwrap();
3201 drop(dynamic_a);
3202 settle().await;
3203 dynamic_b.assert_no_request();
3204
3205 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3207 late.assert_closed();
3208 }
3209
3210 #[tokio::test]
3213 async fn test_serve_resets_retry_budget() {
3214 tokio::time::pause();
3215
3216 let origin = Origin::random().produce();
3217 let consumer = origin.consume();
3218
3219 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3220 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3221 let mut dynamic = source.dynamic();
3222 settle().await;
3223 settle().await;
3224 let broadcast = consumer.request_broadcast("test").await.unwrap();
3225
3226 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3228
3229 for _ in 0..2 * MAX_TRACK_RETRIES {
3232 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3233 .await
3234 .expect("timed out waiting for a retry")
3235 .unwrap();
3236 request.reject(Error::NotFound);
3237 let producer = accept_track(&mut dynamic, "video").await;
3238 settle().await;
3239 drop(producer);
3240 }
3241
3242 let _producer = accept_track(&mut dynamic, "video").await;
3243 settle().await;
3244 let mut sub = subscribing.await.unwrap();
3245 sub.assert_not_closed();
3246 }
3247
3248 #[tokio::test]
3252 async fn test_route_handover() {
3253 tokio::time::pause();
3254
3255 let origin = Origin::random().produce();
3256 let consumer = origin.consume();
3257 let mut announced = consumer.announced();
3258
3259 let hops_long = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
3260 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3261
3262 let source_a = origin
3263 .create_broadcast("test", announce().with_hops(hops_long))
3264 .unwrap();
3265 let mut dynamic_a = source_a.dynamic();
3266 settle().await;
3267 settle().await;
3268 let broadcast = consumer.request_broadcast("test").await.unwrap();
3269 announced.assert_next_some("test");
3270
3271 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3272 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3273 settle().await;
3274 let mut sub = subscribing.await.unwrap();
3275 producer_a.append_group().unwrap();
3276 producer_a.append_group().unwrap();
3277 assert_eq!(sub.assert_group().sequence, 0);
3278 assert_eq!(sub.assert_group().sequence, 1);
3279
3280 let source_b = origin
3283 .create_broadcast("test", announce().with_hops(hops_short))
3284 .unwrap();
3285 let mut dynamic_b = source_b.dynamic();
3286 settle().await;
3287 settle().await;
3288 announced.assert_next_wait();
3289
3290 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3291 settle().await;
3292
3293 sub.assert_no_group();
3296 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3297 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3298
3299 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3301 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3302 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3303 assert_eq!(sub.assert_group().sequence, 2);
3304 assert_eq!(sub.assert_group().sequence, 3);
3305 sub.assert_no_group();
3306 sub.assert_not_closed();
3307 }
3308
3309 #[tokio::test(start_paused = true)]
3312 async fn test_route_unannounce_immediate() {
3313 let origin = Origin::random().produce();
3314 let consumer = origin.consume();
3315 let mut announced = consumer.announced();
3316
3317 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3318 let mut source = origin
3319 .create_broadcast("test", announce().with_hops(hops.clone()))
3320 .unwrap();
3321 settle().await;
3322 let broadcast = consumer.request_broadcast("test").await.unwrap();
3323 announced.assert_next_some("test");
3324
3325 source.finish();
3328 settle().await;
3329 announced.assert_next_none("test");
3330
3331 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3333 settle().await;
3334 let fresh = consumer.request_broadcast("test").await.unwrap();
3335 announced.assert_next_some("test");
3336 assert!(
3337 !fresh.is_clone(&broadcast),
3338 "re-create must not splice the old broadcast"
3339 );
3340 }
3341
3342 #[tokio::test(start_paused = true)]
3347 async fn test_route_detach_immediate() {
3348 let origin = Origin::random().produce();
3349 let consumer = origin.consume();
3350 let mut announced = consumer.announced();
3351
3352 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3353 let source = origin
3354 .create_broadcast("test", announce().with_hops(hops.clone()))
3355 .unwrap();
3356 let mut dynamic = source.dynamic();
3357 settle().await;
3358 settle().await;
3359 let broadcast = consumer.request_broadcast("test").await.unwrap();
3360 announced.assert_next_some("test");
3361
3362 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3363 let producer = accept_track(&mut dynamic, "video").await;
3364 settle().await;
3365 let mut sub = subscribing.await.unwrap();
3366
3367 drop(producer);
3369 source.abort(Error::Dropped).unwrap();
3370 drop(dynamic);
3371
3372 settle().await;
3373 announced.assert_next_none("test");
3374 sub.assert_error();
3375
3376 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3379 settle().await;
3380 settle().await;
3381 let fresh = consumer.request_broadcast("test").await.unwrap();
3382 announced.assert_next_some("test");
3383 assert!(
3384 !fresh.is_clone(&broadcast),
3385 "re-create must not splice the old broadcast"
3386 );
3387 }
3388
3389 #[tokio::test(start_paused = true)]
3393 async fn test_linger_reconnect_splices() {
3394 let origin = Info::new(Origin::random())
3395 .with_linger(Duration::from_secs(5))
3396 .produce();
3397 let consumer = origin.consume();
3398 let mut announced = consumer.announced();
3399
3400 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3401 let source = origin
3402 .create_broadcast("test", announce().with_hops(hops.clone()))
3403 .unwrap();
3404 let mut dynamic = source.dynamic();
3405 settle().await;
3406 settle().await;
3407 let broadcast = consumer.request_broadcast("test").await.unwrap();
3408 announced.assert_next_some("test");
3409
3410 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3411 let mut producer = accept_track(&mut dynamic, "video").await;
3412 settle().await;
3413 let mut sub = subscribing.await.unwrap();
3414
3415 producer.append_group().unwrap();
3416 producer.append_group().unwrap();
3417 assert_eq!(sub.assert_group().sequence, 0);
3418 assert_eq!(sub.assert_group().sequence, 1);
3419
3420 drop(producer);
3423 source.abort(Error::Dropped).unwrap();
3424 drop(dynamic);
3425 settle().await;
3426
3427 announced.assert_next_wait();
3429 sub.assert_no_group();
3430 sub.assert_not_closed();
3431
3432 let during = consumer.request_broadcast("test").await.unwrap();
3434 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
3435
3436 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3439 let mut dynamic = source.dynamic();
3440 settle().await;
3441 settle().await;
3442 announced.assert_next_wait();
3443 let again = consumer.request_broadcast("test").await.unwrap();
3444 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3445
3446 let mut producer = accept_track(&mut dynamic, "video").await;
3450 settle().await;
3451 sub.assert_no_group();
3452 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3453 producer.create_group(group::Info { sequence: 2 }).unwrap();
3454 assert_eq!(sub.assert_group().sequence, 2);
3455 sub.assert_not_closed();
3456 }
3457
3458 #[tokio::test(start_paused = true)]
3461 async fn test_linger_expiry_closes() {
3462 let origin = Info::new(Origin::random())
3463 .with_linger(Duration::from_secs(5))
3464 .produce();
3465 let consumer = origin.consume();
3466 let mut announced = consumer.announced();
3467
3468 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3469 let source = origin
3470 .create_broadcast("test", announce().with_hops(hops.clone()))
3471 .unwrap();
3472 let mut dynamic = source.dynamic();
3473 settle().await;
3474 settle().await;
3475 let broadcast = consumer.request_broadcast("test").await.unwrap();
3476 announced.assert_next_some("test");
3477
3478 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3479 let producer = accept_track(&mut dynamic, "video").await;
3480 settle().await;
3481 let mut sub = subscribing.await.unwrap();
3482
3483 drop(producer);
3484 source.abort(Error::Dropped).unwrap();
3485 drop(dynamic);
3486 settle().await;
3487 announced.assert_next_wait();
3488
3489 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
3491 settle().await;
3492 announced.assert_next_none("test");
3493 sub.assert_error();
3494
3495 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3497 settle().await;
3498 settle().await;
3499 let fresh = consumer.request_broadcast("test").await.unwrap();
3500 announced.assert_next_some("test");
3501 assert!(
3502 !fresh.is_clone(&broadcast),
3503 "a late re-create must not splice the expired broadcast"
3504 );
3505 }
3506
3507 #[tokio::test(start_paused = true)]
3511 async fn test_linger_forever() {
3512 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
3513 let consumer = origin.consume();
3514 let mut announced = consumer.announced();
3515
3516 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3517 let source = origin
3518 .create_broadcast("test", announce().with_hops(hops.clone()))
3519 .unwrap();
3520 settle().await;
3521 let broadcast = consumer.request_broadcast("test").await.unwrap();
3522 announced.assert_next_some("test");
3523
3524 source.abort(Error::Dropped).unwrap();
3525 settle().await;
3526
3527 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
3529 announced.assert_next_wait();
3530 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3531 settle().await;
3532 settle().await;
3533 let again = consumer.request_broadcast("test").await.unwrap();
3534 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
3535 drop(source);
3536 }
3537
3538 #[tokio::test(start_paused = true)]
3541 async fn test_linger_skipped_on_finish() {
3542 let origin = Info::new(Origin::random())
3543 .with_linger(Duration::from_secs(5))
3544 .produce();
3545 let consumer = origin.consume();
3546 let mut announced = consumer.announced();
3547
3548 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3549 let mut source = origin
3550 .create_broadcast("test", announce().with_hops(hops.clone()))
3551 .unwrap();
3552 settle().await;
3553 let broadcast = consumer.request_broadcast("test").await.unwrap();
3554 announced.assert_next_some("test");
3555
3556 source.finish();
3559 settle().await;
3560 announced.assert_next_none("test");
3561
3562 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3564 settle().await;
3565 let fresh = consumer.request_broadcast("test").await.unwrap();
3566 announced.assert_next_some("test");
3567 assert!(
3568 !fresh.is_clone(&broadcast),
3569 "a finish must not leave a lingering broadcast to splice into"
3570 );
3571 }
3572
3573 #[tokio::test]
3576 async fn test_announce_toggle() {
3577 tokio::time::pause();
3578
3579 let origin = Origin::random().produce();
3580 let consumer = origin.consume();
3581 let mut announced = consumer.announced();
3582
3583 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3584 settle().await;
3585
3586 announced.assert_next_wait();
3588 let broadcast = consumer
3589 .get_broadcast("test")
3590 .expect("offline broadcast is still routable");
3591 assert!(!broadcast.route().announce);
3592
3593 let requested = consumer.request_broadcast("test").await.unwrap();
3595 assert!(requested.is_clone(&broadcast));
3596
3597 source.set_route(announce()).unwrap();
3599 settle().await;
3600 let face = announced.assert_next_some("test");
3601 assert!(face.is_clone(&broadcast));
3602
3603 let mut fresh = origin.consume().announced();
3605 fresh.assert_next_some("test");
3606 fresh.assert_next_wait();
3607
3608 source.set_route(broadcast::Route::new()).unwrap();
3610 settle().await;
3611 announced.assert_next_none("test");
3612 assert!(consumer.get_broadcast("test").is_some());
3613 let mut fresh = origin.consume().announced();
3614 fresh.assert_next_wait();
3615
3616 source.finish();
3617 settle().await;
3618 assert!(consumer.get_broadcast("test").is_none());
3619 }
3620
3621 #[tokio::test]
3624 async fn test_announce_beats_offline() {
3625 tokio::time::pause();
3626
3627 let origin = Origin::random().produce();
3628 let consumer = origin.consume();
3629 let mut announced = consumer.announced();
3630
3631 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
3633 settle().await;
3634 announced.assert_next_wait();
3635
3636 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
3639 settle().await;
3640 announced.assert_next_some("test");
3641 let face = consumer.get_broadcast("test").unwrap();
3642 assert!(face.route().announce);
3643 assert_eq!(face.route().cost, 10);
3644
3645 announced_source.finish();
3648 settle().await;
3649 announced.assert_next_none("test");
3650 assert!(consumer.get_broadcast("test").is_some());
3651 }
3652
3653 #[tokio::test]
3656 async fn test_better_source_no_churn() {
3657 tokio::time::pause();
3658
3659 let origin = Origin::random().produce();
3660 let mut announced = origin.consume().announced();
3661
3662 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3664 let _a = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3665 settle().await;
3666 let face = announced.assert_next_some("test");
3667
3668 let _b = origin.create_broadcast("test", announce()).unwrap();
3669 settle().await;
3670 announced.assert_next_wait();
3671 let current = origin.consume().get_broadcast("test").unwrap();
3672 assert!(current.is_clone(&face), "the broadcast identity must not change");
3673 assert!(current.route().hops.is_empty());
3675 }
3676
3677 #[tokio::test]
3678 async fn test_duplicate_reverse() {
3679 tokio::time::pause();
3680
3681 let origin = Origin::random().produce();
3682
3683 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3684 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3685 settle().await;
3686 assert!(origin.consume().get_broadcast("test").is_some());
3687
3688 broadcast2.finish();
3690 settle().await;
3691 assert!(origin.consume().get_broadcast("test").is_some());
3692
3693 broadcast1.finish();
3694 settle().await;
3695 assert!(origin.consume().get_broadcast("test").is_none());
3696 }
3697
3698 #[tokio::test]
3699 async fn test_deterministic_tiebreak() {
3700 tokio::time::pause();
3701
3702 fn hops(ids: &[u64]) -> OriginList {
3703 OriginList::try_from(
3704 ids.iter()
3705 .copied()
3706 .map(|id| Origin::new(id).unwrap())
3707 .collect::<Vec<_>>(),
3708 )
3709 .unwrap()
3710 }
3711
3712 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
3715 let origin = Origin::random().produce();
3716 let _a = origin
3717 .create_broadcast("test", announce().with_hops(hops(first)))
3718 .unwrap();
3719 let _b = origin
3720 .create_broadcast("test", announce().with_hops(hops(second)))
3721 .unwrap();
3722 settle().await;
3723 origin.consume().get_broadcast("test").unwrap().route().hops
3724 }
3725
3726 let forward = winner(&[10, 20], &[30, 40]).await;
3729 let reverse = winner(&[30, 40], &[10, 20]).await;
3730 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
3731
3732 assert_eq!(winner(&[10, 20], &[30]).await.len(), 1);
3734 assert_eq!(winner(&[30], &[10, 20]).await.len(), 1);
3735 }
3736
3737 #[tokio::test]
3742 async fn test_many_announces() {
3743 let origin = Origin::random().produce();
3744
3745 let mut consumer = origin.consume().announced();
3746 let mut broadcasts = Vec::new();
3748 for i in 0..256 {
3749 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3750 settle().await;
3751 }
3752
3753 for i in 0..256 {
3754 consumer.assert_next_some(format!("test{i:03}"));
3755 }
3756 consumer.assert_next_wait();
3757 }
3758
3759 #[tokio::test]
3760 async fn test_many_announces_try() {
3761 let origin = Origin::random().produce();
3762
3763 let mut consumer = origin.consume().announced();
3764 let mut broadcasts = Vec::new();
3766 for i in 0..256 {
3767 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
3768 settle().await;
3769 }
3770
3771 for i in 0..256 {
3772 consumer.assert_try_next_some(format!("test{i:03}"));
3773 }
3774 }
3775
3776 #[tokio::test]
3777 async fn test_with_root_basic() {
3778 let origin = Origin::random().produce();
3779
3780 let foo_producer = origin.with_root("foo").expect("should create root");
3782 assert_eq!(foo_producer.root().as_str(), "foo");
3783
3784 let mut consumer = origin.consume().announced();
3785
3786 let _broadcast = foo_producer
3788 .create_broadcast("bar/baz", announce())
3789 .expect("publish allowed");
3790 settle().await;
3791 consumer.assert_next_some("foo/bar/baz");
3793
3794 let mut foo_consumer = foo_producer.consume().announced();
3796 foo_consumer.assert_next_some("bar/baz");
3797 }
3798
3799 #[tokio::test]
3800 async fn test_with_root_nested() {
3801 let origin = Origin::random().produce();
3802
3803 let foo_producer = origin.with_root("foo").expect("should create foo root");
3805 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
3806 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
3807
3808 let mut consumer = origin.consume().announced();
3809
3810 let _broadcast = foo_bar_producer
3812 .create_broadcast("baz", announce())
3813 .expect("publish allowed");
3814 settle().await;
3815 consumer.assert_next_some("foo/bar/baz");
3817
3818 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
3820 foo_bar_consumer.assert_next_some("baz");
3821 }
3822
3823 #[tokio::test]
3824 async fn test_publish_scope_allows() {
3825 let origin = Origin::random().produce();
3826
3827 let limited_producer = origin
3829 .scope(&["allowed/path1".into(), "allowed/path2".into()])
3830 .expect("should create limited producer");
3831
3832 let _broadcast = limited_producer
3834 .create_broadcast("allowed/path1", announce())
3835 .expect("publish allowed");
3836 let _keep2 = limited_producer
3837 .create_broadcast("allowed/path1/nested", announce())
3838 .expect("publish allowed");
3839 let _keep3 = limited_producer
3840 .create_broadcast("allowed/path2", announce())
3841 .expect("publish allowed");
3842 settle().await;
3843
3844 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
3846 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
3848 }
3849
3850 #[tokio::test]
3851 async fn test_publish_max_parts() {
3852 let origin = Origin::random().produce();
3853
3854 let at_limit = (0..Path::MAX_PARTS)
3855 .map(|i| i.to_string())
3856 .collect::<Vec<_>>()
3857 .join("/");
3858 let _broadcast = origin
3859 .create_broadcast(at_limit.as_str(), announce())
3860 .expect("publish allowed");
3861 settle().await;
3862
3863 let too_deep = format!("{at_limit}/extra");
3864 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
3865
3866 let rooted = origin.with_root("root").expect("wildcard allows any root");
3868 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
3869 }
3870
3871 #[tokio::test]
3872 async fn test_publish_scope_empty() {
3873 let origin = Origin::random().produce();
3874
3875 assert!(origin.scope(&[]).is_none());
3877 }
3878
3879 #[tokio::test]
3880 async fn test_consume_scope_filters() {
3881 let origin = Origin::random().produce();
3882
3883 let mut consumer = origin.consume().announced();
3884
3885 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
3887 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
3888 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
3889 settle().await;
3890
3891 let mut limited_consumer = origin
3893 .consume()
3894 .scope(&["allowed".into()])
3895 .expect("should create limited consumer")
3896 .announced();
3897
3898 limited_consumer.assert_next_some("allowed");
3900 limited_consumer.assert_next_some("allowed/nested");
3901 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
3905 consumer.assert_next_some("allowed/nested");
3906 consumer.assert_next_some("notallowed");
3907 }
3908
3909 #[tokio::test]
3910 async fn test_consume_scope_multiple_prefixes() {
3911 let origin = Origin::random().produce();
3912
3913 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
3914 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
3915 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
3916 settle().await;
3917
3918 let mut limited_consumer = origin
3920 .consume()
3921 .scope(&["foo".into(), "bar".into()])
3922 .expect("should create limited consumer")
3923 .announced();
3924
3925 limited_consumer.assert_next_some("bar/test");
3927 limited_consumer.assert_next_some("foo/test");
3928 limited_consumer.assert_next_wait(); }
3930
3931 #[tokio::test]
3932 async fn test_with_root_and_publish_scope() {
3933 let origin = Origin::random().produce();
3934
3935 let foo_producer = origin.with_root("foo").expect("should create foo root");
3937
3938 let limited_producer = foo_producer
3940 .scope(&["bar".into(), "goop/pee".into()])
3941 .expect("should create limited producer");
3942
3943 let mut consumer = origin.consume().announced();
3944
3945 let _broadcast = limited_producer
3947 .create_broadcast("bar", announce())
3948 .expect("publish allowed");
3949 let _keep2 = limited_producer
3950 .create_broadcast("bar/nested", announce())
3951 .expect("publish allowed");
3952 let _keep3 = limited_producer
3953 .create_broadcast("goop/pee", announce())
3954 .expect("publish allowed");
3955 let _keep4 = limited_producer
3956 .create_broadcast("goop/pee/nested", announce())
3957 .expect("publish allowed");
3958 settle().await;
3959
3960 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
3962 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
3964
3965 consumer.assert_next_some("foo/bar");
3967 consumer.assert_next_some("foo/bar/nested");
3968 consumer.assert_next_some("foo/goop/pee");
3969 consumer.assert_next_some("foo/goop/pee/nested");
3970 }
3971
3972 #[tokio::test]
3973 async fn test_with_root_and_consume_scope() {
3974 let origin = Origin::random().produce();
3975
3976 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
3978 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
3979 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
3980 settle().await;
3981
3982 let foo_producer = origin.with_root("foo").expect("should create foo root");
3984
3985 let mut limited_consumer = foo_producer
3987 .consume()
3988 .scope(&["bar".into(), "goop/pee".into()])
3989 .expect("should create limited consumer")
3990 .announced();
3991
3992 limited_consumer.assert_next_some("bar/test");
3994 limited_consumer.assert_next_some("goop/pee/test");
3995 limited_consumer.assert_next_wait(); }
3997
3998 #[tokio::test]
3999 async fn test_with_root_unauthorized() {
4000 let origin = Origin::random().produce();
4001
4002 let limited_producer = origin
4004 .scope(&["allowed".into()])
4005 .expect("should create limited producer");
4006
4007 assert!(limited_producer.with_root("notallowed").is_none());
4009
4010 let allowed_root = limited_producer
4012 .with_root("allowed")
4013 .expect("should create allowed root");
4014 assert_eq!(allowed_root.root().as_str(), "allowed");
4015 }
4016
4017 #[tokio::test]
4018 async fn test_wildcard_permission() {
4019 let origin = Origin::random().produce();
4020
4021 let root_producer = origin.clone();
4023
4024 let _broadcast = root_producer
4026 .create_broadcast("any/path", announce())
4027 .expect("publish allowed");
4028 let _keep2 = root_producer
4029 .create_broadcast("other/path", announce())
4030 .expect("publish allowed");
4031 settle().await;
4032
4033 let foo_producer = root_producer.with_root("foo").expect("should create any root");
4035 assert_eq!(foo_producer.root().as_str(), "foo");
4036 }
4037
4038 #[tokio::test]
4039 async fn test_consume_broadcast_with_permissions() {
4040 let origin = Origin::random().produce();
4041
4042 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
4043 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
4044 settle().await;
4045
4046 let limited_consumer = origin
4048 .consume()
4049 .scope(&["allowed".into()])
4050 .expect("should create limited consumer");
4051
4052 let result = limited_consumer.get_broadcast("allowed/test");
4054 assert!(result.is_some());
4055 assert!(
4056 result
4057 .unwrap()
4058 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
4059 );
4060
4061 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
4063
4064 let consumer = origin.consume();
4066 assert!(consumer.get_broadcast("allowed/test").is_some());
4067 assert!(consumer.get_broadcast("notallowed/test").is_some());
4068 }
4069
4070 #[tokio::test]
4071 async fn test_nested_paths_with_permissions() {
4072 let origin = Origin::random().produce();
4073
4074 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
4076
4077 let _broadcast = limited_producer
4079 .create_broadcast("a/b/c", announce())
4080 .expect("publish allowed");
4081 let _keep2 = limited_producer
4082 .create_broadcast("a/b/c/d", announce())
4083 .expect("publish allowed");
4084 let _keep3 = limited_producer
4085 .create_broadcast("a/b/c/d/e", announce())
4086 .expect("publish allowed");
4087 settle().await;
4088
4089 assert!(limited_producer.create_broadcast("a", announce()).is_err());
4091 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
4092 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
4093 }
4094
4095 #[tokio::test]
4096 async fn test_multiple_consumers_with_different_permissions() {
4097 let origin = Origin::random().produce();
4098
4099 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
4101 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
4102 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
4103 settle().await;
4104
4105 let mut foo_consumer = origin
4107 .consume()
4108 .scope(&["foo".into()])
4109 .expect("should create foo consumer")
4110 .announced();
4111
4112 let mut bar_consumer = origin
4113 .consume()
4114 .scope(&["bar".into()])
4115 .expect("should create bar consumer")
4116 .announced();
4117
4118 let mut foobar_consumer = origin
4119 .consume()
4120 .scope(&["foo".into(), "bar".into()])
4121 .expect("should create foobar consumer")
4122 .announced();
4123
4124 foo_consumer.assert_next_some("foo/test");
4126 foo_consumer.assert_next_wait();
4127
4128 bar_consumer.assert_next_some("bar/test");
4129 bar_consumer.assert_next_wait();
4130
4131 foobar_consumer.assert_next_some("bar/test");
4132 foobar_consumer.assert_next_some("foo/test");
4133 foobar_consumer.assert_next_wait();
4134 }
4135
4136 #[tokio::test]
4137 async fn test_select_with_empty_prefix() {
4138 let origin = Origin::random().produce();
4139
4140 let demo_producer = origin.with_root("demo").expect("should create demo root");
4142 let limited_producer = demo_producer
4143 .scope(&["worm-node".into(), "foobar".into()])
4144 .expect("should create limited producer");
4145
4146 let _broadcast1 = limited_producer
4148 .create_broadcast("worm-node/test", announce())
4149 .expect("publish allowed");
4150 let _broadcast2 = limited_producer
4151 .create_broadcast("foobar/test", announce())
4152 .expect("publish allowed");
4153 settle().await;
4154
4155 let mut consumer = limited_producer
4157 .consume()
4158 .scope(&["".into()])
4159 .expect("should create consumer with empty prefix")
4160 .announced();
4161
4162 let a1 = consumer.try_next().expect("expected first announcement");
4164 let a2 = consumer.try_next().expect("expected second announcement");
4165 consumer.assert_next_wait();
4166
4167 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4168 paths.sort();
4169 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
4170 }
4171
4172 #[tokio::test]
4173 async fn test_select_narrowing_scope() {
4174 let origin = Origin::random().produce();
4175
4176 let demo_producer = origin.with_root("demo").expect("should create demo root");
4178 let limited_producer = demo_producer
4179 .scope(&["worm-node".into(), "foobar".into()])
4180 .expect("should create limited producer");
4181
4182 let _broadcast1 = limited_producer
4184 .create_broadcast("worm-node", announce())
4185 .expect("publish allowed");
4186 let _broadcast2 = limited_producer
4187 .create_broadcast("worm-node/foo", announce())
4188 .expect("publish allowed");
4189 let _broadcast3 = limited_producer
4190 .create_broadcast("foobar/bar", announce())
4191 .expect("publish allowed");
4192 settle().await;
4193
4194 let mut worm_consumer = limited_producer
4196 .consume()
4197 .scope(&["worm-node".into()])
4198 .expect("should create worm-node consumer")
4199 .announced();
4200
4201 worm_consumer.assert_next_some("worm-node");
4203 worm_consumer.assert_next_some("worm-node/foo");
4204 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
4208 .consume()
4209 .scope(&["worm-node/foo".into()])
4210 .expect("should create worm-node/foo consumer")
4211 .announced();
4212
4213 foo_consumer.assert_next_some("worm-node/foo");
4214 foo_consumer.assert_next_wait(); }
4216
4217 #[tokio::test]
4218 async fn test_select_multiple_roots_with_empty_prefix() {
4219 let origin = Origin::random().produce();
4220
4221 let limited_producer = origin
4223 .scope(&["app1".into(), "app2".into(), "shared".into()])
4224 .expect("should create limited producer");
4225
4226 let _broadcast1 = limited_producer
4228 .create_broadcast("app1/data", announce())
4229 .expect("publish allowed");
4230 let _broadcast2 = limited_producer
4231 .create_broadcast("app2/config", announce())
4232 .expect("publish allowed");
4233 let _broadcast3 = limited_producer
4234 .create_broadcast("shared/resource", announce())
4235 .expect("publish allowed");
4236 settle().await;
4237
4238 let mut consumer = limited_producer
4240 .consume()
4241 .scope(&["".into()])
4242 .expect("should create consumer with empty prefix")
4243 .announced();
4244
4245 consumer.assert_next_some("app1/data");
4247 consumer.assert_next_some("app2/config");
4248 consumer.assert_next_some("shared/resource");
4249 consumer.assert_next_wait();
4250 }
4251
4252 #[tokio::test]
4253 async fn test_publish_scope_with_empty_prefix() {
4254 let origin = Origin::random().produce();
4255
4256 let limited_producer = origin
4258 .scope(&["services/api".into(), "services/web".into()])
4259 .expect("should create limited producer");
4260
4261 let same_producer = limited_producer
4263 .scope(&["".into()])
4264 .expect("should create producer with empty prefix");
4265
4266 let _broadcast = same_producer
4268 .create_broadcast("services/api", announce())
4269 .expect("publish allowed");
4270 let _keep2 = same_producer
4271 .create_broadcast("services/web", announce())
4272 .expect("publish allowed");
4273 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
4274 assert!(same_producer.create_broadcast("other", announce()).is_err());
4275 }
4276
4277 #[tokio::test]
4278 async fn test_select_narrowing_to_deeper_path() {
4279 let origin = Origin::random().produce();
4280
4281 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
4283
4284 let _broadcast1 = limited_producer
4286 .create_broadcast("org/team1/project1", announce())
4287 .expect("publish allowed");
4288 let _broadcast2 = limited_producer
4289 .create_broadcast("org/team1/project2", announce())
4290 .expect("publish allowed");
4291 let _broadcast3 = limited_producer
4292 .create_broadcast("org/team2/project1", announce())
4293 .expect("publish allowed");
4294 settle().await;
4295
4296 let mut team2_consumer = limited_producer
4298 .consume()
4299 .scope(&["org/team2".into()])
4300 .expect("should create team2 consumer")
4301 .announced();
4302
4303 team2_consumer.assert_next_some("org/team2/project1");
4304 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
4308 .consume()
4309 .scope(&["org/team1/project1".into()])
4310 .expect("should create project1 consumer")
4311 .announced();
4312
4313 project1_consumer.assert_next_some("org/team1/project1");
4315 project1_consumer.assert_next_wait();
4316 }
4317
4318 #[tokio::test]
4319 async fn test_select_with_non_matching_prefix() {
4320 let origin = Origin::random().produce();
4321
4322 let limited_producer = origin
4324 .scope(&["allowed/path".into()])
4325 .expect("should create limited producer");
4326
4327 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
4329
4330 assert!(limited_producer.scope(&["other/path".into()]).is_none());
4332 }
4333
4334 #[tokio::test]
4337 async fn test_with_root_trailing_slash_consumer() {
4338 let origin = Origin::random().produce();
4339
4340 let prefix = "some_prefix/".to_string();
4342 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4343
4344 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4345 settle().await;
4346 consumer.assert_next_some("test");
4347 }
4348
4349 #[tokio::test]
4351 async fn test_with_root_trailing_slash_producer() {
4352 let origin = Origin::random().produce();
4353
4354 let prefix = "some_prefix/".to_string();
4356 let rooted = origin.with_root(prefix).unwrap();
4357
4358 let _b = rooted.create_broadcast("test", announce()).unwrap();
4359 settle().await;
4360
4361 let mut consumer = rooted.consume().announced();
4362 consumer.assert_next_some("test");
4363 }
4364
4365 #[tokio::test]
4367 async fn test_with_root_trailing_slash_unannounce() {
4368 tokio::time::pause();
4369
4370 let origin = Origin::random().produce();
4371
4372 let prefix = "some_prefix/".to_string();
4373 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
4374
4375 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
4376 settle().await;
4377 consumer.assert_next_some("test");
4378
4379 b.finish();
4381 settle().await;
4382
4383 consumer.assert_next_none("test");
4385 }
4386
4387 #[tokio::test]
4388 async fn test_select_maintains_access_with_wider_prefix() {
4389 let origin = Origin::random().produce();
4390
4391 let demo_producer = origin.with_root("demo").expect("should create demo root");
4393 let user_producer = demo_producer
4394 .scope(&["worm-node".into(), "foobar".into()])
4395 .expect("should create user producer");
4396
4397 let _broadcast1 = user_producer
4399 .create_broadcast("worm-node/data", announce())
4400 .expect("publish allowed");
4401 let _broadcast2 = user_producer
4402 .create_broadcast("foobar", announce())
4403 .expect("publish allowed");
4404 settle().await;
4405
4406 let mut consumer = user_producer
4408 .consume()
4409 .scope(&["".into()])
4410 .expect("scope with empty prefix should not fail when user has specific permissions")
4411 .announced();
4412
4413 let a1 = consumer.try_next().expect("expected first announcement");
4415 let a2 = consumer.try_next().expect("expected second announcement");
4416 consumer.assert_next_wait();
4417
4418 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4419 paths.sort();
4420 assert_eq!(paths, ["foobar", "worm-node/data"]);
4421
4422 let mut narrow_consumer = user_producer
4424 .consume()
4425 .scope(&["worm-node".into()])
4426 .expect("should be able to narrow scope to worm-node")
4427 .announced();
4428
4429 narrow_consumer.assert_next_some("worm-node/data");
4430 narrow_consumer.assert_next_wait(); }
4432
4433 #[tokio::test]
4434 async fn test_duplicate_prefixes_deduped() {
4435 let origin = Origin::random().produce();
4436
4437 let producer = origin
4439 .scope(&["demo".into(), "demo".into()])
4440 .expect("should create producer");
4441
4442 let _broadcast = producer
4443 .create_broadcast("demo/stream", announce())
4444 .expect("publish allowed");
4445 settle().await;
4446
4447 let mut consumer = producer.consume().announced();
4448 consumer.assert_next_some("demo/stream");
4449 consumer.assert_next_wait();
4450 }
4451
4452 #[tokio::test]
4453 async fn test_overlapping_prefixes_deduped() {
4454 let origin = Origin::random().produce();
4455
4456 let producer = origin
4458 .scope(&["demo".into(), "demo/foo".into()])
4459 .expect("should create producer");
4460
4461 let _broadcast = producer
4463 .create_broadcast("demo/bar/stream", announce())
4464 .expect("publish allowed");
4465 settle().await;
4466
4467 let mut consumer = producer.consume().announced();
4468 consumer.assert_next_some("demo/bar/stream");
4469 consumer.assert_next_wait();
4470 }
4471
4472 #[tokio::test]
4473 async fn test_overlapping_prefixes_no_duplicate_announcements() {
4474 let origin = Origin::random().produce();
4475
4476 let producer = origin
4478 .scope(&["demo".into(), "demo/foo".into()])
4479 .expect("should create producer");
4480
4481 let _broadcast = producer
4482 .create_broadcast("demo/foo/stream", announce())
4483 .expect("publish allowed");
4484 settle().await;
4485
4486 let mut consumer = producer.consume().announced();
4487 consumer.assert_next_some("demo/foo/stream");
4489 consumer.assert_next_wait();
4490 }
4491
4492 #[tokio::test]
4493 async fn test_allowed_returns_deduped_prefixes() {
4494 let origin = Origin::random().produce();
4495
4496 let producer = origin
4497 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
4498 .expect("should create producer");
4499
4500 let allowed: Vec<_> = producer.allowed().collect();
4501 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
4502 }
4503
4504 #[tokio::test]
4505 async fn test_announced_broadcast_already_announced() {
4506 let origin = Origin::random().produce();
4507
4508 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4509 settle().await;
4510
4511 let consumer = origin.consume();
4512 let result = consumer.announced_broadcast("test").await.expect("should find it");
4513 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4514 }
4515
4516 #[tokio::test]
4517 async fn test_announced_broadcast_delayed() {
4518 tokio::time::pause();
4519
4520 let origin = Origin::random().produce();
4521
4522 let consumer = origin.consume();
4523
4524 let wait = tokio::spawn({
4526 let consumer = consumer.clone();
4527 async move { consumer.announced_broadcast("test").await }
4528 });
4529
4530 tokio::task::yield_now().await;
4532
4533 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
4534 settle().await;
4535
4536 let result = wait.await.unwrap().expect("should find it");
4537 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
4538 }
4539
4540 #[tokio::test]
4541 async fn test_announced_broadcast_ignores_unrelated_paths() {
4542 tokio::time::pause();
4543
4544 let origin = Origin::random().produce();
4545
4546 let consumer = origin.consume();
4547
4548 let wait = tokio::spawn({
4549 let consumer = consumer.clone();
4550 async move { consumer.announced_broadcast("target").await }
4551 });
4552
4553 tokio::task::yield_now().await;
4554
4555 let _other = origin.create_broadcast("other", announce()).unwrap();
4557 settle().await;
4558 tokio::task::yield_now().await;
4559 assert!(!wait.is_finished(), "must not resolve on unrelated path");
4560
4561 let _target = origin.create_broadcast("target", announce()).unwrap();
4562 settle().await;
4563 let result = wait.await.unwrap().expect("should find target");
4564 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
4565 }
4566
4567 #[tokio::test]
4568 async fn test_announced_broadcast_skips_nested_paths() {
4569 tokio::time::pause();
4570
4571 let origin = Origin::random().produce();
4572
4573 let consumer = origin.consume();
4574
4575 let wait = tokio::spawn({
4576 let consumer = consumer.clone();
4577 async move { consumer.announced_broadcast("foo").await }
4578 });
4579
4580 tokio::task::yield_now().await;
4581
4582 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
4584 settle().await;
4585 tokio::task::yield_now().await;
4586 assert!(!wait.is_finished(), "must not resolve on a nested path");
4587
4588 let _exact = origin.create_broadcast("foo", announce()).unwrap();
4589 settle().await;
4590 let result = wait.await.unwrap().expect("should find foo exactly");
4591 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
4592 }
4593
4594 #[tokio::test]
4595 async fn test_announced_broadcast_disallowed() {
4596 let origin = Origin::random().produce();
4597 let limited = origin
4598 .consume()
4599 .scope(&["allowed".into()])
4600 .expect("should create limited");
4601
4602 assert!(limited.announced_broadcast("notallowed").await.is_none());
4604 }
4605
4606 #[tokio::test]
4607 async fn test_announced_broadcast_scope_too_narrow() {
4608 let origin = Origin::random().produce();
4611 let limited = origin
4612 .consume()
4613 .scope(&["foo/specific".into()])
4614 .expect("should create limited");
4615
4616 let result = limited
4618 .announced_broadcast("foo")
4619 .now_or_never()
4620 .expect("must not block");
4621 assert!(result.is_none());
4622 }
4623
4624 #[tokio::test]
4628 async fn test_coalesce_announce_then_unannounce() {
4629 tokio::time::pause();
4631
4632 let origin = Origin::random().produce();
4633 let mut announced = origin.consume().announced();
4634
4635 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4636 settle().await;
4637 broadcast.finish();
4638
4639 settle().await;
4640
4641 announced.assert_next_wait();
4642 }
4643
4644 #[tokio::test]
4645 async fn test_coalesce_announce_unannounce_announce() {
4646 tokio::time::pause();
4649
4650 let origin = Origin::random().produce();
4651 let mut announced = origin.consume().announced();
4652
4653 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4654 settle().await;
4655 broadcast1.finish();
4656 settle().await;
4657 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4658 settle().await;
4659
4660 announced.assert_next_some("test");
4661 announced.assert_next_wait();
4662 }
4663
4664 #[tokio::test]
4665 async fn test_coalesce_unannounce_announce_preserved() {
4666 tokio::time::pause();
4669
4670 let origin = Origin::random().produce();
4671 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4672 settle().await;
4673
4674 let mut announced = origin.consume().announced();
4675 announced.assert_next_some("test");
4676
4677 broadcast1.finish();
4679 settle().await;
4680
4681 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4682 settle().await;
4683
4684 announced.assert_next_none("test");
4686 announced.assert_next_some("test");
4687 announced.assert_next_wait();
4688 }
4689
4690 #[tokio::test]
4691 async fn test_coalesce_unannounce_announce_unannounce() {
4692 tokio::time::pause();
4695
4696 let origin = Origin::random().produce();
4697 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4698 settle().await;
4699
4700 let mut announced = origin.consume().announced();
4701 announced.assert_next_some("test");
4702
4703 broadcast1.finish();
4704 settle().await;
4705
4706 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4707 settle().await;
4708 broadcast2.finish();
4709 settle().await;
4710
4711 announced.assert_next_none("test");
4712 announced.assert_next_wait();
4713 }
4714
4715 #[tokio::test]
4716 async fn test_coalesce_churn_bounded() {
4717 tokio::time::pause();
4722
4723 let origin = Origin::random().produce();
4724 let mut announced = origin.consume().announced();
4725
4726 for _ in 0..1000 {
4727 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
4728 settle().await;
4729 broadcast.finish();
4730 }
4731 settle().await;
4732
4733 let mut collected = Vec::new();
4734 while let Some(update) = announced.try_next() {
4735 collected.push(update);
4736 }
4737 assert!(
4738 collected.len() <= 1,
4739 "expected at most one pending update, got {}",
4740 collected.len()
4741 );
4742 assert!(
4743 collected.iter().all(|a| a.path == Path::new("test")),
4744 "unexpected path in pending updates",
4745 );
4746 }
4747
4748 #[tokio::test]
4752 async fn test_consumer_clone_is_side_effect_free() {
4753 let origin = Origin::random().produce();
4754
4755 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
4756 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
4757 settle().await;
4758
4759 let consumer = origin.consume();
4760 let mut announced = consumer.announced();
4761
4762 for _ in 0..16 {
4765 let cloned = consumer.clone();
4766 assert!(cloned.get_broadcast("test1").is_some());
4767 assert!(cloned.get_broadcast("test2").is_some());
4768 }
4769
4770 let a1 = announced.try_next().expect("first announcement");
4773 let a2 = announced.try_next().expect("second announcement");
4774 announced.assert_next_wait();
4775
4776 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
4777 paths.sort();
4778 assert_eq!(paths, ["test1", "test2"]);
4779
4780 let mut fresh = consumer.announced();
4782 let b1 = fresh.try_next().expect("backlog: first");
4783 let b2 = fresh.try_next().expect("backlog: second");
4784 fresh.assert_next_wait();
4785
4786 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
4787 paths.sort();
4788 assert_eq!(paths, ["test1", "test2"]);
4789 }
4790
4791 #[tokio::test]
4793 async fn dynamic_request_unroutable_without_handler() {
4794 let origin = Origin::random().produce();
4795 let consumer = origin.consume();
4796 assert!(matches!(
4797 consumer.request_broadcast("missing").await,
4798 Err(Error::Unroutable)
4799 ));
4800 }
4801
4802 #[tokio::test(start_paused = true)]
4805 async fn dynamic_request_served_not_announced() {
4806 let origin = Origin::random().produce();
4807 let mut dynamic = origin.dynamic();
4808 let consumer = origin.consume();
4809
4810 let mut announced = origin.consume().announced();
4812 announced.assert_next_wait();
4813
4814 let served = broadcast::Info::new().produce();
4815 let request_fut = consumer.request_broadcast("fallback");
4818
4819 let mut served_dynamic = served.dynamic();
4821
4822 let request = dynamic.requested_broadcast().await.unwrap();
4823 assert_eq!(request.path(), &Path::new("fallback"));
4824 request.accept(&served);
4825
4826 let broadcast = request_fut.await.unwrap();
4827 assert!(broadcast.is_clone(&served.consume()));
4828
4829 let track_fut = broadcast.track("video").unwrap().subscribe(None);
4831 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
4832 let mut track = track_fut.await.unwrap();
4833 producer.append_group().unwrap();
4834 track.assert_group();
4835
4836 announced.assert_next_wait();
4838 }
4839
4840 #[tokio::test(start_paused = true)]
4842 async fn dynamic_request_coalesces() {
4843 let origin = Origin::random().produce();
4844 let mut dynamic = origin.dynamic();
4845 let consumer = origin.consume();
4846
4847 let f1 = consumer.request_broadcast("dup");
4849 let f2 = consumer.request_broadcast("dup");
4850
4851 let request = dynamic.requested_broadcast().await.unwrap();
4853 assert_eq!(request.path(), &Path::new("dup"));
4854 assert!(
4855 dynamic.requested_broadcast().now_or_never().is_none(),
4856 "a coalesced request must not be served twice"
4857 );
4858
4859 let served = broadcast::Info::new().produce();
4861 request.accept(&served);
4862 assert!(f1.await.unwrap().is_clone(&served.consume()));
4863 assert!(f2.await.unwrap().is_clone(&served.consume()));
4864 }
4865
4866 #[tokio::test(start_paused = true)]
4869 async fn dynamic_request_dedups_served() {
4870 let origin = Origin::random().produce();
4871 let mut dynamic = origin.dynamic();
4872 let consumer = origin.consume();
4873
4874 let request_fut = consumer.request_broadcast("fallback");
4875 let request = dynamic.requested_broadcast().await.unwrap();
4876 let served = broadcast::Info::new().produce();
4877 request.accept(&served);
4878 let first = request_fut.await.unwrap();
4879 assert!(first.is_clone(&served.consume()));
4880
4881 let second = consumer.request_broadcast("fallback").await.unwrap();
4883 assert!(second.is_clone(&served.consume()));
4884
4885 assert!(
4887 dynamic.requested_broadcast().now_or_never().is_none(),
4888 "a still-live served broadcast must not be re-requested from the handler"
4889 );
4890 }
4891
4892 #[tokio::test(start_paused = true)]
4894 async fn dynamic_request_reserves_after_close() {
4895 let origin = Origin::random().produce();
4896 let mut dynamic = origin.dynamic();
4897 let consumer = origin.consume();
4898
4899 let request_fut = consumer.request_broadcast("fallback");
4900 let request = dynamic.requested_broadcast().await.unwrap();
4901 let served = broadcast::Info::new().produce();
4902 request.accept(&served);
4903 request_fut.await.unwrap();
4904
4905 drop(served);
4907
4908 let request_fut = consumer.request_broadcast("fallback");
4910 let request = dynamic.requested_broadcast().await.unwrap();
4911 assert_eq!(request.path(), &Path::new("fallback"));
4912 let served = broadcast::Info::new().produce();
4913 request.accept(&served);
4914 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
4915 }
4916
4917 #[tokio::test(start_paused = true)]
4920 async fn dynamic_request_served_cache_bounded() {
4921 let origin = Origin::random().produce();
4922 let mut dynamic = origin.dynamic();
4923 let consumer = origin.consume();
4924
4925 for i in 0..100 {
4926 let path = format!("one-shot/{i}");
4927 let request_fut = consumer.request_broadcast(&path);
4928 let request = dynamic.requested_broadcast().await.unwrap();
4929 let served = broadcast::Info::new().produce();
4930 request.accept(&served);
4931 request_fut.await.unwrap();
4932 drop(served);
4934 }
4935
4936 assert!(
4939 origin.dynamic.read().served.len() <= 4,
4940 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
4941 origin.dynamic.read().served.len()
4942 );
4943 }
4944
4945 #[tokio::test(start_paused = true)]
4948 async fn dynamic_request_coalesces_after_handoff() {
4949 let origin = Origin::random().produce();
4950 let mut dynamic = origin.dynamic();
4951 let consumer = origin.consume();
4952
4953 let f1 = consumer.request_broadcast("fallback");
4954 let request = dynamic.requested_broadcast().await.unwrap();
4956
4957 let f2 = consumer.request_broadcast("fallback");
4959 assert!(
4960 dynamic.requested_broadcast().now_or_never().is_none(),
4961 "a repeat request during hand-off must coalesce, not re-queue"
4962 );
4963
4964 let served = broadcast::Info::new().produce();
4966 request.accept(&served);
4967 assert!(f1.await.unwrap().is_clone(&served.consume()));
4968 assert!(f2.await.unwrap().is_clone(&served.consume()));
4969 }
4970
4971 #[tokio::test(start_paused = true)]
4973 async fn dynamic_request_dropped_after_handoff() {
4974 let origin = Origin::random().produce();
4975 let mut dynamic = origin.dynamic();
4976 let consumer = origin.consume();
4977
4978 let f1 = consumer.request_broadcast("fallback");
4979 let request = dynamic.requested_broadcast().await.unwrap();
4980 let f2 = consumer.request_broadcast("fallback");
4981
4982 drop(request);
4984 assert!(matches!(f1.await, Err(Error::Unroutable)));
4985 assert!(matches!(f2.await, Err(Error::Unroutable)));
4986 }
4987
4988 #[tokio::test(start_paused = true)]
4990 async fn dynamic_request_rejected() {
4991 let origin = Origin::random().produce();
4992 let mut dynamic = origin.dynamic();
4993 let consumer = origin.consume();
4994
4995 let request_fut = consumer.request_broadcast("fallback");
4996
4997 let request = dynamic.requested_broadcast().await.unwrap();
4998 request.reject(Error::Cancel);
4999
5000 assert!(matches!(request_fut.await, Err(Error::Cancel)));
5001 }
5002
5003 #[tokio::test(start_paused = true)]
5007 async fn dynamic_request_rerequest_after_reject() {
5008 let origin = Origin::random().produce();
5009 let mut dynamic = origin.dynamic();
5010 let consumer = origin.consume();
5011
5012 let f1 = consumer.request_broadcast("fallback");
5013 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
5014 assert!(matches!(f1.await, Err(Error::Unroutable)));
5015
5016 let served = broadcast::Info::new().produce();
5017 let f2 = consumer.request_broadcast("fallback");
5019 let request = dynamic.requested_broadcast().await.unwrap();
5020 assert_eq!(request.path(), &Path::new("fallback"));
5021 request.accept(&served);
5022 assert!(f2.await.unwrap().is_clone(&served.consume()));
5023 }
5024
5025 #[tokio::test(start_paused = true)]
5028 async fn dynamic_request_handler_dropped() {
5029 let origin = Origin::random().produce();
5030 let dynamic = origin.dynamic();
5031 let consumer = origin.consume();
5032
5033 let request_fut = consumer.request_broadcast("fallback");
5034 drop(dynamic);
5035 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
5036
5037 assert!(matches!(
5039 consumer.request_broadcast("again").await,
5040 Err(Error::Unroutable)
5041 ));
5042 }
5043
5044 #[tokio::test(start_paused = true)]
5048 async fn dynamic_request_accept_after_handler_dropped() {
5049 let origin = Origin::random().produce();
5050 let mut dynamic = origin.dynamic();
5051 let consumer = origin.consume();
5052
5053 let request_fut = consumer.request_broadcast("fallback");
5054
5055 let request = dynamic.requested_broadcast().await.unwrap();
5057 drop(dynamic);
5058
5059 let served = broadcast::Info::new().produce();
5060 request.accept(&served);
5062 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
5063 }
5064
5065 #[tokio::test(start_paused = true)]
5067 async fn dynamic_request_prefers_announced() {
5068 let origin = Origin::random().produce();
5069 let mut dynamic = origin.dynamic();
5070 let consumer = origin.consume();
5071
5072 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
5073 settle().await;
5074
5075 let got = consumer.request_broadcast("live").await.unwrap();
5076 assert!(
5077 got.is_clone(&consumer.get_broadcast("live").unwrap()),
5078 "should return the published broadcast"
5079 );
5080 assert!(
5081 dynamic.requested_broadcast().now_or_never().is_none(),
5082 "a published path must not queue a fallback request"
5083 );
5084 }
5085
5086 #[tokio::test(start_paused = true)]
5088 async fn dynamic_clone_keeps_alive() {
5089 let origin = Origin::random().produce();
5090 let dynamic = origin.dynamic();
5091 let consumer = origin.consume();
5092
5093 drop(dynamic.clone());
5094
5095 let request_fut = consumer.request_broadcast("fallback");
5098 assert!(
5099 request_fut.now_or_never().is_none(),
5100 "request should stay pending until served"
5101 );
5102 }
5103}