1use crate::{broadcast, cache, stats, track};
2use kio::Pollable;
3use std::{
4 cmp::Reverse,
5 collections::{BTreeMap, HashMap, HashSet},
6 fmt,
7 sync::Arc,
8 sync::atomic::{AtomicU64, Ordering},
9 task::{Poll, ready},
10 time::Duration,
11};
12
13use rand::RngExt;
14use web_async::Lock;
15
16use super::{Requests, WeakCache};
17use crate::{
18 AsPath, Error, Path, PathOwned, PathPrefixes,
19 coding::{BoundsExceeded, Decode, DecodeError, Encode, EncodeError},
20};
21
22#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
28pub struct Origin {
29 id: u64,
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
35#[non_exhaustive]
36pub struct InvalidOrigin;
37
38impl fmt::Display for InvalidOrigin {
39 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
40 write!(f, "local origin id must be non-zero and below 2^62")
41 }
42}
43
44impl std::error::Error for InvalidOrigin {}
45
46impl Origin {
47 pub(crate) const UNKNOWN: Self = Self { id: 0 };
50
51 pub fn new(id: u64) -> Result<Self, InvalidOrigin> {
57 if id == 0 || id >= 1u64 << 62 {
58 return Err(InvalidOrigin);
59 }
60 Ok(Self { id })
61 }
62
63 pub fn random() -> Self {
72 let mut rng = rand::rng();
73 let id = rng.random_range(1..(1u64 << 53));
74 Self { id }
75 }
76
77 pub fn id(self) -> u64 {
79 self.id
80 }
81
82 pub fn produce(self) -> Producer {
85 Info::new(self).produce()
86 }
87}
88
89#[derive(Clone, Debug)]
99#[non_exhaustive]
100pub struct Info {
101 pub id: Origin,
104
105 pub pool: cache::Pool,
111
112 pub cache_duration: Duration,
119
120 pub latency_default: Duration,
130
131 pub linger: Duration,
140}
141
142impl Default for Info {
143 fn default() -> Self {
146 Self {
147 id: Origin::UNKNOWN,
148 pool: cache::Pool::default(),
149 cache_duration: Duration::MAX,
150 latency_default: track::DEFAULT_LATENCY_MAX,
151 linger: Duration::ZERO,
152 }
153 }
154}
155
156impl Info {
157 pub fn new(id: Origin) -> Self {
159 Self { id, ..Self::default() }
160 }
161
162 pub fn with_pool(mut self, pool: cache::Pool) -> Self {
164 self.pool = pool;
165 self
166 }
167
168 pub fn with_cache_duration(mut self, cache_duration: Duration) -> Self {
171 self.cache_duration = cache_duration;
172 self
173 }
174
175 pub fn with_latency_default(mut self, latency_default: Duration) -> Self {
178 self.latency_default = latency_default;
179 self
180 }
181
182 pub fn with_linger(mut self, linger: Duration) -> Self {
185 self.linger = linger;
186 self
187 }
188
189 pub fn produce(self) -> Producer {
191 Producer::new(self)
192 }
193}
194
195impl TryFrom<u64> for Origin {
196 type Error = InvalidOrigin;
197
198 fn try_from(id: u64) -> Result<Self, Self::Error> {
199 Self::new(id)
200 }
201}
202
203impl fmt::Display for Origin {
204 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
205 self.id.fmt(f)
206 }
207}
208
209impl<V: Copy> Encode<V> for Origin
210where
211 u64: Encode<V>,
212{
213 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
214 self.id.encode(w, version)
215 }
216}
217
218impl<V: Copy> Decode<V> for Origin
219where
220 u64: Decode<V>,
221{
222 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
223 let id = u64::decode(r, version)?;
224 if id >= 1u64 << 62 {
225 return Err(DecodeError::InvalidValue);
226 }
227 Ok(Self { id })
228 }
229}
230
231pub(crate) const MAX_HOPS: usize = 32;
237
238#[derive(Debug, Clone, Default, PartialEq, Eq)]
243pub struct OriginList(Vec<Origin>);
244
245#[derive(Debug, Clone, Copy, PartialEq, Eq)]
247#[non_exhaustive]
248pub struct TooManyOrigins;
249
250impl fmt::Display for TooManyOrigins {
251 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
252 write!(f, "too many origins (max {MAX_HOPS})")
253 }
254}
255
256impl std::error::Error for TooManyOrigins {}
257
258impl From<TooManyOrigins> for DecodeError {
259 fn from(_: TooManyOrigins) -> Self {
260 DecodeError::BoundsExceeded
261 }
262}
263
264impl OriginList {
265 pub fn new() -> Self {
267 Self(Vec::new())
268 }
269
270 pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
272 if self.0.len() >= MAX_HOPS {
273 return Err(TooManyOrigins);
274 }
275 self.0.push(origin);
276 Ok(())
277 }
278
279 pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
282 for entry in &mut self.0 {
283 if *entry == target {
284 *entry = replacement;
285 return true;
286 }
287 }
288 false
289 }
290
291 pub fn contains(&self, origin: &Origin) -> bool {
293 self.0.contains(origin)
294 }
295
296 pub fn len(&self) -> usize {
298 self.0.len()
299 }
300
301 pub fn is_empty(&self) -> bool {
303 self.0.is_empty()
304 }
305
306 pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
308 self.0.iter()
309 }
310
311 pub fn as_slice(&self) -> &[Origin] {
313 &self.0
314 }
315}
316
317impl TryFrom<Vec<Origin>> for OriginList {
318 type Error = TooManyOrigins;
319
320 fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
321 if v.len() > MAX_HOPS {
322 return Err(TooManyOrigins);
323 }
324 Ok(Self(v))
325 }
326}
327
328impl<'a> IntoIterator for &'a OriginList {
329 type Item = &'a Origin;
330 type IntoIter = std::slice::Iter<'a, Origin>;
331
332 fn into_iter(self) -> Self::IntoIter {
333 self.iter()
334 }
335}
336
337impl<V: Copy> Encode<V> for OriginList
338where
339 u64: Encode<V>,
340 Origin: Encode<V>,
341{
342 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
343 (self.0.len() as u64).encode(w, version)?;
344 for origin in &self.0 {
345 origin.encode(w, version)?;
346 }
347 Ok(())
348 }
349}
350
351impl<V: Copy> Decode<V> for OriginList
352where
353 u64: Decode<V>,
354 Origin: Decode<V>,
355{
356 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
357 let count = u64::decode(r, version)? as usize;
358 if count > MAX_HOPS {
359 return Err(DecodeError::BoundsExceeded);
360 }
361 let mut list = Vec::with_capacity(count);
362 for _ in 0..count {
363 list.push(Origin::decode(r, version)?);
364 }
365 Ok(Self(list))
366 }
367}
368
369static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
370
371#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
372struct ConsumerId(u64);
373
374impl ConsumerId {
375 fn new() -> Self {
376 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
377 }
378}
379
380struct OriginBroadcast {
386 path: PathOwned,
387 broadcast: broadcast::Producer,
389 state: kio::Producer<FrontState>,
392 announced: bool,
393}
394
395fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
404 (hops.len(), fnv_key(name, hops.iter().copied()))
405}
406
407fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
420 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
422
423 let mut hash = SEED;
424 for &byte in name.as_str().as_bytes() {
425 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
426 }
427 for origin in origins {
428 for &byte in &origin.id().to_le_bytes() {
429 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
430 }
431 }
432
433 hash
434}
435
436fn route_order(name: &Path, route: &FrontRoute) -> (bool, u64, usize, u64, Reverse<u64>) {
455 let (len, hash) = route_key(name, &route.route.hops);
456 (!route.route.announce, route.route.cost, len, hash, Reverse(route.id))
457}
458
459enum PendingUpdate {
468 Announce(broadcast::Consumer),
469 Unannounce,
470 UnannounceAnnounce(broadcast::Consumer),
471}
472
473#[derive(Default)]
478struct OriginConsumerState {
479 pending: BTreeMap<PathOwned, PendingUpdate>,
480}
481
482impl OriginConsumerState {
483 fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
484 let new = match self.pending.remove(&path) {
485 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
487 Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
489 PendingUpdate::UnannounceAnnounce(broadcast)
490 }
491 };
492 self.pending.insert(path, new);
493 }
494
495 fn apply_unannounce(&mut self, path: PathOwned) {
496 match self.pending.remove(&path) {
497 Some(PendingUpdate::Announce(_)) => {}
499 None | Some(PendingUpdate::Unannounce) => {
500 self.pending.insert(path, PendingUpdate::Unannounce);
501 }
502 Some(PendingUpdate::UnannounceAnnounce(_)) => {
505 self.pending.insert(path, PendingUpdate::Unannounce);
506 }
507 }
508 }
509
510 fn take(&mut self) -> Option<OriginAnnounce> {
512 let path = self.pending.keys().next()?.clone();
513 let broadcast = match self.pending.remove(&path).unwrap() {
514 PendingUpdate::Announce(broadcast) => Some(broadcast),
515 PendingUpdate::Unannounce => None,
516 PendingUpdate::UnannounceAnnounce(broadcast) => {
517 self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
520 None
521 }
522 };
523 Some(OriginAnnounce { path, broadcast })
524 }
525}
526
527#[derive(Clone)]
528struct AnnounceConsumerNotify {
529 root: PathOwned,
530 state: kio::Producer<OriginConsumerState>,
531 exclude: Option<Origin>,
536}
537
538impl AnnounceConsumerNotify {
539 fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer, front: &kio::Producer<FrontState>) {
540 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
541
542 let broadcast = match self.exclude.and_then(|peer| ExclusionGuard::new(front, peer)) {
547 Some(guard) => broadcast.with_exclusion(guard),
548 None => broadcast,
549 };
550
551 self.state
552 .write()
553 .ok()
554 .expect("consumer closed")
555 .apply_announce(path, broadcast);
556 }
557
558 fn unannounce(&self, path: impl AsPath) {
559 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
560 self.state.write().ok().expect("consumer closed").apply_unannounce(path);
561 }
562}
563
564struct NotifyNode {
565 parent: Option<Lock<NotifyNode>>,
566
567 consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
570}
571
572impl NotifyNode {
573 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
574 Self {
575 parent,
576 consumers: HashMap::new(),
577 }
578 }
579
580 fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer, state: &kio::Producer<FrontState>) {
584 for consumer in self.consumers.values() {
585 consumer.announce(path.as_path(), broadcast.clone(), state);
586 }
587
588 if let Some(parent) = &self.parent {
589 parent.lock().announce(path, broadcast, state);
590 }
591 }
592
593 fn unannounce(&mut self, path: impl AsPath) {
594 for consumer in self.consumers.values() {
595 consumer.unannounce(path.as_path());
596 }
597
598 if let Some(parent) = &self.parent {
599 parent.lock().unannounce(path);
600 }
601 }
602}
603
604pub(crate) struct ExclusionGuard {
611 state: kio::Producer<FrontState>,
612 peer: Origin,
613}
614
615impl ExclusionGuard {
616 fn new(state: &kio::Producer<FrontState>, peer: Origin) -> Option<Arc<Self>> {
619 let mut s = state.write().ok()?;
620 if s.closed {
621 return None;
622 }
623 *s.excluded.entry(peer).or_default() += 1;
624 drop(s);
625 Some(Arc::new(Self {
626 state: state.clone(),
627 peer,
628 }))
629 }
630}
631
632impl Drop for ExclusionGuard {
633 fn drop(&mut self) {
634 let Ok(mut state) = self.state.write() else { return };
635 if let std::collections::hash_map::Entry::Occupied(mut entry) = state.excluded.entry(self.peer) {
636 match entry.get() {
637 1 => drop(entry.remove()),
638 n => *entry.get_mut() = n - 1,
639 }
640 }
641 }
642}
643
644enum Resolved {
651 Found(broadcast::Consumer),
654 Excluded,
656 Missing,
658}
659
660struct OriginNode {
661 broadcast: Option<OriginBroadcast>,
664
665 nested: HashMap<String, Lock<OriginNode>>,
667
668 notify: Lock<NotifyNode>,
670}
671
672impl OriginNode {
673 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
674 Self {
675 broadcast: None,
676 nested: HashMap::new(),
677 notify: Lock::new(NotifyNode::new(parent)),
678 }
679 }
680
681 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
682 let (dir, rest) = path.next_part().expect("leaf called with empty path");
683
684 let next = self.entry(dir);
685 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
686 }
687
688 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
689 match self.nested.get(dir) {
690 Some(next) => next.clone(),
691 None => {
692 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
693 self.nested.insert(dir.to_string(), next.clone());
694 next
695 }
696 }
697 }
698
699 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
703 let Some(existing) = &mut self.broadcast else { return };
704 if !existing.state.same_channel(expect) || existing.announced == announce {
705 return;
706 }
707 existing.announced = announce;
708 let path = existing.path.clone();
709 let consumer = existing.broadcast.consume();
710 let state = existing.state.clone();
711 let mut notify = self.notify.lock();
712 if announce {
713 notify.announce(&path, &consumer, &state);
714 } else {
715 notify.unannounce(&path);
716 }
717 }
718
719 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
720 self.consume_initial(&mut notify);
721 self.notify.lock().consumers.insert(id, notify);
722 }
723
724 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
725 if let Some(broadcast) = &self.broadcast
728 && broadcast.announced
729 {
730 notify.announce(&broadcast.path, broadcast.broadcast.consume(), &broadcast.state);
731 }
732
733 for nested in self.nested.values() {
735 nested.lock().consume_initial(notify);
736 }
737 }
738
739 fn resolve_broadcast(&self, rest: impl AsPath, exclude: Option<Origin>) -> Resolved {
740 let rest = rest.as_path();
741
742 if let Some((dir, rest)) = rest.next_part() {
743 let Some(node) = self.nested.get(dir) else {
744 return Resolved::Missing;
745 };
746 let node = node.lock();
747 return node.resolve_broadcast(&rest, exclude);
748 }
749
750 let Some(broadcast) = self.broadcast.as_ref() else {
751 return Resolved::Missing;
752 };
753 let Some(origin) = exclude else {
754 return Resolved::Found(broadcast.broadcast.consume());
755 };
756
757 let state = broadcast.state.read();
773 if !state.routes.iter().any(|r| r.route.hops.contains(&origin)) {
774 drop(state);
775 let shared = broadcast.broadcast.consume();
776 return match ExclusionGuard::new(&broadcast.state, origin) {
777 Some(guard) => Resolved::Found(shared.with_exclusion(guard)),
778 None => Resolved::Found(shared),
781 };
782 }
783 match state
784 .dispatch(Some(origin))
785 .and_then(|clean| state.routes.iter().find(|r| r.id == clean))
786 {
787 Some(route) => Resolved::Found(route.source.clone()),
788 None => Resolved::Excluded,
789 }
790 }
791
792 fn unconsume(&mut self, id: ConsumerId) {
793 self.notify.lock().consumers.remove(&id).expect("consumer not found");
794 if self.is_empty() {
795 }
798 }
799
800 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
804 let relative = relative.as_path();
805
806 if let Some((dir, relative)) = relative.next_part() {
807 let Some(nested) = self.nested.get(dir) else { return };
808 let nested = nested.clone();
809 let mut locked = nested.lock();
810 locked.remove(expect, &relative);
811
812 if locked.is_empty() {
813 drop(locked);
814 self.nested.remove(dir);
815 }
816 } else if let Some(existing) = &self.broadcast
817 && existing.state.same_channel(expect)
818 {
819 let existing = self.broadcast.take().expect("checked above");
820 if existing.announced {
821 self.notify.lock().unannounce(&existing.path);
822 }
823 }
824 }
825
826 fn is_empty(&self) -> bool {
827 self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
828 }
829}
830
831#[derive(Clone)]
832struct OriginNodes {
833 nodes: Vec<(PathOwned, Lock<OriginNode>)>,
834}
835
836impl OriginNodes {
837 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
840 let mut roots = Vec::new();
841
842 for (root, state) in &self.nodes {
843 for prefix in prefixes {
844 if root.has_prefix(prefix) {
845 roots.push((root.to_owned(), state.clone()));
847 continue;
848 }
849
850 if let Some(suffix) = prefix.strip_prefix(root) {
851 let nested = state.lock().leaf(&suffix);
853 roots.push((prefix.to_owned(), nested));
854 }
855 }
856 }
857
858 if roots.is_empty() {
859 None
860 } else {
861 Some(Self { nodes: roots })
862 }
863 }
864
865 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
866 let new_root = new_root.as_path();
867 let mut roots = Vec::new();
868
869 if new_root.is_empty() {
870 return Some(self.clone());
871 }
872
873 for (root, state) in &self.nodes {
874 if let Some(suffix) = root.strip_prefix(&new_root) {
875 roots.push((suffix.to_owned(), state.clone()));
877 } else if let Some(suffix) = new_root.strip_prefix(root) {
878 let nested = state.lock().leaf(&suffix);
881 roots.push(("".into(), nested));
882 }
883 }
884
885 if roots.is_empty() {
886 None
887 } else {
888 Some(Self { nodes: roots })
889 }
890 }
891
892 pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
894 let path = path.as_path();
895
896 for (root, state) in &self.nodes {
897 if let Some(suffix) = path.strip_prefix(root) {
898 return Some((state.clone(), suffix.to_owned()));
899 }
900 }
901
902 None
903 }
904}
905
906impl Default for OriginNodes {
907 fn default() -> Self {
908 Self {
909 nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
910 }
911 }
912}
913
914#[derive(Clone)]
916pub struct OriginAnnounce {
917 pub path: PathOwned,
919 pub broadcast: Option<broadcast::Consumer>,
925}
926
927#[derive(Clone)]
929pub struct Producer {
930 info: Origin,
934
935 nodes: OriginNodes,
938
939 root: PathOwned,
941
942 dynamic: kio::Shared<OriginDynamicState>,
946
947 pool: cache::Pool,
950
951 cache_duration: Duration,
954
955 latency_default: Duration,
958
959 linger: Duration,
962
963 stats: stats::Session,
967}
968
969impl std::ops::Deref for Producer {
970 type Target = Origin;
971
972 fn deref(&self) -> &Self::Target {
973 &self.info
974 }
975}
976
977impl Producer {
978 pub fn new(info: Info) -> Self {
982 Self {
983 info: info.id,
984 nodes: OriginNodes::default(),
985 root: PathOwned::default(),
986 dynamic: kio::Shared::default(),
987 pool: info.pool,
988 cache_duration: info.cache_duration,
989 latency_default: info.latency_default,
990 linger: info.linger,
991 stats: stats::Session::default(),
992 }
993 }
994
995 pub fn with_stats(mut self, session: stats::Session) -> Self {
999 self.stats = session;
1000 self
1001 }
1002
1003 pub fn with_linger(mut self, linger: Duration) -> Self {
1012 self.linger = linger;
1013 self
1014 }
1015
1016 pub fn info(&self) -> Info {
1019 Info {
1020 id: self.info,
1021 pool: self.pool.clone(),
1022 cache_duration: self.cache_duration,
1023 latency_default: self.latency_default,
1024 linger: self.linger,
1025 }
1026 }
1027
1028 pub(crate) fn latency_default(&self) -> Duration {
1031 self.latency_default
1032 }
1033
1034 pub(crate) fn empty(info: Origin) -> Self {
1039 Self {
1040 info,
1041 nodes: OriginNodes { nodes: Vec::new() },
1042 root: PathOwned::default(),
1043 dynamic: kio::Shared::default(),
1044 pool: cache::Pool::default(),
1045 cache_duration: Duration::MAX,
1046 latency_default: track::DEFAULT_LATENCY_MAX,
1047 linger: Duration::ZERO,
1048 stats: stats::Session::default(),
1049 }
1050 }
1051
1052 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
1103 let path = path.as_path();
1104
1105 debug_assert!(
1106 !route.hops.contains(&self.info),
1107 "create_broadcast called with a looping hop chain",
1108 );
1109
1110 let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
1111 let full = self.root.join(&path).to_owned();
1112
1113 if full.parts().count() > Path::MAX_PARTS {
1117 return Err(BoundsExceeded.into());
1118 }
1119
1120 let ingress = self.stats.ingress(&full);
1124
1125 let mut source = broadcast::Info { origin: self.info() }
1126 .produce()
1127 .with_stats(ingress.clone());
1128 source.set_route(route).expect("fresh producer");
1129
1130 web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
1131
1132 Ok(source)
1133 }
1134
1135 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
1141 let prefixes = PathPrefixes::new(prefixes);
1142 Some(Producer {
1143 info: self.info,
1144 nodes: self.nodes.select(&prefixes)?,
1145 root: self.root.clone(),
1146 dynamic: self.dynamic.clone(),
1147 pool: self.pool.clone(),
1148 cache_duration: self.cache_duration,
1149 latency_default: self.latency_default,
1150 linger: self.linger,
1151 stats: self.stats.clone(),
1152 })
1153 }
1154
1155 pub fn dynamic(&self) -> Dynamic {
1164 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
1165 }
1166
1167 pub fn consume(&self) -> Consumer {
1172 Consumer::new(
1175 self.info,
1176 self.root.clone(),
1177 self.nodes.clone(),
1178 self.dynamic.clone(),
1179 stats::Session::default(),
1180 )
1181 }
1182
1183 pub fn announces(&self) -> AnnounceProducer {
1189 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
1190 }
1191
1192 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1197 let prefix = prefix.as_path();
1198
1199 Some(Self {
1200 info: self.info,
1201 root: self.root.join(&prefix).to_owned(),
1202 nodes: self.nodes.root(&prefix)?,
1203 dynamic: self.dynamic.clone(),
1204 pool: self.pool.clone(),
1205 cache_duration: self.cache_duration,
1206 latency_default: self.latency_default,
1207 linger: self.linger,
1208 stats: self.stats.clone(),
1209 })
1210 }
1211
1212 pub fn root(&self) -> &Path<'_> {
1214 &self.root
1215 }
1216
1217 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1220 self.nodes.nodes.iter().map(|(root, _)| root)
1221 }
1222
1223 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1225 self.root.join(path)
1226 }
1227}
1228
1229const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1242
1243struct FrontRoute {
1245 id: u64,
1246 route: broadcast::Route,
1249 source: broadcast::Consumer,
1251}
1252
1253struct FrontState {
1255 path: PathOwned,
1257 self_origin: Origin,
1259 publisher: Option<Origin>,
1267 next_route: u64,
1270 routes: Vec<FrontRoute>,
1271 excluded: HashMap<Origin, usize>,
1284 active: Option<u64>,
1286 linger: Duration,
1289 closed: bool,
1294}
1295
1296impl FrontState {
1297 fn pick(&self, keep: impl Fn(&FrontRoute) -> bool, untainted: bool) -> Option<u64> {
1304 let candidates: Vec<&FrontRoute> = self.routes.iter().filter(|r| keep(r)).collect();
1305 let candidates = match untainted {
1306 true => self.prefer_untainted(&candidates),
1307 false => candidates,
1308 };
1309 candidates
1310 .into_iter()
1311 .min_by_key(|r| route_order(&self.path.as_path(), r))
1312 .map(|r| r.id)
1313 }
1314
1315 fn best_route(&self) -> Option<u64> {
1320 self.pick(|_| true, true)
1321 }
1322
1323 fn dispatch(&self, exclude: Option<Origin>) -> Option<u64> {
1330 self.pick(|r| exclude.is_none_or(|origin| !r.route.hops.contains(&origin)), false)
1331 }
1332
1333 fn taints_a_reader(&self, route: &broadcast::Route) -> bool {
1340 route.hops.iter().any(|hop| self.excluded.contains_key(hop))
1341 }
1342
1343 fn prefer_untainted<'a>(&self, candidates: &[&'a FrontRoute]) -> Vec<&'a FrontRoute> {
1353 if self.excluded.is_empty() {
1354 return candidates.to_vec();
1355 }
1356 let clean: Vec<&FrontRoute> = candidates
1357 .iter()
1358 .copied()
1359 .filter(|r| !self.taints_a_reader(&r.route))
1360 .collect();
1361 match clean.is_empty() {
1362 true => candidates.to_vec(),
1363 false => clean,
1364 }
1365 }
1366
1367 fn serve_route(&self, skip: impl Fn(u64) -> bool) -> Option<u64> {
1376 if let Some(active) = self.active
1377 && !skip(active)
1378 && let Some(route) = self.routes.iter().find(|r| r.id == active)
1379 && !self.taints_a_reader(&route.route)
1380 {
1381 return Some(active);
1382 }
1383 self.pick(|r| !skip(r.id), true)
1384 }
1385
1386 fn reselect(&mut self, carrying: bool) {
1404 let best = self.best_route();
1405 if carrying
1406 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1407 && best_id != cur_id
1408 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1409 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1410 && incumbent.route.announce
1411 && candidate.route.cost < incumbent.route.cost
1412 && candidate.route.advertised == 0
1413 && candidate.route.hops.len() >= 2
1414 && !self.handover_allowed(&candidate.route)
1415 {
1416 return;
1418 }
1419 self.active = best;
1420 }
1421
1422 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1434 let name = self.path.as_path();
1435 match route.hops.iter().last() {
1436 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1437 None => true,
1438 }
1439 }
1440
1441 fn routes_snapshot(&self) -> Vec<broadcast::Route> {
1447 let mut routes: Vec<&FrontRoute> = self.routes.iter().collect();
1448 routes.sort_by_key(|r| route_order(&self.path.as_path(), r));
1449 routes.sort_by_key(|r| Some(r.id) != self.active);
1450 routes.into_iter().map(|r| r.route.clone()).collect()
1451 }
1452}
1453
1454fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1464 let mut leaf_guard = leaf.lock();
1469 let routes = state.read().routes_snapshot();
1470 if let Some(advert) = routes.first() {
1471 let announce = advert.announce;
1472 broadcast.clone().set_routes(routes);
1473 leaf_guard.set_announced(state, announce);
1474 }
1475}
1476
1477fn detach_source(
1490 state: &kio::Producer<FrontState>,
1491 broadcast: &broadcast::Producer,
1492 leaf: &Lock<OriginNode>,
1493 id: u64,
1494 graceful: bool,
1495) {
1496 let close = {
1497 let carrying = broadcast.demand().is_used();
1501 let Ok(mut s) = state.write() else { return };
1502 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1503 return;
1504 };
1505 s.routes.remove(pos);
1506 s.reselect(carrying);
1507 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1508 s.closed = true;
1511 true
1512 } else {
1513 false
1514 }
1515 };
1516 if close {
1517 broadcast.abort_spliced(Error::Dropped);
1518 }
1519 sync_front(state, broadcast, leaf);
1520}
1521
1522fn sync_announce(guard: &mut Option<stats::Announce>, announced: bool, ingress: &stats::Scope) {
1525 match (announced, guard.is_some()) {
1526 (true, false) => *guard = Some(ingress.announce()),
1527 (false, true) => *guard = None,
1528 _ => {}
1529 }
1530}
1531
1532async fn run_source(
1542 origin: Info,
1543 node: Lock<OriginNode>,
1544 full: PathOwned,
1545 rest: PathOwned,
1546 mut source: broadcast::Consumer,
1547 ingress: stats::Scope,
1548) {
1549 let ctx = AttachContext {
1550 origin: &origin,
1551 node: &node,
1552 full: &full,
1553 rest: &rest,
1554 };
1555
1556 let Ok(mut route) = source.route_changed().await else {
1560 return;
1562 };
1563
1564 let mut announce = route.announce.then(|| ingress.announce());
1569
1570 'attach: loop {
1571 let leaf = if rest.is_empty() {
1576 node.clone()
1577 } else {
1578 node.lock().leaf(&rest)
1579 };
1580
1581 let (state, broadcast, id) = match attach_source(&ctx, &leaf, &source, route.clone()) {
1582 Attach::Ready(state, broadcast, id) => (state, broadcast, id),
1583 Attach::Parked(incumbent) => {
1584 tracing::debug!(
1585 broadcast = %full,
1586 "path already live with a different publisher; parking this source until it ends",
1587 );
1588 let update = kio::wait(|waiter| {
1592 if let Poll::Ready(update) = source.poll_route_changed(waiter) {
1593 return Poll::Ready(Some(update));
1594 }
1595 match incumbent.poll(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending }) {
1598 Poll::Ready(_) => Poll::Ready(None),
1599 Poll::Pending => Poll::Pending,
1600 }
1601 })
1602 .await;
1603 match update {
1604 Some(Ok(update)) => {
1606 sync_announce(&mut announce, update.announce, &ingress);
1607 route = update;
1608 }
1609 Some(Err(_)) => return,
1611 None => {}
1613 }
1614 continue 'attach;
1615 }
1616 };
1617 let publisher = route.hops.iter().next().copied();
1618
1619 loop {
1620 match source.route_changed().await {
1621 Ok(update) => {
1622 let announced = update.announce;
1623 if update.hops.iter().next().copied() != publisher {
1634 detach_source(&state, &broadcast, &leaf, id, true);
1635 sync_announce(&mut announce, announced, &ingress);
1636 route = update;
1637 continue 'attach;
1638 }
1639 {
1640 let carrying = broadcast.demand().is_used();
1641 let Ok(mut s) = state.write() else { return };
1642 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1643 return;
1644 };
1645 if entry.route == update {
1646 continue;
1647 }
1648 entry.route = update;
1649 s.reselect(carrying);
1650 }
1651 sync_announce(&mut announce, announced, &ingress);
1653 sync_front(&state, &broadcast, &leaf);
1654 }
1655 Err(_) => {
1656 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1659 return;
1660 }
1661 }
1662 }
1663 }
1664}
1665
1666enum Attach {
1668 Ready(kio::Producer<FrontState>, broadcast::Producer, u64),
1671 Parked(kio::Producer<FrontState>),
1677}
1678
1679struct AttachContext<'a> {
1681 origin: &'a Info,
1682 node: &'a Lock<OriginNode>,
1683 full: &'a PathOwned,
1685 rest: &'a PathOwned,
1687}
1688
1689fn same_publisher(a: Option<Origin>, b: Option<Origin>) -> bool {
1698 if a == Some(Origin::UNKNOWN) || b == Some(Origin::UNKNOWN) {
1699 return false;
1700 }
1701 a == b
1702}
1703
1704fn attach_source(
1725 ctx: &AttachContext,
1726 leaf: &Lock<OriginNode>,
1727 source: &broadcast::Consumer,
1728 route: broadcast::Route,
1729) -> Attach {
1730 let publisher = route.hops.iter().next().copied();
1731 let mut leaf_guard = leaf.lock();
1732
1733 if let Some(existing) = &leaf_guard.broadcast {
1736 let mut joined = None;
1737 let carrying = existing.broadcast.demand().is_used();
1738 if let Ok(mut s) = existing.state.write()
1739 && !s.closed
1740 {
1741 if same_publisher(s.publisher, publisher) {
1742 let id = s.next_route;
1743 s.next_route += 1;
1744 s.routes.push(FrontRoute {
1745 id,
1746 route: route.clone(),
1747 source: source.clone(),
1748 });
1749 s.reselect(carrying);
1750 joined = Some(id);
1751 } else if !route.announce || s.taints_a_reader(&route) {
1752 return Attach::Parked(existing.state.clone());
1753 } else {
1754 s.closed = true;
1762 tracing::warn!(broadcast = %ctx.full, "replacing a live broadcast from a different publisher");
1763 }
1764 }
1765 if let Some(id) = joined {
1766 let state = existing.state.clone();
1767 let broadcast = existing.broadcast.clone();
1768 drop(leaf_guard);
1769 sync_front(&state, &broadcast, leaf);
1770 return Attach::Ready(state, broadcast, id);
1771 }
1772 }
1773
1774 let announce = route.announce;
1776 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
1777 origin: ctx.origin.clone(),
1778 });
1779 let _ = broadcast.clone().set_route(route.clone());
1780 let state = kio::Producer::new(FrontState {
1781 path: ctx.full.clone(),
1782 self_origin: ctx.origin.id,
1783 publisher,
1784 next_route: 1,
1785 excluded: HashMap::new(),
1786 routes: vec![FrontRoute {
1787 id: 0,
1788 route,
1789 source: source.clone(),
1790 }],
1791 active: Some(0),
1792 linger: ctx.origin.linger,
1793 closed: false,
1794 });
1795
1796 if let Some(stale) = leaf_guard.broadcast.take()
1800 && stale.announced
1801 {
1802 leaf_guard.notify.lock().unannounce(&stale.path);
1803 }
1804 let entry = OriginBroadcast {
1805 path: ctx.full.clone(),
1806 broadcast: broadcast.clone(),
1807 state: state.clone(),
1808 announced: announce,
1809 };
1810 if entry.announced {
1811 leaf_guard
1812 .notify
1813 .lock()
1814 .announce(ctx.full, &broadcast.consume(), &state);
1815 }
1816 leaf_guard.broadcast = Some(entry);
1817 drop(leaf_guard);
1818
1819 web_async::spawn(run_front(
1820 state.clone(),
1821 broadcast.clone(),
1822 ctx.node.clone(),
1823 ctx.rest.clone(),
1824 ));
1825
1826 Attach::Ready(state, broadcast, 0)
1827}
1828
1829async fn run_front(
1832 state: kio::Producer<FrontState>,
1833 mut broadcast: broadcast::Producer,
1834 node: Lock<OriginNode>,
1835 rest: PathOwned,
1836) {
1837 enum Step {
1838 Serve(Arc<str>, super::resume::Producer),
1839 Changed,
1841 Expired,
1843 Closed,
1844 }
1845
1846 let linger = state.read().linger;
1847 let mut deadline = kio::time::Deadline::new();
1852
1853 loop {
1854 let empty = {
1855 let s = state.read();
1856 !s.closed && s.routes.is_empty()
1857 };
1858 deadline.set(match (empty, deadline.deadline()) {
1859 (true, None) => web_async::time::Instant::now().checked_add(linger),
1862 (true, at) => at,
1863 (false, _) => None,
1864 });
1865
1866 let step = {
1867 kio::wait(|waiter| {
1868 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1869 return Poll::Ready(Step::Serve(name, resume));
1870 }
1871 match state.poll(waiter, |s| {
1874 if s.closed || s.routes.is_empty() != empty {
1875 Poll::Ready(())
1876 } else {
1877 Poll::Pending
1878 }
1879 }) {
1880 Poll::Ready(Ok(guard)) => {
1881 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1882 }
1883 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1884 Poll::Pending => {}
1885 }
1886 deadline.poll(waiter).map(|_| Step::Expired)
1887 })
1888 .await
1889 };
1890
1891 match step {
1892 Step::Serve(name, resume) => {
1893 web_async::spawn(serve_track(state.clone(), name, resume));
1896 }
1897 Step::Changed => {}
1898 Step::Expired => {
1899 let close = {
1903 let Ok(mut s) = state.write() else { break };
1904 if !s.closed && s.routes.is_empty() {
1905 s.closed = true;
1906 true
1907 } else {
1908 false
1909 }
1910 };
1911 if close {
1912 break;
1913 }
1914 }
1915 Step::Closed => break,
1916 }
1917 }
1918
1919 broadcast.abort_spliced(Error::Dropped);
1921
1922 broadcast.finish();
1924
1925 node.lock().remove(&state, &rest);
1928}
1929
1930async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1943 enum Step {
1944 Closed,
1945 Splice(u64, broadcast::Consumer),
1946 Complete,
1947 Failed(Error),
1948 NoRoute,
1953 Idle,
1955 Demand,
1957 }
1958
1959 let mut serving: Option<(u64, track::Consumer)> = None;
1961 let mut spliced_edge: Option<u64> = None;
1967 let mut refused: HashSet<u64> = HashSet::new();
1972 let mut refusal: Option<Error> = None;
1973 let mut dead: HashSet<u64> = HashSet::new();
1978 let mut idle_since: Option<web_async::time::Instant> = None;
1980 let mut deadline = kio::time::Deadline::new();
1981
1982 loop {
1983 let serving_id = serving.as_ref().map(|(id, _)| *id);
1984
1985 {
1992 let s = state.read();
1993 refused.retain(|id| s.routes.iter().any(|r| r.id == *id));
1994 dead.retain(|id| s.routes.iter().any(|r| r.id == *id));
1995 let exhausted = !s.routes.is_empty()
1996 && s.serve_route(|id| refused.contains(&id) || dead.contains(&id))
1997 .is_none();
1998 if exhausted && dead.is_empty() {
1999 drop(s);
2000 let err = refusal.take().unwrap_or(Error::NotFound);
2001 tracing::debug!(name = %name, %err, "every source refused track; aborting");
2002 let _ = resume.abort(err);
2003 return;
2004 }
2005 }
2006
2007 let used = resume.is_used();
2019 idle_since = match (resume.is_spliced(), used) {
2020 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
2021 _ => None,
2022 };
2023 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
2024
2025 let step = {
2026 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2027 kio::wait(|waiter| {
2028 match state.poll(waiter, |s| {
2034 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
2035 if s.closed
2036 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
2037 {
2038 Poll::Ready(())
2039 } else {
2040 Poll::Pending
2041 }
2042 }) {
2043 Poll::Ready(Ok(guard)) => {
2044 if guard.closed {
2045 return Poll::Ready(Step::Closed);
2046 }
2047 let Some(next) = guard.serve_route(skip) else {
2048 return Poll::Ready(Step::NoRoute);
2049 };
2050 let source = guard
2051 .routes
2052 .iter()
2053 .find(|r| r.id == next)
2054 .expect("servable source in table")
2055 .source
2056 .clone();
2057 return Poll::Ready(Step::Splice(next, source));
2058 }
2059 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
2060 Poll::Pending => {}
2061 }
2062
2063 let edge = match used {
2068 true => resume.poll_unused(waiter),
2069 false => resume.poll_used(waiter),
2070 };
2071 if edge.is_ready() {
2072 return Poll::Ready(Step::Demand);
2073 }
2074
2075 if let Some((_, track)) = &serving
2078 && let Poll::Ready(result) = track.poll_complete(waiter)
2079 {
2080 return Poll::Ready(match result {
2081 Ok(()) => Step::Complete,
2082 Err(err) => Step::Failed(err),
2083 });
2084 }
2085
2086 deadline.poll(waiter).map(|_| Step::Idle)
2087 })
2088 .await
2089 };
2090
2091 match step {
2092 Step::Closed => return,
2094 Step::Complete => {
2095 let _ = resume.finish();
2096 return;
2097 }
2098 Step::Failed(err) => {
2099 if resume.latest() == spliced_edge
2107 && let Some(id) = serving_id
2108 {
2109 let closing = state
2110 .read()
2111 .routes
2112 .iter()
2113 .find(|r| r.id == id)
2114 .is_some_and(|r| r.source.is_closing());
2115 if closing {
2116 dead.insert(id);
2117 } else {
2118 refused.insert(id);
2119 refusal = Some(err);
2120 }
2121 }
2122 serving = None;
2123 }
2124 Step::Demand => {}
2126 Step::NoRoute => serving = None,
2132 Step::Idle => {
2133 if resume.release().is_err() {
2138 return;
2140 }
2141 serving = None;
2142 }
2143 Step::Splice(id, source) => {
2144 let attempt = match source.track(&name) {
2148 Ok(track) => {
2149 let query = track.info().into_inner();
2152 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2153 let info = kio::wait(|waiter| {
2154 if let Poll::Ready(result) = query.poll(waiter) {
2155 return Poll::Ready(Some(result));
2156 }
2157 match state.poll(waiter, |s| {
2158 if s.closed || s.serve_route(skip) != Some(id) {
2159 Poll::Ready(())
2160 } else {
2161 Poll::Pending
2162 }
2163 }) {
2164 Poll::Ready(_) => Poll::Ready(None),
2165 Poll::Pending => Poll::Pending,
2166 }
2167 })
2168 .await;
2169 match info {
2170 None => continue,
2172 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2175 Poll::Ready(Err(err)) => Err(err),
2176 _ => Ok(track),
2177 },
2178 Some(Err(err)) => Err(err),
2179 }
2180 }
2181 Err(err) => Err(err),
2182 };
2183
2184 match attempt {
2185 Ok(track) => {
2186 if let Err(err) = resume.takeover(&track) {
2187 let _ = resume.abort(err);
2192 return;
2193 }
2194 spliced_edge = resume.latest();
2206 serving = Some((id, track));
2207 }
2208 Err(_) if source.is_closing() => {
2212 dead.insert(id);
2213 serving = None;
2214 }
2215 Err(err) => {
2221 tracing::debug!(name = %name, source = id, %err, "source refused track");
2222 refused.insert(id);
2223 refusal = Some(err);
2224 serving = None;
2225 }
2226 }
2227 }
2228 }
2229 }
2230}
2231
2232#[derive(Default)]
2238struct OriginDynamicState {
2239 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2242
2243 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2249}
2250
2251#[derive(Default)]
2258struct PendingBroadcast {
2259 resolved: Option<Result<broadcast::Consumer, Error>>,
2260}
2261
2262pub struct Dynamic {
2273 info: Origin,
2274 root: PathOwned,
2275 state: kio::Shared<OriginDynamicState>,
2276}
2277
2278impl Clone for Dynamic {
2279 fn clone(&self) -> Self {
2280 self.state.lock().requests.add_handler();
2284
2285 Self {
2286 info: self.info,
2287 root: self.root.clone(),
2288 state: self.state.clone(),
2289 }
2290 }
2291}
2292
2293impl Dynamic {
2294 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2295 state.lock().requests.add_handler();
2296
2297 Self { info, root, state }
2298 }
2299
2300 pub fn info(&self) -> &Origin {
2302 &self.info
2303 }
2304
2305 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2307 let mut state = ready!(self.state.poll(waiter, |state| {
2308 if state.requests.has_queued() {
2309 Poll::Ready(())
2310 } else {
2311 Poll::Pending
2312 }
2313 }));
2314
2315 let path = state.requests.pop().expect("predicate guaranteed a request");
2316 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2322 Poll::Ready(Ok(Request {
2323 path,
2324 producer,
2325 state: self.state.clone(),
2326 }))
2327 }
2328
2329 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2332 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2333 }
2334
2335 pub fn root(&self) -> &Path<'_> {
2337 &self.root
2338 }
2339}
2340
2341impl Drop for Dynamic {
2342 fn drop(&mut self) {
2343 let mut state = self.state.lock();
2346 if state.requests.remove_handler() {
2347 state.requests.drain_queued();
2351 }
2352 }
2353}
2354
2355pub struct Request {
2362 path: PathOwned,
2364
2365 producer: kio::Producer<PendingBroadcast>,
2368
2369 state: kio::Shared<OriginDynamicState>,
2371}
2372
2373impl Request {
2374 pub fn path(&self) -> &Path<'_> {
2376 &self.path
2377 }
2378
2379 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2385 let broadcast = broadcast.consume();
2386
2387 let resolved = {
2393 let mut state = self.state.lock();
2394 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2395 state
2396 .requests
2397 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2398 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2399 };
2400
2401 if let Ok(mut pending) = self.producer.write() {
2402 pending.resolved = Some(Ok(resolved));
2403 }
2404 }
2406
2407 pub fn reject(self, err: Error) {
2409 self.state
2410 .lock()
2411 .requests
2412 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2413 if let Ok(mut state) = self.producer.write() {
2414 state.resolved = Some(Err(err));
2415 }
2416 }
2417}
2418
2419impl Drop for Request {
2420 fn drop(&mut self) {
2421 self.state
2429 .lock()
2430 .requests
2431 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2432 }
2433}
2434
2435pub struct Requesting {
2442 inner: RequestState,
2443 stats: stats::Scope,
2446}
2447
2448enum RequestState {
2449 Ready(broadcast::Consumer),
2451 Failed(Error),
2454 Pending(kio::Consumer<PendingBroadcast>),
2456}
2457
2458impl Requesting {
2459 fn ready(broadcast: broadcast::Consumer) -> Self {
2460 Self {
2461 inner: RequestState::Ready(broadcast),
2462 stats: stats::Scope::default(),
2463 }
2464 }
2465
2466 fn failed(error: Error) -> Self {
2467 Self {
2468 inner: RequestState::Failed(error),
2469 stats: stats::Scope::default(),
2470 }
2471 }
2472
2473 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2474 Self {
2475 inner: RequestState::Pending(consumer),
2476 stats: stats::Scope::default(),
2477 }
2478 }
2479
2480 fn with_stats(mut self, scope: stats::Scope) -> Self {
2481 self.stats = scope;
2482 self
2483 }
2484
2485 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2487 match &self.inner {
2488 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2489 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2490 RequestState::Pending(consumer) => Poll::Ready(
2491 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2492 Some(result) => Poll::Ready(result.clone()),
2493 None => Poll::Pending,
2494 })) {
2495 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2496 Err(_closed) => Err(Error::Unroutable),
2498 },
2499 ),
2500 }
2501 }
2502}
2503
2504impl kio::Pollable for Requesting {
2505 type Output = Result<broadcast::Consumer, Error>;
2506
2507 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2508 self.poll_ok(waiter)
2509 }
2510}
2511
2512pub trait Consume<T> {
2520 fn consume(&self) -> T;
2522}
2523
2524impl<T, U: Consume<T>> Consume<T> for &U {
2525 fn consume(&self) -> T {
2526 (**self).consume()
2527 }
2528}
2529
2530impl Consume<Consumer> for Producer {
2531 fn consume(&self) -> Consumer {
2532 Consumer::new(
2536 self.info,
2537 self.root.clone(),
2538 self.nodes.clone(),
2539 self.dynamic.clone(),
2540 stats::Session::default(),
2541 )
2542 }
2543}
2544
2545impl Consume<Consumer> for Consumer {
2546 fn consume(&self) -> Consumer {
2547 self.clone()
2548 }
2549}
2550
2551impl Consume<broadcast::Consumer> for broadcast::Producer {
2552 fn consume(&self) -> broadcast::Consumer {
2553 self.consume()
2555 }
2556}
2557
2558impl Consume<broadcast::Consumer> for broadcast::Consumer {
2559 fn consume(&self) -> broadcast::Consumer {
2560 self.clone()
2561 }
2562}
2563
2564impl Consume<track::Consumer> for track::Producer {
2565 fn consume(&self) -> track::Consumer {
2566 self.consume()
2567 }
2568}
2569
2570impl Consume<track::Consumer> for track::Consumer {
2571 fn consume(&self) -> track::Consumer {
2572 self.clone()
2573 }
2574}
2575
2576#[derive(Clone)]
2582pub struct Consumer {
2583 info: Origin,
2585 nodes: OriginNodes,
2586
2587 root: PathOwned,
2589
2590 dynamic: kio::Shared<OriginDynamicState>,
2593
2594 stats: stats::Session,
2598
2599 exclude: Option<Origin>,
2603}
2604
2605impl std::ops::Deref for Consumer {
2606 type Target = Origin;
2607
2608 fn deref(&self) -> &Self::Target {
2609 &self.info
2610 }
2611}
2612
2613impl Consumer {
2614 fn new(
2615 info: Origin,
2616 root: PathOwned,
2617 nodes: OriginNodes,
2618 dynamic: kio::Shared<OriginDynamicState>,
2619 stats: stats::Session,
2620 ) -> Self {
2621 Self {
2622 info,
2623 nodes,
2624 root,
2625 dynamic,
2626 stats,
2627 exclude: None,
2628 }
2629 }
2630
2631 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2636 self.exclude = Some(peer);
2637 self
2638 }
2639
2640 pub fn with_stats(mut self, session: stats::Session) -> Self {
2644 self.stats = session;
2645 self
2646 }
2647
2648 fn untagged(&self) -> Self {
2652 Self {
2653 stats: stats::Session::default(),
2654 ..self.clone()
2655 }
2656 }
2657
2658 pub(crate) fn empty(&self) -> Self {
2663 Self {
2664 info: self.info,
2665 nodes: OriginNodes { nodes: Vec::new() },
2666 root: self.root.clone(),
2667 dynamic: self.dynamic.clone(),
2668 stats: self.stats.clone(),
2669 exclude: self.exclude,
2670 }
2671 }
2672
2673 pub fn announced(&self) -> AnnounceConsumer {
2680 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone(), self.exclude)
2681 }
2682
2683 pub fn consume(&self) -> Self {
2685 self.clone()
2686 }
2687
2688 fn resolve(&self, path: impl AsPath) -> Resolved {
2697 let path = path.as_path();
2698 let Some((root, rest)) = self.nodes.get(&path) else {
2699 return Resolved::Missing;
2700 };
2701 let state = root.lock();
2702 state.resolve_broadcast(&rest, self.exclude)
2703 }
2704
2705 #[cfg(test)]
2707 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2708 match self.resolve(path) {
2709 Resolved::Found(broadcast) => Some(broadcast),
2710 Resolved::Excluded | Resolved::Missing => None,
2711 }
2712 }
2713
2714 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2726 let path = path.as_path();
2727
2728 let consumer = self.scope(std::slice::from_ref(&path))?;
2730
2731 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2735 return None;
2736 }
2737
2738 let mut announced = consumer.untagged().announced();
2742 let scope = self.stats.egress(self.root.join(&path).to_owned());
2743 loop {
2744 let OriginAnnounce {
2745 path: announced_path,
2746 broadcast,
2747 } = announced.next().await?;
2748 if announced_path.as_path() == path
2750 && let Some(broadcast) = broadcast
2751 {
2752 return Some(broadcast.with_stats(scope));
2753 }
2754 }
2755 }
2756
2757 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2763 let prefixes = PathPrefixes::new(prefixes);
2764 Some(Consumer {
2765 info: self.info,
2766 root: self.root.clone(),
2767 nodes: self.nodes.select(&prefixes)?,
2768 dynamic: self.dynamic.clone(),
2769 stats: self.stats.clone(),
2770 exclude: self.exclude,
2771 })
2772 }
2773
2774 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2793 let path = path.as_path();
2794
2795 let absolute = self.root.join(&path).to_owned();
2799 let scope = self.stats.egress(&absolute);
2800
2801 match self.resolve(&path) {
2807 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
2808 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
2809 Resolved::Missing => {}
2810 }
2811
2812 let mut state = self.dynamic.lock();
2813
2814 if let Some(weak) = state.served.get(&absolute) {
2818 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2819 }
2820
2821 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2824 producer.consume()
2825 } else {
2826 let producer = kio::Producer::<PendingBroadcast>::default();
2827 let consumer = producer.consume();
2828 if state.requests.insert(absolute, producer).is_err() {
2829 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2830 }
2831 consumer
2832 };
2833
2834 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2835 }
2836
2837 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2842 let prefix = prefix.as_path();
2843
2844 Some(Self {
2845 info: self.info,
2846 root: self.root.join(&prefix).to_owned(),
2847 nodes: self.nodes.root(&prefix)?,
2848 dynamic: self.dynamic.clone(),
2849 stats: self.stats.clone(),
2850 exclude: self.exclude,
2851 })
2852 }
2853
2854 pub fn root(&self) -> &Path<'_> {
2856 &self.root
2857 }
2858
2859 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2862 self.nodes.nodes.iter().map(|(root, _)| root)
2863 }
2864
2865 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2867 self.root.join(path)
2868 }
2869}
2870
2871#[derive(Clone)]
2876pub struct AnnounceProducer {
2877 nodes: OriginNodes,
2878 root: PathOwned,
2879}
2880
2881impl AnnounceProducer {
2882 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2883 Self { nodes, root }
2884 }
2885
2886 pub fn consume(&self) -> AnnounceConsumer {
2892 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default(), None)
2895 }
2896
2897 pub fn root(&self) -> &Path<'_> {
2899 &self.root
2900 }
2901}
2902
2903pub struct AnnounceConsumer {
2908 id: ConsumerId,
2909 nodes: OriginNodes,
2910 root: PathOwned,
2911
2912 state: kio::Producer<OriginConsumerState>,
2915
2916 stats: stats::Session,
2919
2920 guards: HashMap<PathOwned, stats::Announce>,
2924}
2925
2926impl AnnounceConsumer {
2927 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session, exclude: Option<Origin>) -> Self {
2928 let state = kio::Producer::<OriginConsumerState>::default();
2929 let id = ConsumerId::new();
2930
2931 for (_, node) in &nodes.nodes {
2932 let notify = AnnounceConsumerNotify {
2933 root: root.clone(),
2934 state: state.clone(),
2935 exclude,
2936 };
2937 node.lock().consume(id, notify);
2938 }
2939
2940 Self {
2941 id,
2942 nodes,
2943 root,
2944 state,
2945 stats,
2946 guards: HashMap::new(),
2947 }
2948 }
2949
2950 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2956 let OriginAnnounce { path, broadcast } = update;
2957 let absolute = self.root.join(&path).to_owned();
2958 match broadcast {
2959 Some(broadcast) => {
2960 let scope = self.stats.egress(&absolute);
2961 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2962 OriginAnnounce {
2963 path,
2964 broadcast: Some(broadcast.with_stats(scope)),
2965 }
2966 }
2967 None => {
2968 self.guards.remove(&absolute);
2969 OriginAnnounce { path, broadcast: None }
2970 }
2971 }
2972 }
2973
2974 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2981 kio::wait(|waiter| self.poll_next(waiter)).await
2982 }
2983
2984 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2990 let update = {
2991 let mut state = match ready!(self.state.poll(waiter, |state| {
2992 if state.pending.is_empty() {
2993 Poll::Pending
2994 } else {
2995 Poll::Ready(())
2996 }
2997 })) {
2998 Ok(state) => state,
2999 Err(_) => return Poll::Ready(None),
3001 };
3002 state.take().expect("predicate guaranteed an update")
3003 };
3004 Poll::Ready(Some(self.attribute(update)))
3005 }
3006
3007 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
3012 let update = self.state.write().ok()?.take()?;
3013 Some(self.attribute(update))
3014 }
3015
3016 pub fn is_closed(&self) -> bool {
3018 self.state.write().is_err()
3019 }
3020
3021 pub fn root(&self) -> &Path<'_> {
3023 &self.root
3024 }
3025
3026 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
3028 self.root.join(path)
3029 }
3030}
3031
3032impl Drop for AnnounceConsumer {
3033 fn drop(&mut self) {
3034 for (_, root) in &self.nodes.nodes {
3035 root.lock().unconsume(self.id);
3036 }
3037 }
3038}
3039
3040#[cfg(test)]
3041use futures::FutureExt;
3042
3043#[cfg(test)]
3044#[allow(missing_docs)] impl AnnounceConsumer {
3046 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3047 let expected = expected.as_path();
3048 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3049 assert_eq!(announce.path, expected, "wrong path");
3050 let announced = announce.broadcast.expect("should be an active announce");
3051 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3052 }
3053
3054 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3058 let expected = expected.as_path();
3059 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3060 assert_eq!(announce.path, expected, "wrong path");
3061 announce.broadcast.expect("should be an active announce")
3062 }
3063
3064 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3065 let expected = expected.as_path();
3066 let announce = self.try_next().expect("no next");
3067 assert_eq!(announce.path, expected, "wrong path");
3068 let announced = announce.broadcast.expect("should be an active announce");
3069 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3070 }
3071
3072 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3074 let expected = expected.as_path();
3075 let announce = self.try_next().expect("no next");
3076 assert_eq!(announce.path, expected, "wrong path");
3077 announce.broadcast.expect("should be an active announce")
3078 }
3079
3080 pub fn assert_next_none(&mut self, expected: impl AsPath) {
3081 let expected = expected.as_path();
3082 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3083 assert_eq!(announce.path, expected, "wrong path");
3084 assert!(announce.broadcast.is_none(), "should be unannounced");
3085 }
3086
3087 pub fn assert_next_wait(&mut self) {
3088 if let Some(res) = self.next().now_or_never() {
3089 panic!("next should block: got {:?}", res.map(|a| a.path));
3090 }
3091 }
3092
3093 }
3102
3103#[cfg(test)]
3104mod tests {
3105 use crate::coding::Decode;
3106 use crate::group;
3107
3108 use super::*;
3109
3110 fn announce() -> broadcast::Route {
3112 broadcast::Route::new().with_announce(true)
3113 }
3114
3115 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
3121 let name = Path::new(name);
3122 let peer_key = fnv_key(&name, [peer]);
3123 (100u64..)
3124 .map(|id| Origin::new(id).unwrap())
3125 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
3126 .unwrap()
3127 }
3128
3129 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
3132 let source = broadcast::Info::new().produce().consume();
3133 FrontState {
3134 path: Path::new("test").to_owned(),
3135 self_origin,
3136 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
3137 next_route: routes.len() as u64,
3138 excluded: HashMap::new(),
3139 routes: routes
3140 .into_iter()
3141 .enumerate()
3142 .map(|(id, route)| FrontRoute {
3143 id: id as u64,
3144 route,
3145 source: source.clone(),
3146 })
3147 .collect(),
3148 active: Some(0),
3149 linger: Duration::ZERO,
3150 closed: false,
3151 }
3152 }
3153
3154 fn sibling_route(peer: Origin) -> broadcast::Route {
3157 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3158 announce().with_hops(hops)
3159 }
3160
3161 fn upstream_route(cost: u64) -> broadcast::Route {
3163 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3164 announce().with_hops(hops).with_cost(cost)
3165 }
3166
3167 #[test]
3171 fn test_carrying_gate_keys() {
3172 let peer = Origin::new(3).unwrap();
3173
3174 let mut lost = front_state(
3176 origin_keyed("test", peer, false),
3177 vec![upstream_route(10), sibling_route(peer)],
3178 );
3179 lost.reselect(true);
3180 assert_eq!(
3181 lost.active,
3182 Some(0),
3183 "carrying front re-parented onto a higher-keyed peer"
3184 );
3185 lost.reselect(false);
3186 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3187
3188 let mut won = front_state(
3190 origin_keyed("test", peer, true),
3191 vec![upstream_route(10), sibling_route(peer)],
3192 );
3193 won.reselect(true);
3194 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3195 }
3196
3197 #[test]
3202 fn test_carrying_gate_symmetric_race() {
3203 let a = Origin::new(1).unwrap();
3204 let b = Origin::new(2).unwrap();
3205
3206 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3207 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3208 a_view.reselect(true);
3209 b_view.reselect(true);
3210
3211 let a_moved = a_view.active == Some(1);
3212 let b_moved = b_view.active == Some(1);
3213 assert!(
3214 a_moved != b_moved,
3215 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3216 );
3217 }
3218
3219 #[test]
3224 fn test_carrying_switches_to_benign_routes() {
3225 let peer = Origin::new(3).unwrap();
3226 let lost = origin_keyed("test", peer, false);
3227
3228 let mut forwarder = sibling_route(peer).with_cost(4);
3230 forwarder.advertised = 4;
3231 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3232 state.reselect(true);
3233 assert_eq!(
3234 state.active,
3235 Some(1),
3236 "a cheaper forwarder path must win while carrying"
3237 );
3238
3239 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3241 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3242 state.reselect(true);
3243 assert_eq!(
3244 state.active,
3245 Some(1),
3246 "a direct publisher route must win while carrying"
3247 );
3248
3249 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3254 state.reselect(true);
3255 assert_eq!(
3256 state.active,
3257 Some(1),
3258 "a reconnect on an identical chain must win while carrying"
3259 );
3260 }
3261
3262 #[test]
3265 fn test_carrying_gate_ignores_unannounced_incumbent() {
3266 let peer = Origin::new(3).unwrap();
3267 let unannounced = upstream_route(10).with_announce(false);
3268 let mut state = front_state(
3269 origin_keyed("test", peer, false),
3270 vec![unannounced, sibling_route(peer)],
3271 );
3272 state.reselect(true);
3273 assert_eq!(
3274 state.active,
3275 Some(1),
3276 "an unannounced incumbent must always be displaced"
3277 );
3278 }
3279
3280 #[test]
3285 fn test_reflection_through_an_exposed_peer_cannot_take_over() {
3286 let peer = Origin::new(42).unwrap();
3287 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3288 let reflected = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3289
3290 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3291
3292 assert!(!state.taints_a_reader(&reflected));
3294
3295 *state.excluded.entry(peer).or_default() += 1;
3298 assert!(state.taints_a_reader(&reflected));
3299 }
3300
3301 #[test]
3304 fn test_rival_publisher_through_another_peer_still_takes_over() {
3305 let peer = Origin::new(42).unwrap();
3306 let elsewhere = Origin::new(43).unwrap();
3307 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3308 let rival = announce().with_hops(OriginList::try_from(vec![Origin::UNKNOWN, elsewhere]).unwrap());
3309
3310 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3311 *state.excluded.entry(peer).or_default() += 1;
3312
3313 assert!(!state.taints_a_reader(&rival), "only the peer we feed is a reflection");
3314 }
3315
3316 #[test]
3321 fn test_opaque_peer_understates_its_depth() {
3322 let us = Origin::new(1).unwrap();
3323
3324 let direct = || {
3326 announce()
3327 .with_hops(OriginList::try_from(vec![Origin::new(7).unwrap(), Origin::new(8).unwrap()]).unwrap())
3328 .with_cost(2)
3329 };
3330 let opaque = |cost| {
3333 announce()
3334 .with_hops(OriginList::try_from(vec![Origin::new(42).unwrap()]).unwrap())
3335 .with_cost(cost)
3336 };
3337
3338 let mut state = front_state(us, vec![direct(), opaque(1)]);
3340 state.reselect(false);
3341 assert_eq!(
3342 state.active,
3343 Some(1),
3344 "an unpriced opaque link out-ranks a shorter real path"
3345 );
3346
3347 let mut state = front_state(us, vec![direct(), opaque(16)]);
3349 state.reselect(false);
3350 assert_eq!(
3351 state.active,
3352 Some(0),
3353 "pricing the opaque link restores the intended order"
3354 );
3355 }
3356
3357 async fn settle() {
3360 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3361 }
3362
3363 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3366 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3367 .await
3368 .expect("timed out waiting for a track request")
3369 .expect("source closed");
3370 assert_eq!(request.name(), name, "unexpected track dispatched");
3371 request.accept(None)
3372 }
3373
3374 #[tokio::test]
3378 async fn test_stats_tagged_end_to_end() {
3379 use crate::Timestamp;
3380 use crate::stats::{Config, Registry, Tier};
3381 use bytes::Bytes;
3382
3383 tokio::time::pause();
3384
3385 let registry = Registry::new(Config::new());
3386 let ctx = registry.tier(Tier::default()).session("acme");
3387
3388 let origin = Origin::random().produce();
3389 let ingress = origin.clone().with_stats(ctx.clone());
3390 let egress = origin.consume().with_stats(ctx.clone());
3391
3392 let mut announced = egress.announced();
3395
3396 let source = ingress.create_broadcast("demo", announce()).unwrap();
3398 let mut dynamic = source.dynamic();
3399 settle().await;
3400 settle().await;
3401
3402 let update = announced.next().await.unwrap();
3404 assert_eq!(update.path.as_str(), "demo");
3405 let broadcast = update.broadcast.unwrap();
3406
3407 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3409 let mut producer = accept_track(&mut dynamic, "video").await;
3410 settle().await;
3411 let mut sub = subscribing.await.unwrap();
3412
3413 let mut group = producer.append_group().unwrap();
3415 group
3416 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3417 .unwrap();
3418 group
3419 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3420 .unwrap();
3421 group.finish().unwrap();
3422
3423 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3425 let mut frames = 0;
3426 while let Some(frame) = group_c.read_frame().await.unwrap() {
3427 assert_eq!(frame.payload.len(), 5);
3428 frames += 1;
3429 }
3430 assert_eq!(frames, 2);
3431 settle().await;
3432
3433 let report = registry.report();
3434 let entry = report
3435 .traffic
3436 .iter()
3437 .find(|e| e.path.as_str() == "demo")
3438 .expect("demo tracked");
3439 let path_len = "demo".len() as u64;
3440
3441 let egress = &entry.publisher;
3443 assert_eq!(egress.announced, 1, "one egress announce");
3444 assert_eq!(egress.announced_bytes, path_len);
3445 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3446 assert_eq!(egress.broadcasts, 1, "one viewer");
3447 assert_eq!(egress.groups, 1);
3448 assert_eq!(egress.frames, 2);
3449 assert_eq!(egress.bytes, 10);
3450 assert_eq!(egress.fetches, 0);
3451
3452 let ingress = &entry.subscriber;
3454 assert_eq!(ingress.announced, 1, "one ingress announce");
3455 assert_eq!(ingress.announced_bytes, path_len);
3456 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3457 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3458 assert_eq!(ingress.groups, 1);
3459 assert_eq!(ingress.frames, 2);
3460 assert_eq!(ingress.bytes, 10);
3461
3462 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3464 let _ = fetched;
3465 settle().await;
3466 let report = registry.report();
3467 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3468 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3469 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3470 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3471 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3475 }
3476
3477 #[tokio::test]
3482 async fn test_stats_read_frame_counts_once() {
3483 use crate::Timestamp;
3484 use crate::stats::{Config, Registry, Tier};
3485 use bytes::Bytes;
3486
3487 tokio::time::pause();
3488
3489 let registry = Registry::new(Config::new());
3490 let ctx = registry.tier(Tier::default()).session("acme");
3491
3492 let origin = Origin::random().produce();
3493 let ingress = origin.clone().with_stats(ctx.clone());
3494 let egress = origin.consume().with_stats(ctx.clone());
3495
3496 let mut announced = egress.announced();
3497 let source = ingress.create_broadcast("demo", announce()).unwrap();
3498 let mut dynamic = source.dynamic();
3499 settle().await;
3500 settle().await;
3501
3502 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3503 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3504 let mut producer = accept_track(&mut dynamic, "video").await;
3505 settle().await;
3506 let mut sub = subscribing.await.unwrap();
3507
3508 producer
3510 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3511 .unwrap();
3512
3513 let frame = sub.read_frame().await.unwrap().expect("frame");
3514 assert_eq!(frame.payload.len(), 5);
3515 settle().await;
3516
3517 let report = registry.report();
3518 let entry = report
3519 .traffic
3520 .iter()
3521 .find(|e| e.path.as_str() == "demo")
3522 .expect("demo tracked");
3523 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3524 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3525 assert_eq!(
3526 entry.publisher.bytes, 5,
3527 "payload counted once, not zero and not doubled"
3528 );
3529 }
3530
3531 #[tokio::test]
3535 async fn test_stats_datagrams_counted_both_sides() {
3536 use crate::Timestamp;
3537 use crate::stats::{Config, Registry, Tier};
3538
3539 tokio::time::pause();
3540
3541 let registry = Registry::new(Config::new());
3542 let ctx = registry.tier(Tier::default()).session("acme");
3543
3544 let origin = Origin::random().produce();
3545 let ingress = origin.clone().with_stats(ctx.clone());
3546 let egress = origin.consume().with_stats(ctx.clone());
3547
3548 let mut announced = egress.announced();
3549 let source = ingress.create_broadcast("demo", announce()).unwrap();
3550 let mut dynamic = source.dynamic();
3551 settle().await;
3552 settle().await;
3553
3554 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3555 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3556 let mut producer = accept_track(&mut dynamic, "video").await;
3557 settle().await;
3558 let mut sub = subscribing.await.unwrap();
3559
3560 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3561 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3562 assert_eq!(&datagram.payload[..], b"hello");
3563 settle().await;
3564
3565 let report = registry.report();
3566 let entry = report
3567 .traffic
3568 .iter()
3569 .find(|e| e.path.as_str() == "demo")
3570 .expect("demo tracked");
3571
3572 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3573 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3574 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3575 assert_eq!(traffic.frames, 1, "{side}: one frame");
3576 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3577 }
3578 }
3579
3580 #[test]
3581 fn origin_rejects_reserved_ids() {
3582 assert!(Origin::new(0).is_err());
3583 assert!(Origin::new(1u64 << 62).is_err());
3584 assert_eq!(Origin::new(1).unwrap().id(), 1);
3585
3586 let mut zero = [0u8].as_slice();
3587 assert_eq!(
3588 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3589 Origin::UNKNOWN
3590 );
3591 }
3592
3593 #[test]
3594 fn origin_list_push_fails_at_limit() {
3595 let mut list = OriginList::new();
3596 for _ in 0..MAX_HOPS {
3597 list.push(Origin::random()).unwrap();
3598 }
3599 assert_eq!(list.len(), MAX_HOPS);
3600 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3601 }
3602
3603 #[test]
3604 fn origin_list_replace_first() {
3605 let mut list = OriginList::new();
3606 for _ in 0..3 {
3607 list.push(Origin::UNKNOWN).unwrap();
3608 }
3609
3610 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3612 assert_eq!(
3613 list.as_slice(),
3614 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3615 );
3616
3617 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3619 assert_eq!(list.len(), 3);
3620 }
3621
3622 #[test]
3623 fn origin_list_try_from_vec_enforces_limit() {
3624 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3625 assert!(OriginList::try_from(under).is_ok());
3626
3627 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3628 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3629 }
3630
3631 #[tokio::test]
3632 async fn test_announce() {
3633 tokio::time::pause();
3634
3635 let origin = Origin::random().produce();
3636
3637 let mut consumer1 = origin.consume().announced();
3638 consumer1.assert_next_wait();
3639
3640 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3642 settle().await;
3643
3644 consumer1.assert_next_some("test1");
3645 consumer1.assert_next_wait();
3646
3647 let mut consumer2 = origin.consume().announced();
3650
3651 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3653 settle().await;
3654
3655 consumer1.assert_next_some("test2");
3656 consumer1.assert_next_wait();
3657
3658 consumer2.assert_next_some("test1");
3659 consumer2.assert_next_some("test2");
3660 consumer2.assert_next_wait();
3661
3662 broadcast1.finish();
3664 settle().await;
3665
3666 consumer1.assert_next_none("test1");
3668 consumer2.assert_next_none("test1");
3669 consumer1.assert_next_wait();
3670 consumer2.assert_next_wait();
3671
3672 let mut consumer3 = origin.consume().announced();
3674 consumer3.assert_next_some("test2");
3675 consumer3.assert_next_wait();
3676
3677 broadcast2.finish();
3678 settle().await;
3679
3680 consumer1.assert_next_none("test2");
3681 consumer2.assert_next_none("test2");
3682 consumer3.assert_next_none("test2");
3683 }
3684
3685 #[tokio::test]
3689 async fn test_duplicate() {
3690 tokio::time::pause();
3691
3692 let origin = Origin::random().produce();
3693 let consumer = origin.consume();
3694 let mut announced = consumer.announced();
3695
3696 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3697 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3698 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3699 settle().await;
3700 assert!(consumer.get_broadcast("test").is_some());
3701
3702 announced.assert_next_some("test");
3703 announced.assert_next_wait();
3704
3705 broadcast2.finish();
3707 settle().await;
3708 assert!(consumer.get_broadcast("test").is_some());
3709 announced.assert_next_wait();
3710
3711 broadcast1.finish();
3713 settle().await;
3714 assert!(consumer.get_broadcast("test").is_some());
3715 announced.assert_next_wait();
3716
3717 broadcast3.finish();
3719 settle().await;
3720 assert!(consumer.get_broadcast("test").is_none());
3721
3722 announced.assert_next_none("test");
3723 announced.assert_next_wait();
3724 }
3725
3726 #[tokio::test]
3729 async fn test_route_failover() {
3730 tokio::time::pause();
3731
3732 let origin = Origin::random().produce();
3733 let consumer = origin.consume();
3734 let mut announced = consumer.announced();
3735
3736 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3739 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3740
3741 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3743 let mut dynamic_a = source_a.dynamic();
3744 settle().await;
3745 settle().await;
3746 let broadcast = consumer.request_broadcast("test").await.unwrap();
3747 announced.assert_next_some("test");
3748
3749 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3751 let mut dynamic_b = source_b.dynamic();
3752 settle().await;
3753 settle().await;
3754 announced.assert_next_wait();
3755
3756 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3758 let mut producer = accept_track(&mut dynamic_a, "video").await;
3759 settle().await;
3760 dynamic_b.assert_no_request();
3761
3762 let mut sub = subscribing.await.unwrap();
3763 sub.assert_no_group();
3766 assert_eq!(producer.subscription().unwrap().group_start, None);
3767
3768 producer.append_group().unwrap();
3769 producer.append_group().unwrap();
3770 assert_eq!(sub.assert_group().sequence, 0);
3771 assert_eq!(sub.assert_group().sequence, 1);
3772
3773 producer.abort(Error::Dropped).unwrap();
3777 source_a.abort(Error::Dropped).unwrap();
3778 drop(dynamic_a);
3779 settle().await;
3780 announced.assert_next_wait();
3781
3782 let mut producer = accept_track(&mut dynamic_b, "video").await;
3785 settle().await;
3786 sub.assert_no_group();
3787 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3788 producer.create_group(group::Info { sequence: 1 }).unwrap();
3789 producer.create_group(group::Info { sequence: 2 }).unwrap();
3790 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3791 sub.assert_not_closed();
3792 }
3793
3794 #[tokio::test]
3797 async fn test_broadcast_route_watch() {
3798 let mut producer = broadcast::Info::new().produce();
3799 let mut consumer = producer.consume();
3800
3801 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3803
3804 producer.set_route(broadcast::Route::default()).unwrap();
3806 assert!(consumer.route_changed().now_or_never().is_none());
3807
3808 let mut hops = OriginList::new();
3809 hops.push(Origin::new(7).unwrap()).unwrap();
3810 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3811 producer.set_route(route.clone()).unwrap();
3812 assert_eq!(consumer.route_changed().await.unwrap(), route);
3813
3814 let mut fresh = producer.consume();
3816 assert_eq!(fresh.route_changed().await.unwrap(), route);
3817
3818 drop(producer);
3819 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3820 }
3821
3822 #[tokio::test]
3826 async fn test_route_cost_update() {
3827 tokio::time::pause();
3828
3829 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3833 let consumer = origin.consume();
3834 let mut announced = consumer.announced();
3835
3836 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3839 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3840
3841 let mut source_a = origin
3843 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3844 .unwrap();
3845 let mut dynamic_a = source_a.dynamic();
3846 settle().await;
3847 let broadcast = consumer.request_broadcast("test").await.unwrap();
3848 announced.assert_next_some("test");
3849
3850 let mut watch = broadcast.clone();
3851 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3852
3853 let mut source_b = origin
3854 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3855 .unwrap();
3856 let mut dynamic_b = source_b.dynamic();
3857 settle().await;
3858 assert!(
3859 watch.route_changed().now_or_never().is_none(),
3860 "a losing standby must not change the advertised route"
3861 );
3862
3863 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3865 let mut producer = accept_track(&mut dynamic_a, "video").await;
3866 settle().await;
3867 let mut sub = subscribing.await.unwrap();
3868 producer.append_group().unwrap();
3869 assert_eq!(sub.assert_group().sequence, 0);
3870
3871 source_a
3874 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3875 .unwrap();
3876 settle().await;
3877 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3878 announced.assert_next_wait();
3879
3880 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3881 settle().await;
3882 sub.assert_no_group();
3885 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3886 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3887 assert_eq!(sub.assert_group().sequence, 1);
3888 sub.assert_not_closed();
3889
3890 source_b
3892 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3893 .unwrap();
3894 settle().await;
3895 let advertised = watch.route_changed().await.unwrap();
3896 assert_eq!(advertised.hops, hops_b);
3897 assert_eq!(advertised.cost, 5);
3898 announced.assert_next_wait();
3899 }
3900
3901 #[tokio::test]
3904 async fn test_completed_track_survives_route_churn() {
3905 tokio::time::pause();
3906
3907 let origin = Origin::random().produce();
3908 let consumer = origin.consume();
3909
3910 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3912 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3913
3914 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3915 let mut dynamic_a = source_a.dynamic();
3916 settle().await;
3917 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3918 let mut dynamic_b = source_b.dynamic();
3919 settle().await;
3920 settle().await;
3921 let broadcast = consumer.request_broadcast("test").await.unwrap();
3922
3923 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3925 let mut producer = accept_track(&mut dynamic_a, "video").await;
3926 settle().await;
3927 let mut sub = subscribing.await.unwrap();
3928 producer.append_group().unwrap();
3929 assert_eq!(sub.assert_group().sequence, 0);
3930 producer.finish().unwrap();
3931 drop(producer);
3932 settle().await;
3933 sub.assert_closed();
3934
3935 source_a.abort(Error::Dropped).unwrap();
3937 drop(dynamic_a);
3938 settle().await;
3939 dynamic_b.assert_no_request();
3940
3941 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3943 late.assert_closed();
3944 }
3945
3946 #[tokio::test]
3950 async fn test_refused_track_aborts_instantly() {
3951 tokio::time::pause();
3952
3953 let origin = Origin::random().produce();
3954 let consumer = origin.consume();
3955
3956 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3957 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3958 let mut dynamic = source.dynamic();
3959 settle().await;
3960 settle().await;
3961 let broadcast = consumer.request_broadcast("test").await.unwrap();
3962
3963 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3964 let request = dynamic.requested_track().await.unwrap();
3965 request.reject(Error::NotFound);
3966 settle().await;
3967
3968 assert!(matches!(subscribing.await, Err(Error::NotFound)));
3970 dynamic.assert_no_request();
3971 }
3972
3973 #[tokio::test]
3978 async fn test_stale_rejection_does_not_abort_a_handover() {
3979 tokio::time::pause();
3980
3981 let origin = Origin::random().produce();
3982 let consumer = origin.consume();
3983
3984 let publisher = Origin::new(1).unwrap();
3985 let peer = Origin::new(5).unwrap();
3986 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
3987 let local = OriginList::try_from(vec![publisher]).unwrap();
3988
3989 let source_remote = origin
3991 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
3992 .unwrap();
3993 let mut dynamic_remote = source_remote.dynamic();
3994 settle().await;
3995 settle().await;
3996 let broadcast = consumer.request_broadcast("test").await.unwrap();
3997 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3998 let request_remote = dynamic_remote.requested_track().await.unwrap();
3999
4000 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4003 let mut dynamic_local = source_local.dynamic();
4004 request_remote.reject(Error::NotFound);
4005 settle().await;
4006
4007 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4009 settle().await;
4010 let mut sub = subscribing
4011 .await
4012 .expect("the handover must win over the stale rejection");
4013 producer_local.append_group().unwrap();
4014 assert_eq!(sub.assert_group().sequence, 0);
4015 sub.assert_not_closed();
4016 }
4017
4018 #[tokio::test]
4022 async fn test_route_handover() {
4023 tokio::time::pause();
4024
4025 let origin = Origin::random().produce();
4026 let consumer = origin.consume();
4027 let mut announced = consumer.announced();
4028
4029 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4031 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4032
4033 let source_a = origin
4034 .create_broadcast("test", announce().with_hops(hops_long))
4035 .unwrap();
4036 let mut dynamic_a = source_a.dynamic();
4037 settle().await;
4038 settle().await;
4039 let broadcast = consumer.request_broadcast("test").await.unwrap();
4040 announced.assert_next_some("test");
4041
4042 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4043 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
4044 settle().await;
4045 let mut sub = subscribing.await.unwrap();
4046 producer_a.append_group().unwrap();
4047 producer_a.append_group().unwrap();
4048 assert_eq!(sub.assert_group().sequence, 0);
4049 assert_eq!(sub.assert_group().sequence, 1);
4050
4051 let source_b = origin
4054 .create_broadcast("test", announce().with_hops(hops_short))
4055 .unwrap();
4056 let mut dynamic_b = source_b.dynamic();
4057 settle().await;
4058 settle().await;
4059 announced.assert_next_wait();
4060
4061 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
4062 settle().await;
4063
4064 sub.assert_no_group();
4067 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
4068 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
4069
4070 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
4072 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
4073 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
4074 assert_eq!(sub.assert_group().sequence, 2);
4075 assert_eq!(sub.assert_group().sequence, 3);
4076 sub.assert_no_group();
4077 sub.assert_not_closed();
4078 }
4079
4080 #[tokio::test(start_paused = true)]
4083 async fn test_route_unannounce_immediate() {
4084 let origin = Origin::random().produce();
4085 let consumer = origin.consume();
4086 let mut announced = consumer.announced();
4087
4088 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4089 let mut source = origin
4090 .create_broadcast("test", announce().with_hops(hops.clone()))
4091 .unwrap();
4092 settle().await;
4093 let broadcast = consumer.request_broadcast("test").await.unwrap();
4094 announced.assert_next_some("test");
4095
4096 source.finish();
4099 settle().await;
4100 announced.assert_next_none("test");
4101
4102 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4104 settle().await;
4105 let fresh = consumer.request_broadcast("test").await.unwrap();
4106 announced.assert_next_some("test");
4107 assert!(
4108 !fresh.is_clone(&broadcast),
4109 "re-create must not splice the old broadcast"
4110 );
4111 }
4112
4113 #[tokio::test(start_paused = true)]
4118 async fn test_route_detach_immediate() {
4119 let origin = Origin::random().produce();
4120 let consumer = origin.consume();
4121 let mut announced = consumer.announced();
4122
4123 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4124 let source = origin
4125 .create_broadcast("test", announce().with_hops(hops.clone()))
4126 .unwrap();
4127 let mut dynamic = source.dynamic();
4128 settle().await;
4129 settle().await;
4130 let broadcast = consumer.request_broadcast("test").await.unwrap();
4131 announced.assert_next_some("test");
4132
4133 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4134 let producer = accept_track(&mut dynamic, "video").await;
4135 settle().await;
4136 let mut sub = subscribing.await.unwrap();
4137
4138 drop(producer);
4140 source.abort(Error::Dropped).unwrap();
4141 drop(dynamic);
4142
4143 settle().await;
4144 announced.assert_next_none("test");
4145 sub.assert_error();
4146
4147 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4150 settle().await;
4151 settle().await;
4152 let fresh = consumer.request_broadcast("test").await.unwrap();
4153 announced.assert_next_some("test");
4154 assert!(
4155 !fresh.is_clone(&broadcast),
4156 "re-create must not splice the old broadcast"
4157 );
4158 }
4159
4160 #[tokio::test(start_paused = true)]
4165 async fn test_idle_track_releases_without_respinning() {
4166 let origin = Info::new(Origin::random()).produce();
4167 let consumer = origin.consume();
4168
4169 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4170 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4171 let mut dynamic = source.dynamic();
4172 settle().await;
4173 let broadcast = consumer.request_broadcast("test").await.unwrap();
4174
4175 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4176 let producer = accept_track(&mut dynamic, "video").await;
4177 settle().await;
4178 let sub = subscribing.await.unwrap();
4179
4180 drop(sub);
4183 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
4184 settle().await;
4185 assert!(
4186 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
4187 "the copy must stay spliced inside the linger",
4188 );
4189
4190 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4193 settle().await;
4194 assert!(
4195 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4196 "an idle copy must be released after the linger",
4197 );
4198
4199 for _ in 0..3 {
4203 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4204 settle().await;
4205 assert!(
4206 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4207 "an unread copy must stay released, not be re-spliced",
4208 );
4209 }
4210 assert!(
4211 dynamic.requested_track().now_or_never().is_none(),
4212 "an unread track must not be re-requested",
4213 );
4214 drop(producer);
4215
4216 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4218 let mut producer = accept_track(&mut dynamic, "video").await;
4219 settle().await;
4220 let mut sub = subscribing.await.unwrap();
4221 producer.append_group().unwrap();
4222 assert_eq!(sub.assert_group().sequence, 0);
4223 }
4224
4225 #[tokio::test(start_paused = true)]
4229 async fn test_back_to_back_fetches_reuse_the_track() {
4230 let origin = Info::new(Origin::random()).produce();
4231 let consumer = origin.consume();
4232
4233 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4234 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4235 let mut dynamic = source.dynamic();
4236 settle().await;
4237 let broadcast = consumer.request_broadcast("test").await.unwrap();
4238
4239 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4241 let mut producer = accept_track(&mut dynamic, "video").await;
4242 producer.append_group().unwrap().finish().unwrap();
4243 settle().await;
4244 let first = fetching.await.expect("first fetch");
4245 drop(first);
4246
4247 settle().await;
4249 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4250 settle().await;
4251 assert!(
4252 dynamic.requested_track().now_or_never().is_none(),
4253 "a fetch inside the linger must reuse the track, not re-request it",
4254 );
4255 drop(fetching.await.expect("second fetch"));
4256
4257 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4259 settle().await;
4260 assert!(
4261 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4262 "the copy must be released once the fetches stop",
4263 );
4264 drop(producer);
4265
4266 settle().await;
4268 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4269 let mut producer = accept_track(&mut dynamic, "video").await;
4270 producer.append_group().unwrap().finish().unwrap();
4271 settle().await;
4272 fetching.await.expect("fetch after the linger");
4273 }
4274
4275 #[tokio::test(start_paused = true)]
4279 async fn test_linger_reconnect_splices() {
4280 let origin = Info::new(Origin::random())
4281 .with_linger(Duration::from_secs(5))
4282 .produce();
4283 let consumer = origin.consume();
4284 let mut announced = consumer.announced();
4285
4286 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4287 let source = origin
4288 .create_broadcast("test", announce().with_hops(hops.clone()))
4289 .unwrap();
4290 let mut dynamic = source.dynamic();
4291 settle().await;
4292 settle().await;
4293 let broadcast = consumer.request_broadcast("test").await.unwrap();
4294 announced.assert_next_some("test");
4295
4296 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4297 let mut producer = accept_track(&mut dynamic, "video").await;
4298 settle().await;
4299 let mut sub = subscribing.await.unwrap();
4300
4301 producer.append_group().unwrap();
4302 producer.append_group().unwrap();
4303 assert_eq!(sub.assert_group().sequence, 0);
4304 assert_eq!(sub.assert_group().sequence, 1);
4305
4306 drop(producer);
4309 source.abort(Error::Dropped).unwrap();
4310 drop(dynamic);
4311 settle().await;
4312
4313 announced.assert_next_wait();
4315 sub.assert_no_group();
4316 sub.assert_not_closed();
4317
4318 let during = consumer.request_broadcast("test").await.unwrap();
4320 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4321
4322 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4325 let mut dynamic = source.dynamic();
4326 settle().await;
4327 settle().await;
4328 announced.assert_next_wait();
4329 let again = consumer.request_broadcast("test").await.unwrap();
4330 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4331
4332 let mut producer = accept_track(&mut dynamic, "video").await;
4336 settle().await;
4337 sub.assert_no_group();
4338 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
4339 producer.create_group(group::Info { sequence: 2 }).unwrap();
4340 assert_eq!(sub.assert_group().sequence, 2);
4341 sub.assert_not_closed();
4342 }
4343
4344 #[tokio::test(start_paused = true)]
4347 async fn test_linger_expiry_closes() {
4348 let origin = Info::new(Origin::random())
4349 .with_linger(Duration::from_secs(5))
4350 .produce();
4351 let consumer = origin.consume();
4352 let mut announced = consumer.announced();
4353
4354 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4355 let source = origin
4356 .create_broadcast("test", announce().with_hops(hops.clone()))
4357 .unwrap();
4358 let mut dynamic = source.dynamic();
4359 settle().await;
4360 settle().await;
4361 let broadcast = consumer.request_broadcast("test").await.unwrap();
4362 announced.assert_next_some("test");
4363
4364 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4365 let producer = accept_track(&mut dynamic, "video").await;
4366 settle().await;
4367 let mut sub = subscribing.await.unwrap();
4368
4369 drop(producer);
4370 source.abort(Error::Dropped).unwrap();
4371 drop(dynamic);
4372 settle().await;
4373 announced.assert_next_wait();
4374
4375 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4377 settle().await;
4378 announced.assert_next_none("test");
4379 sub.assert_error();
4380
4381 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4383 settle().await;
4384 settle().await;
4385 let fresh = consumer.request_broadcast("test").await.unwrap();
4386 announced.assert_next_some("test");
4387 assert!(
4388 !fresh.is_clone(&broadcast),
4389 "a late re-create must not splice the expired broadcast"
4390 );
4391 }
4392
4393 #[tokio::test(start_paused = true)]
4397 async fn test_linger_forever() {
4398 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4399 let consumer = origin.consume();
4400 let mut announced = consumer.announced();
4401
4402 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4403 let source = origin
4404 .create_broadcast("test", announce().with_hops(hops.clone()))
4405 .unwrap();
4406 settle().await;
4407 let broadcast = consumer.request_broadcast("test").await.unwrap();
4408 announced.assert_next_some("test");
4409
4410 source.abort(Error::Dropped).unwrap();
4411 settle().await;
4412
4413 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4415 announced.assert_next_wait();
4416 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4417 settle().await;
4418 settle().await;
4419 let again = consumer.request_broadcast("test").await.unwrap();
4420 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4421 drop(source);
4422 }
4423
4424 #[tokio::test(start_paused = true)]
4438 async fn test_linger_parks_a_live_subscription() {
4439 let origin = Info::new(Origin::random())
4440 .with_linger(Duration::from_secs(5))
4441 .produce();
4442 let consumer = origin.consume();
4443
4444 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4445 let source = origin
4446 .create_broadcast("test", announce().with_hops(hops.clone()))
4447 .unwrap();
4448 let mut dynamic = source.dynamic();
4449 settle().await;
4450 settle().await;
4451 let broadcast = consumer.request_broadcast("test").await.unwrap();
4452
4453 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4456 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4457 settle().await;
4458 let mut sub = subscribing.await.unwrap();
4459 producer.append_group().unwrap();
4460 assert_eq!(sub.assert_group().sequence, 0);
4461
4462 source.abort(Error::Dropped).unwrap();
4467 settle().await;
4468 settle().await;
4469 sub.assert_not_closed();
4470
4471 tokio::time::sleep(Duration::from_secs(4)).await;
4474 settle().await;
4475 sub.assert_not_closed();
4476
4477 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4479 let mut dynamic = source.dynamic();
4480 settle().await;
4481 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4482 settle().await;
4483 producer.create_group(group::Info { sequence: 1 }).unwrap();
4484 assert_eq!(sub.assert_group().sequence, 1);
4485 sub.assert_not_closed();
4486 }
4487
4488 #[tokio::test(start_paused = true)]
4499 async fn test_idle_release_survives_the_route_leaving() {
4500 let origin = Info::new(Origin::random())
4501 .with_linger(Duration::from_secs(600))
4502 .produce();
4503 let consumer = origin.consume();
4504
4505 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4506 let source = origin
4507 .create_broadcast("test", announce().with_hops(hops.clone()))
4508 .unwrap();
4509 let mut dynamic = source.dynamic();
4510 settle().await;
4511 settle().await;
4512 let broadcast = consumer.request_broadcast("test").await.unwrap();
4513
4514 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4515 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4516 settle().await;
4517 let mut sub = subscribing.await.unwrap();
4518 producer.append_group().unwrap();
4519 assert_eq!(sub.assert_group().sequence, 0);
4520
4521 source.abort(Error::Dropped).unwrap();
4524 settle().await;
4525 settle().await;
4526 drop(sub);
4527 drop(producer);
4528 drop(dynamic);
4529 settle().await;
4530 tokio::time::sleep(TRACK_IDLE_LINGER + Duration::from_secs(1)).await;
4531 settle().await;
4532
4533 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4536 let mut dynamic = source.dynamic();
4537 settle().await;
4538 settle().await;
4539 let broadcast = consumer.request_broadcast("test").await.unwrap();
4540 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4541 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4542 settle().await;
4543 let mut sub = subscribing.await.unwrap();
4544 producer.append_group().unwrap();
4545 assert_eq!(
4546 sub.assert_group().sequence,
4547 0,
4548 "the reconnect's first group must not be filtered by a stale boundary"
4549 );
4550 }
4551
4552 #[tokio::test(start_paused = true)]
4555 async fn test_linger_skipped_on_finish() {
4556 let origin = Info::new(Origin::random())
4557 .with_linger(Duration::from_secs(5))
4558 .produce();
4559 let consumer = origin.consume();
4560 let mut announced = consumer.announced();
4561
4562 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4563 let mut source = origin
4564 .create_broadcast("test", announce().with_hops(hops.clone()))
4565 .unwrap();
4566 settle().await;
4567 let broadcast = consumer.request_broadcast("test").await.unwrap();
4568 announced.assert_next_some("test");
4569
4570 source.finish();
4573 settle().await;
4574 announced.assert_next_none("test");
4575
4576 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4578 settle().await;
4579 let fresh = consumer.request_broadcast("test").await.unwrap();
4580 announced.assert_next_some("test");
4581 assert!(
4582 !fresh.is_clone(&broadcast),
4583 "a finish must not leave a lingering broadcast to splice into"
4584 );
4585 }
4586
4587 #[tokio::test]
4590 async fn test_announce_toggle() {
4591 tokio::time::pause();
4592
4593 let origin = Origin::random().produce();
4594 let consumer = origin.consume();
4595 let mut announced = consumer.announced();
4596
4597 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4598 settle().await;
4599
4600 announced.assert_next_wait();
4602 let broadcast = consumer
4603 .get_broadcast("test")
4604 .expect("offline broadcast is still routable");
4605 assert!(!broadcast.route().announce);
4606
4607 let requested = consumer.request_broadcast("test").await.unwrap();
4609 assert!(requested.is_clone(&broadcast));
4610
4611 source.set_route(announce()).unwrap();
4613 settle().await;
4614 let face = announced.assert_next_some("test");
4615 assert!(face.is_clone(&broadcast));
4616
4617 let mut fresh = origin.consume().announced();
4619 fresh.assert_next_some("test");
4620 fresh.assert_next_wait();
4621
4622 source.set_route(broadcast::Route::new()).unwrap();
4624 settle().await;
4625 announced.assert_next_none("test");
4626 assert!(consumer.get_broadcast("test").is_some());
4627 let mut fresh = origin.consume().announced();
4628 fresh.assert_next_wait();
4629
4630 source.finish();
4631 settle().await;
4632 assert!(consumer.get_broadcast("test").is_none());
4633 }
4634
4635 #[tokio::test]
4638 async fn test_announce_beats_offline() {
4639 tokio::time::pause();
4640
4641 let origin = Origin::random().produce();
4642 let consumer = origin.consume();
4643 let mut announced = consumer.announced();
4644
4645 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4647 settle().await;
4648 announced.assert_next_wait();
4649
4650 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
4653 settle().await;
4654 announced.assert_next_some("test");
4655 let face = consumer.get_broadcast("test").unwrap();
4656 assert!(face.route().announce);
4657 assert_eq!(face.route().cost, 10);
4658
4659 announced_source.finish();
4662 settle().await;
4663 announced.assert_next_none("test");
4664 assert!(consumer.get_broadcast("test").is_some());
4665 }
4666
4667 #[tokio::test]
4670 async fn test_better_source_no_churn() {
4671 tokio::time::pause();
4672
4673 let origin = Origin::random().produce();
4674 let mut announced = origin.consume().announced();
4675
4676 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4679 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4680 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4681 settle().await;
4682 let face = announced.assert_next_some("test");
4683
4684 let _b = origin
4685 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4686 .unwrap();
4687 settle().await;
4688 announced.assert_next_wait();
4689 let current = origin.consume().get_broadcast("test").unwrap();
4690 assert!(current.is_clone(&face), "the broadcast identity must not change");
4691 assert_eq!(current.route().hops, hops_b);
4693 }
4694
4695 #[tokio::test]
4701 async fn test_publisher_mismatch_replaces() {
4702 tokio::time::pause();
4703
4704 let origin = Origin::random().produce();
4705 let consumer = origin.consume();
4706 let mut announced = consumer.announced();
4707
4708 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4709 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4710
4711 let mut source_a = origin
4712 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4713 .unwrap();
4714 settle().await;
4715 let face_a = announced.assert_next_some("test");
4716
4717 let _source_b = origin
4720 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4721 .unwrap();
4722 settle().await;
4723 settle().await;
4724 announced.assert_next_none("test");
4725 let face_b = announced.assert_next_some("test");
4726 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
4727 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4728 assert!(face_a.is_closed(), "the displaced front must close");
4732
4733 source_a.finish();
4735 settle().await;
4736 settle().await;
4737 announced.assert_next_wait();
4738 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4739 }
4740
4741 #[tokio::test]
4746 async fn test_reconnect_wins_over_stale_route() {
4747 tokio::time::pause();
4748
4749 let origin = Origin::random().produce();
4750 let consumer = origin.consume();
4751
4752 let publisher = Origin::new(1).unwrap();
4753 let hops = OriginList::try_from(vec![publisher]).unwrap();
4754
4755 let stale = origin
4758 .create_broadcast("test", announce().with_hops(hops.clone()))
4759 .unwrap();
4760 let mut stale_dynamic = stale.dynamic();
4761 settle().await;
4762
4763 let fresh = origin
4765 .create_broadcast("test", announce().with_hops(hops.clone()))
4766 .unwrap();
4767 let mut fresh_dynamic = fresh.dynamic();
4768 settle().await;
4769 settle().await;
4770
4771 let broadcast = consumer.request_broadcast("test").await.unwrap();
4773 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4774 settle().await;
4775 let _producer = accept_track(&mut fresh_dynamic, "video").await;
4776 settle().await;
4777 subscribing.await.unwrap();
4778 stale_dynamic.assert_no_request();
4779 }
4780
4781 #[tokio::test]
4787 async fn test_carrying_reconnect_switches_immediately() {
4788 tokio::time::pause();
4789
4790 let origin = Origin::random().produce();
4791 let consumer = origin.consume();
4792
4793 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4794
4795 let stale = origin
4796 .create_broadcast("test", announce().with_hops(hops.clone()))
4797 .unwrap();
4798 let mut stale_dynamic = stale.dynamic();
4799 settle().await;
4800
4801 let broadcast = consumer.request_broadcast("test").await.unwrap();
4803 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4804 settle().await;
4805 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
4806 settle().await;
4807 let _subscription = subscribing.await.unwrap();
4809
4810 let fresh = origin
4812 .create_broadcast("test", announce().with_hops(hops.clone()))
4813 .unwrap();
4814 let mut fresh_dynamic = fresh.dynamic();
4815 settle().await;
4816 settle().await;
4817
4818 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
4821 }
4822
4823 #[tokio::test]
4833 async fn test_offline_mismatch_never_evicts_a_live_front() {
4834 tokio::time::pause();
4835
4836 let origin = Origin::random().produce();
4837 let consumer = origin.consume();
4838 let mut announced = consumer.announced();
4839
4840 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4841 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4842
4843 let mut live = origin
4844 .create_broadcast("test", announce().with_hops(hops_live.clone()))
4845 .unwrap();
4846 let mut live_dynamic = live.dynamic();
4847 settle().await;
4848 let face = announced.assert_next_some("test");
4849
4850 let broadcast = consumer.request_broadcast("test").await.unwrap();
4852 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4853 settle().await;
4854 let _producer = accept_track(&mut live_dynamic, "video").await;
4855 settle().await;
4856 subscribing.await.unwrap();
4857
4858 let cache = origin
4861 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
4862 .unwrap();
4863 settle().await;
4864 settle().await;
4865 announced.assert_next_wait();
4866 assert!(!face.is_closed(), "the live front must survive");
4867 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
4868
4869 live.finish();
4872 settle().await;
4873 settle().await;
4874 announced.assert_next_none("test");
4875 let taken = consumer
4876 .get_broadcast("test")
4877 .expect("the parked source must take over");
4878 assert_eq!(taken.route().hops, hops_cache);
4879 announced.assert_next_wait();
4881 drop(cache);
4882 }
4883
4884 #[tokio::test]
4890 async fn test_dispatch_excludes_requester() {
4891 tokio::time::pause();
4892
4893 let origin = Origin::random().produce();
4894 let consumer = origin.consume();
4895
4896 let peer = Origin::new(5).unwrap();
4897 let publisher = Origin::new(1).unwrap();
4898 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
4900 let clean = OriginList::try_from(vec![publisher]).unwrap();
4901
4902 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4903 let mut dynamic_a = source_a.dynamic();
4904 settle().await;
4905 let source_b = origin
4906 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
4907 .unwrap();
4908 let mut dynamic_b = source_b.dynamic();
4909 settle().await;
4910 settle().await;
4911
4912 let shared = consumer.request_broadcast("test").await.unwrap();
4915 let subscribing = shared.track("video").unwrap().subscribe(None);
4916 let _producer_a = accept_track(&mut dynamic_a, "video").await;
4917 settle().await;
4918 subscribing.await.unwrap();
4919
4920 let scoped = consumer.clone().excluding(peer);
4925 let pinned = scoped.request_broadcast("test").await.unwrap();
4926 let subscribing = pinned.track("video").unwrap().subscribe(None);
4927 let _producer_b = accept_track(&mut dynamic_b, "video").await;
4928 settle().await;
4929 subscribing.await.unwrap();
4930 dynamic_a.assert_no_request();
4931 }
4932
4933 #[tokio::test]
4938 async fn test_unknown_publishers_do_not_splice() {
4939 tokio::time::pause();
4940
4941 let origin = Origin::random().produce();
4942 let consumer = origin.consume();
4943 let mut announced = consumer.announced();
4944
4945 let unknown_a = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
4946 let unknown_b = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
4947
4948 let source_a = origin
4949 .create_broadcast("test", announce().with_hops(unknown_a))
4950 .unwrap();
4951 settle().await;
4952 settle().await;
4953 announced.assert_next_some("test");
4954
4955 let source_b = origin
4959 .create_broadcast("test", announce().with_hops(unknown_b))
4960 .unwrap();
4961 settle().await;
4962 settle().await;
4963 announced.assert_next_none("test");
4964 announced.assert_next_some("test");
4965
4966 drop(source_a);
4967 drop(source_b);
4968 }
4969
4970 #[tokio::test]
4973 async fn test_known_publishers_still_splice() {
4974 tokio::time::pause();
4975
4976 let origin = Origin::random().produce();
4977 let consumer = origin.consume();
4978 let mut announced = consumer.announced();
4979
4980 let publisher = Origin::new(1).unwrap();
4981 let hops_a = OriginList::try_from(vec![publisher]).unwrap();
4982 let hops_b = OriginList::try_from(vec![publisher, Origin::new(3).unwrap()]).unwrap();
4983
4984 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4985 settle().await;
4986 settle().await;
4987 announced.assert_next_some("test");
4988
4989 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4991 settle().await;
4992 settle().await;
4993 announced.assert_next_wait();
4994
4995 drop(source_a);
4996 drop(source_b);
4997 }
4998
4999 #[tokio::test]
5005 async fn test_standby_join_splices_live_subscriber() {
5006 tokio::time::pause();
5007
5008 let origin = Origin::random().produce();
5009 let consumer = origin.consume();
5010
5011 let publisher = Origin::new(1).unwrap();
5012 let peer = Origin::new(5).unwrap();
5013 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5014 let local = OriginList::try_from(vec![publisher]).unwrap();
5015
5016 let source_remote = origin
5018 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5019 .unwrap();
5020 let mut dynamic_remote = source_remote.dynamic();
5021 settle().await;
5022 settle().await;
5023 let broadcast = consumer.request_broadcast("test").await.unwrap();
5024 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5025 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
5026 settle().await;
5027 let mut sub = subscribing.await.unwrap();
5028 producer_remote.append_group().unwrap();
5029 assert_eq!(sub.assert_group().sequence, 0);
5030
5031 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5034 let mut dynamic_local = source_local.dynamic();
5035 settle().await;
5036 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
5037 settle().await;
5038 sub.assert_no_group();
5039 assert_eq!(producer_local.subscription().unwrap().group_start, Some(1));
5040 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
5041 assert_eq!(sub.assert_group().sequence, 1);
5042 sub.assert_not_closed();
5043 }
5044
5045 #[tokio::test]
5052 async fn test_standby_missing_track_keeps_incumbent() {
5053 tokio::time::pause();
5054
5055 let origin = Origin::random().produce();
5056 let consumer = origin.consume();
5057
5058 let publisher = Origin::new(1).unwrap();
5059 let peer = Origin::new(5).unwrap();
5060 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5061 let local = OriginList::try_from(vec![publisher]).unwrap();
5062
5063 let source_remote = origin
5065 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5066 .unwrap();
5067 let mut dynamic_remote = source_remote.dynamic();
5068 settle().await;
5069 settle().await;
5070 let broadcast = consumer.request_broadcast("test").await.unwrap();
5071 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5072 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
5073 settle().await;
5074 let mut sub = subscribing.await.unwrap();
5075 producer_remote.append_group().unwrap();
5076 assert_eq!(sub.assert_group().sequence, 0);
5077
5078 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5081 let mut dynamic_local = source_local.dynamic();
5082 settle().await;
5083 let request = dynamic_local.requested_track().await.unwrap();
5084 assert_eq!(request.name(), "audio");
5085 request.reject(Error::NotFound);
5086 settle().await;
5087
5088 producer_remote.append_group().unwrap();
5090 assert_eq!(sub.assert_group().sequence, 1);
5091 sub.assert_not_closed();
5092
5093 source_remote.abort(Error::Dropped).unwrap();
5096 settle().await;
5097 settle().await;
5098 sub.assert_closed();
5099 dynamic_local.assert_no_request();
5100
5101 let retry = broadcast.track("audio").unwrap().subscribe(None);
5103 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
5104 settle().await;
5105 let mut sub = retry.await.expect("a fresh request must reach the standby");
5106 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
5107 assert_eq!(sub.assert_group().sequence, 2);
5108 }
5109
5110 #[tokio::test]
5115 async fn test_unservable_track_retried_by_a_later_request() {
5116 tokio::time::pause();
5117
5118 let origin = Origin::random().produce();
5119 let consumer = origin.consume();
5120
5121 let source = origin.create_broadcast("test", announce()).unwrap();
5122 let mut dynamic = source.dynamic();
5123 settle().await;
5124 settle().await;
5125 let broadcast = consumer.request_broadcast("test").await.unwrap();
5126
5127 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5129 let request = dynamic.requested_track().await.unwrap();
5130 request.reject(Error::NotFound);
5131 settle().await;
5132 assert!(matches!(subscribing.await, Err(Error::NotFound)));
5133
5134 let retry = broadcast.track("audio").unwrap().subscribe(None);
5136 let mut producer = accept_track(&mut dynamic, "audio").await;
5137 settle().await;
5138 let mut sub = retry.await.expect("a fresh request must reach the source");
5139 producer.append_group().unwrap();
5140 assert_eq!(sub.assert_group().sequence, 0);
5141 }
5142
5143 #[tokio::test]
5149 async fn test_track_dying_without_progress_aborts() {
5150 tokio::time::pause();
5151
5152 let origin = Origin::random().produce();
5153 let consumer = origin.consume();
5154
5155 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5156 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5157 let mut dynamic = source.dynamic();
5158 settle().await;
5159 settle().await;
5160 let broadcast = consumer.request_broadcast("test").await.unwrap();
5161
5162 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5163 let producer = accept_track(&mut dynamic, "video").await;
5164 settle().await;
5165 let mut sub = subscribing.await.unwrap();
5166
5167 drop(producer);
5170 settle().await;
5171 sub.assert_closed();
5172 dynamic.assert_no_request();
5173
5174 let retry = broadcast.track("video").unwrap().subscribe(None);
5176 let mut producer = accept_track(&mut dynamic, "video").await;
5177 settle().await;
5178 let mut sub = retry.await.expect("a fresh request must reach the source");
5179 producer.append_group().unwrap();
5180 assert_eq!(sub.assert_group().sequence, 0);
5181 }
5182
5183 #[tokio::test]
5188 async fn test_delivered_copy_death_survives_unrelated_wakes() {
5189 tokio::time::pause();
5190
5191 let origin = Origin::random().produce();
5192 let consumer = origin.consume();
5193
5194 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5195 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5196 let mut dynamic = source.dynamic();
5197 settle().await;
5198 settle().await;
5199 let broadcast = consumer.request_broadcast("test").await.unwrap();
5200
5201 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5202 let mut producer = accept_track(&mut dynamic, "video").await;
5203 settle().await;
5204 let mut sub = subscribing.await.unwrap();
5205 producer.append_group().unwrap();
5206 assert_eq!(sub.assert_group().sequence, 0);
5207
5208 drop(sub);
5211 settle().await;
5212 let resubscribing = broadcast.track("video").unwrap().subscribe(None);
5213 settle().await;
5214 let mut sub = resubscribing.await.unwrap();
5215 assert_eq!(sub.assert_group().sequence, 0, "cached group re-served");
5216
5217 drop(producer);
5220 let mut producer = accept_track(&mut dynamic, "video").await;
5221 settle().await;
5222 producer.create_group(group::Info { sequence: 1 }).unwrap();
5223 assert_eq!(sub.assert_group().sequence, 1);
5224 sub.assert_not_closed();
5225 }
5226
5227 #[tokio::test]
5233 async fn test_per_track_fallback_respects_exclusion() {
5234 tokio::time::pause();
5235
5236 let origin = Origin::random().produce();
5237 let consumer = origin.consume();
5238
5239 let publisher = Origin::new(1).unwrap();
5240 let peer = Origin::new(5).unwrap();
5241 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5242 let local = OriginList::try_from(vec![publisher]).unwrap();
5243
5244 let source_tainted = origin
5246 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5247 .unwrap();
5248 let mut dynamic_tainted = source_tainted.dynamic();
5249 settle().await;
5250 settle().await;
5251
5252 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5255 let mut dynamic_clean = source_clean.dynamic();
5256 settle().await;
5257
5258 let scoped = consumer.clone().excluding(peer);
5259 let broadcast = scoped.request_broadcast("test").await.unwrap();
5260 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5261 settle().await;
5262
5263 let request = dynamic_clean.requested_track().await.unwrap();
5267 request.reject(Error::NotFound);
5268 settle().await;
5269 dynamic_tainted.assert_no_request();
5270 }
5271
5272 #[tokio::test]
5276 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
5277 tokio::time::pause();
5278
5279 let origin = Origin::random().produce();
5280 let consumer = origin.consume();
5281
5282 let publisher = Origin::new(1).unwrap();
5283 let peer = Origin::new(5).unwrap();
5284 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5285 let local = OriginList::try_from(vec![publisher]).unwrap();
5286
5287 let source_tainted = origin
5288 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5289 .unwrap();
5290 let mut dynamic_tainted = source_tainted.dynamic();
5291 settle().await;
5292 settle().await;
5293 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5294 let mut dynamic_clean = source_clean.dynamic();
5295 settle().await;
5296
5297 let scoped = consumer.clone().excluding(peer);
5298 let broadcast = scoped.request_broadcast("test").await.unwrap();
5299 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5300 let _clean = accept_track(&mut dynamic_clean, "video").await;
5301 settle().await;
5302
5303 source_clean.abort(Error::Dropped).unwrap();
5305 settle().await;
5306 settle().await;
5307 dynamic_tainted.assert_no_request();
5308
5309 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5312 }
5313
5314 #[tokio::test]
5319 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
5320 tokio::time::pause();
5321
5322 let origin = Origin::random().produce();
5323 let consumer = origin.consume();
5324
5325 let publisher = Origin::new(1).unwrap();
5326 let peer = Origin::new(5).unwrap();
5327 let local = OriginList::try_from(vec![publisher]).unwrap();
5328 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5329
5330 let source_clean = origin
5334 .create_broadcast("test", announce().with_hops(local).with_cost(5))
5335 .unwrap();
5336 let mut dynamic_clean = source_clean.dynamic();
5337 settle().await;
5338 settle().await;
5339 let scoped = consumer.clone().excluding(peer);
5340 let broadcast = scoped.request_broadcast("test").await.unwrap();
5341 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5342 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
5343 settle().await;
5344 let mut sub = subscribing.await.unwrap();
5345 producer_clean.append_group().unwrap();
5346 assert_eq!(sub.assert_group().sequence, 0);
5347
5348 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
5355 tainted.advertised = 1;
5356 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
5357 let mut dynamic_tainted = source_tainted.dynamic();
5358 settle().await;
5359 settle().await;
5360 dynamic_tainted.assert_no_request();
5361 producer_clean.append_group().unwrap();
5362 assert_eq!(sub.assert_group().sequence, 1);
5363 sub.assert_not_closed();
5364
5365 drop(sub);
5368 drop(broadcast);
5369 drop(scoped);
5370 settle().await;
5371 let mut bumped = announce().with_hops(via_peer).with_cost(1);
5372 bumped.advertised = 1;
5373 source_tainted.set_route(bumped).unwrap();
5374 settle().await;
5375 let plain = consumer.request_broadcast("test").await.unwrap();
5376 let _plain_track = plain.track("video").unwrap().subscribe(None);
5377 settle().await;
5378 settle().await;
5379 assert!(
5380 dynamic_tainted.requested_track().now_or_never().is_some(),
5381 "the front must be free to use the route again once the peer is gone"
5382 );
5383 }
5384
5385 #[tokio::test]
5389 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
5390 tokio::time::pause();
5391
5392 let origin = Origin::random().produce();
5393 let consumer = origin.consume();
5394 let mut dynamic = origin.dynamic();
5395
5396 let peer = Origin::new(5).unwrap();
5397 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5398 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5399 settle().await;
5400 settle().await;
5401
5402 let scoped = consumer.clone().excluding(peer);
5403 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5404 assert!(
5405 dynamic.requested_broadcast().now_or_never().is_none(),
5406 "the dynamic handler was asked to route around the exclusion"
5407 );
5408
5409 let _pending = scoped.request_broadcast("other");
5411 settle().await;
5412 assert!(
5413 dynamic.requested_broadcast().now_or_never().is_some(),
5414 "a genuinely missing path must still fall back"
5415 );
5416 }
5417
5418 #[tokio::test]
5421 async fn test_dispatch_all_tainted_unroutable() {
5422 tokio::time::pause();
5423
5424 let origin = Origin::random().produce();
5425 let consumer = origin.consume();
5426
5427 let peer = Origin::new(5).unwrap();
5428 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5429 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5430 settle().await;
5431 settle().await;
5432
5433 let scoped = consumer.clone().excluding(peer);
5434 match scoped.request_broadcast("test").await {
5435 Err(Error::Unroutable) => {}
5436 Err(err) => panic!("expected Unroutable, got {err:?}"),
5437 Ok(_) => panic!("expected Unroutable, got a broadcast"),
5438 }
5439
5440 consumer.request_broadcast("test").await.unwrap();
5442 }
5443
5444 #[tokio::test]
5445 async fn test_duplicate_reverse() {
5446 tokio::time::pause();
5447
5448 let origin = Origin::random().produce();
5449
5450 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5451 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5452 settle().await;
5453 assert!(origin.consume().get_broadcast("test").is_some());
5454
5455 broadcast2.finish();
5457 settle().await;
5458 assert!(origin.consume().get_broadcast("test").is_some());
5459
5460 broadcast1.finish();
5461 settle().await;
5462 assert!(origin.consume().get_broadcast("test").is_none());
5463 }
5464
5465 #[tokio::test]
5466 async fn test_deterministic_tiebreak() {
5467 tokio::time::pause();
5468
5469 fn hops(ids: &[u64]) -> OriginList {
5470 OriginList::try_from(
5471 ids.iter()
5472 .copied()
5473 .map(|id| Origin::new(id).unwrap())
5474 .collect::<Vec<_>>(),
5475 )
5476 .unwrap()
5477 }
5478
5479 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
5482 let origin = Origin::random().produce();
5483 let _a = origin
5484 .create_broadcast("test", announce().with_hops(hops(first)))
5485 .unwrap();
5486 let _b = origin
5487 .create_broadcast("test", announce().with_hops(hops(second)))
5488 .unwrap();
5489 settle().await;
5490 origin.consume().get_broadcast("test").unwrap().route().hops
5491 }
5492
5493 let forward = winner(&[5, 20], &[5, 40]).await;
5497 let reverse = winner(&[5, 40], &[5, 20]).await;
5498 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
5499
5500 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
5502 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
5503 }
5504
5505 #[tokio::test]
5510 async fn test_many_announces() {
5511 let origin = Origin::random().produce();
5512
5513 let mut consumer = origin.consume().announced();
5514 let mut broadcasts = Vec::new();
5516 for i in 0..256 {
5517 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5518 settle().await;
5519 }
5520
5521 for i in 0..256 {
5522 consumer.assert_next_some(format!("test{i:03}"));
5523 }
5524 consumer.assert_next_wait();
5525 }
5526
5527 #[tokio::test]
5528 async fn test_many_announces_try() {
5529 let origin = Origin::random().produce();
5530
5531 let mut consumer = origin.consume().announced();
5532 let mut broadcasts = Vec::new();
5534 for i in 0..256 {
5535 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5536 settle().await;
5537 }
5538
5539 for i in 0..256 {
5540 consumer.assert_try_next_some(format!("test{i:03}"));
5541 }
5542 }
5543
5544 #[tokio::test]
5545 async fn test_with_root_basic() {
5546 let origin = Origin::random().produce();
5547
5548 let foo_producer = origin.with_root("foo").expect("should create root");
5550 assert_eq!(foo_producer.root().as_str(), "foo");
5551
5552 let mut consumer = origin.consume().announced();
5553
5554 let _broadcast = foo_producer
5556 .create_broadcast("bar/baz", announce())
5557 .expect("publish allowed");
5558 settle().await;
5559 consumer.assert_next_some("foo/bar/baz");
5561
5562 let mut foo_consumer = foo_producer.consume().announced();
5564 foo_consumer.assert_next_some("bar/baz");
5565 }
5566
5567 #[tokio::test]
5568 async fn test_with_root_nested() {
5569 let origin = Origin::random().produce();
5570
5571 let foo_producer = origin.with_root("foo").expect("should create foo root");
5573 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
5574 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
5575
5576 let mut consumer = origin.consume().announced();
5577
5578 let _broadcast = foo_bar_producer
5580 .create_broadcast("baz", announce())
5581 .expect("publish allowed");
5582 settle().await;
5583 consumer.assert_next_some("foo/bar/baz");
5585
5586 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
5588 foo_bar_consumer.assert_next_some("baz");
5589 }
5590
5591 #[tokio::test]
5592 async fn test_publish_scope_allows() {
5593 let origin = Origin::random().produce();
5594
5595 let limited_producer = origin
5597 .scope(&["allowed/path1".into(), "allowed/path2".into()])
5598 .expect("should create limited producer");
5599
5600 let _broadcast = limited_producer
5602 .create_broadcast("allowed/path1", announce())
5603 .expect("publish allowed");
5604 let _keep2 = limited_producer
5605 .create_broadcast("allowed/path1/nested", announce())
5606 .expect("publish allowed");
5607 let _keep3 = limited_producer
5608 .create_broadcast("allowed/path2", announce())
5609 .expect("publish allowed");
5610 settle().await;
5611
5612 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
5614 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
5616 }
5617
5618 #[tokio::test]
5619 async fn test_publish_max_parts() {
5620 let origin = Origin::random().produce();
5621
5622 let at_limit = (0..Path::MAX_PARTS)
5623 .map(|i| i.to_string())
5624 .collect::<Vec<_>>()
5625 .join("/");
5626 let _broadcast = origin
5627 .create_broadcast(at_limit.as_str(), announce())
5628 .expect("publish allowed");
5629 settle().await;
5630
5631 let too_deep = format!("{at_limit}/extra");
5632 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
5633
5634 let rooted = origin.with_root("root").expect("wildcard allows any root");
5636 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
5637 }
5638
5639 #[tokio::test]
5640 async fn test_publish_scope_empty() {
5641 let origin = Origin::random().produce();
5642
5643 assert!(origin.scope(&[]).is_none());
5645 }
5646
5647 #[tokio::test]
5648 async fn test_consume_scope_filters() {
5649 let origin = Origin::random().produce();
5650
5651 let mut consumer = origin.consume().announced();
5652
5653 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
5655 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
5656 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
5657 settle().await;
5658
5659 let mut limited_consumer = origin
5661 .consume()
5662 .scope(&["allowed".into()])
5663 .expect("should create limited consumer")
5664 .announced();
5665
5666 limited_consumer.assert_next_some("allowed");
5668 limited_consumer.assert_next_some("allowed/nested");
5669 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
5673 consumer.assert_next_some("allowed/nested");
5674 consumer.assert_next_some("notallowed");
5675 }
5676
5677 #[tokio::test]
5678 async fn test_consume_scope_multiple_prefixes() {
5679 let origin = Origin::random().produce();
5680
5681 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5682 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5683 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5684 settle().await;
5685
5686 let mut limited_consumer = origin
5688 .consume()
5689 .scope(&["foo".into(), "bar".into()])
5690 .expect("should create limited consumer")
5691 .announced();
5692
5693 limited_consumer.assert_next_some("bar/test");
5695 limited_consumer.assert_next_some("foo/test");
5696 limited_consumer.assert_next_wait(); }
5698
5699 #[tokio::test]
5700 async fn test_with_root_and_publish_scope() {
5701 let origin = Origin::random().produce();
5702
5703 let foo_producer = origin.with_root("foo").expect("should create foo root");
5705
5706 let limited_producer = foo_producer
5708 .scope(&["bar".into(), "goop/pee".into()])
5709 .expect("should create limited producer");
5710
5711 let mut consumer = origin.consume().announced();
5712
5713 let _broadcast = limited_producer
5715 .create_broadcast("bar", announce())
5716 .expect("publish allowed");
5717 let _keep2 = limited_producer
5718 .create_broadcast("bar/nested", announce())
5719 .expect("publish allowed");
5720 let _keep3 = limited_producer
5721 .create_broadcast("goop/pee", announce())
5722 .expect("publish allowed");
5723 let _keep4 = limited_producer
5724 .create_broadcast("goop/pee/nested", announce())
5725 .expect("publish allowed");
5726 settle().await;
5727
5728 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
5730 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
5732
5733 consumer.assert_next_some("foo/bar");
5735 consumer.assert_next_some("foo/bar/nested");
5736 consumer.assert_next_some("foo/goop/pee");
5737 consumer.assert_next_some("foo/goop/pee/nested");
5738 }
5739
5740 #[tokio::test]
5741 async fn test_with_root_and_consume_scope() {
5742 let origin = Origin::random().produce();
5743
5744 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
5746 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
5747 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
5748 settle().await;
5749
5750 let foo_producer = origin.with_root("foo").expect("should create foo root");
5752
5753 let mut limited_consumer = foo_producer
5755 .consume()
5756 .scope(&["bar".into(), "goop/pee".into()])
5757 .expect("should create limited consumer")
5758 .announced();
5759
5760 limited_consumer.assert_next_some("bar/test");
5762 limited_consumer.assert_next_some("goop/pee/test");
5763 limited_consumer.assert_next_wait(); }
5765
5766 #[tokio::test]
5767 async fn test_with_root_unauthorized() {
5768 let origin = Origin::random().produce();
5769
5770 let limited_producer = origin
5772 .scope(&["allowed".into()])
5773 .expect("should create limited producer");
5774
5775 assert!(limited_producer.with_root("notallowed").is_none());
5777
5778 let allowed_root = limited_producer
5780 .with_root("allowed")
5781 .expect("should create allowed root");
5782 assert_eq!(allowed_root.root().as_str(), "allowed");
5783 }
5784
5785 #[tokio::test]
5786 async fn test_wildcard_permission() {
5787 let origin = Origin::random().produce();
5788
5789 let root_producer = origin.clone();
5791
5792 let _broadcast = root_producer
5794 .create_broadcast("any/path", announce())
5795 .expect("publish allowed");
5796 let _keep2 = root_producer
5797 .create_broadcast("other/path", announce())
5798 .expect("publish allowed");
5799 settle().await;
5800
5801 let foo_producer = root_producer.with_root("foo").expect("should create any root");
5803 assert_eq!(foo_producer.root().as_str(), "foo");
5804 }
5805
5806 #[tokio::test]
5807 async fn test_consume_broadcast_with_permissions() {
5808 let origin = Origin::random().produce();
5809
5810 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
5811 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
5812 settle().await;
5813
5814 let limited_consumer = origin
5816 .consume()
5817 .scope(&["allowed".into()])
5818 .expect("should create limited consumer");
5819
5820 let result = limited_consumer.get_broadcast("allowed/test");
5822 assert!(result.is_some());
5823 assert!(
5824 result
5825 .unwrap()
5826 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
5827 );
5828
5829 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
5831
5832 let consumer = origin.consume();
5834 assert!(consumer.get_broadcast("allowed/test").is_some());
5835 assert!(consumer.get_broadcast("notallowed/test").is_some());
5836 }
5837
5838 #[tokio::test]
5839 async fn test_nested_paths_with_permissions() {
5840 let origin = Origin::random().produce();
5841
5842 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
5844
5845 let _broadcast = limited_producer
5847 .create_broadcast("a/b/c", announce())
5848 .expect("publish allowed");
5849 let _keep2 = limited_producer
5850 .create_broadcast("a/b/c/d", announce())
5851 .expect("publish allowed");
5852 let _keep3 = limited_producer
5853 .create_broadcast("a/b/c/d/e", announce())
5854 .expect("publish allowed");
5855 settle().await;
5856
5857 assert!(limited_producer.create_broadcast("a", announce()).is_err());
5859 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
5860 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
5861 }
5862
5863 #[tokio::test]
5864 async fn test_multiple_consumers_with_different_permissions() {
5865 let origin = Origin::random().produce();
5866
5867 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5869 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5870 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5871 settle().await;
5872
5873 let mut foo_consumer = origin
5875 .consume()
5876 .scope(&["foo".into()])
5877 .expect("should create foo consumer")
5878 .announced();
5879
5880 let mut bar_consumer = origin
5881 .consume()
5882 .scope(&["bar".into()])
5883 .expect("should create bar consumer")
5884 .announced();
5885
5886 let mut foobar_consumer = origin
5887 .consume()
5888 .scope(&["foo".into(), "bar".into()])
5889 .expect("should create foobar consumer")
5890 .announced();
5891
5892 foo_consumer.assert_next_some("foo/test");
5894 foo_consumer.assert_next_wait();
5895
5896 bar_consumer.assert_next_some("bar/test");
5897 bar_consumer.assert_next_wait();
5898
5899 foobar_consumer.assert_next_some("bar/test");
5900 foobar_consumer.assert_next_some("foo/test");
5901 foobar_consumer.assert_next_wait();
5902 }
5903
5904 #[tokio::test]
5905 async fn test_select_with_empty_prefix() {
5906 let origin = Origin::random().produce();
5907
5908 let demo_producer = origin.with_root("demo").expect("should create demo root");
5910 let limited_producer = demo_producer
5911 .scope(&["worm-node".into(), "foobar".into()])
5912 .expect("should create limited producer");
5913
5914 let _broadcast1 = limited_producer
5916 .create_broadcast("worm-node/test", announce())
5917 .expect("publish allowed");
5918 let _broadcast2 = limited_producer
5919 .create_broadcast("foobar/test", announce())
5920 .expect("publish allowed");
5921 settle().await;
5922
5923 let mut consumer = limited_producer
5925 .consume()
5926 .scope(&["".into()])
5927 .expect("should create consumer with empty prefix")
5928 .announced();
5929
5930 let a1 = consumer.try_next().expect("expected first announcement");
5932 let a2 = consumer.try_next().expect("expected second announcement");
5933 consumer.assert_next_wait();
5934
5935 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5936 paths.sort();
5937 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
5938 }
5939
5940 #[tokio::test]
5941 async fn test_select_narrowing_scope() {
5942 let origin = Origin::random().produce();
5943
5944 let demo_producer = origin.with_root("demo").expect("should create demo root");
5946 let limited_producer = demo_producer
5947 .scope(&["worm-node".into(), "foobar".into()])
5948 .expect("should create limited producer");
5949
5950 let _broadcast1 = limited_producer
5952 .create_broadcast("worm-node", announce())
5953 .expect("publish allowed");
5954 let _broadcast2 = limited_producer
5955 .create_broadcast("worm-node/foo", announce())
5956 .expect("publish allowed");
5957 let _broadcast3 = limited_producer
5958 .create_broadcast("foobar/bar", announce())
5959 .expect("publish allowed");
5960 settle().await;
5961
5962 let mut worm_consumer = limited_producer
5964 .consume()
5965 .scope(&["worm-node".into()])
5966 .expect("should create worm-node consumer")
5967 .announced();
5968
5969 worm_consumer.assert_next_some("worm-node");
5971 worm_consumer.assert_next_some("worm-node/foo");
5972 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
5976 .consume()
5977 .scope(&["worm-node/foo".into()])
5978 .expect("should create worm-node/foo consumer")
5979 .announced();
5980
5981 foo_consumer.assert_next_some("worm-node/foo");
5982 foo_consumer.assert_next_wait(); }
5984
5985 #[tokio::test]
5986 async fn test_select_multiple_roots_with_empty_prefix() {
5987 let origin = Origin::random().produce();
5988
5989 let limited_producer = origin
5991 .scope(&["app1".into(), "app2".into(), "shared".into()])
5992 .expect("should create limited producer");
5993
5994 let _broadcast1 = limited_producer
5996 .create_broadcast("app1/data", announce())
5997 .expect("publish allowed");
5998 let _broadcast2 = limited_producer
5999 .create_broadcast("app2/config", announce())
6000 .expect("publish allowed");
6001 let _broadcast3 = limited_producer
6002 .create_broadcast("shared/resource", announce())
6003 .expect("publish allowed");
6004 settle().await;
6005
6006 let mut consumer = limited_producer
6008 .consume()
6009 .scope(&["".into()])
6010 .expect("should create consumer with empty prefix")
6011 .announced();
6012
6013 consumer.assert_next_some("app1/data");
6015 consumer.assert_next_some("app2/config");
6016 consumer.assert_next_some("shared/resource");
6017 consumer.assert_next_wait();
6018 }
6019
6020 #[tokio::test]
6021 async fn test_publish_scope_with_empty_prefix() {
6022 let origin = Origin::random().produce();
6023
6024 let limited_producer = origin
6026 .scope(&["services/api".into(), "services/web".into()])
6027 .expect("should create limited producer");
6028
6029 let same_producer = limited_producer
6031 .scope(&["".into()])
6032 .expect("should create producer with empty prefix");
6033
6034 let _broadcast = same_producer
6036 .create_broadcast("services/api", announce())
6037 .expect("publish allowed");
6038 let _keep2 = same_producer
6039 .create_broadcast("services/web", announce())
6040 .expect("publish allowed");
6041 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
6042 assert!(same_producer.create_broadcast("other", announce()).is_err());
6043 }
6044
6045 #[tokio::test]
6046 async fn test_select_narrowing_to_deeper_path() {
6047 let origin = Origin::random().produce();
6048
6049 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
6051
6052 let _broadcast1 = limited_producer
6054 .create_broadcast("org/team1/project1", announce())
6055 .expect("publish allowed");
6056 let _broadcast2 = limited_producer
6057 .create_broadcast("org/team1/project2", announce())
6058 .expect("publish allowed");
6059 let _broadcast3 = limited_producer
6060 .create_broadcast("org/team2/project1", announce())
6061 .expect("publish allowed");
6062 settle().await;
6063
6064 let mut team2_consumer = limited_producer
6066 .consume()
6067 .scope(&["org/team2".into()])
6068 .expect("should create team2 consumer")
6069 .announced();
6070
6071 team2_consumer.assert_next_some("org/team2/project1");
6072 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
6076 .consume()
6077 .scope(&["org/team1/project1".into()])
6078 .expect("should create project1 consumer")
6079 .announced();
6080
6081 project1_consumer.assert_next_some("org/team1/project1");
6083 project1_consumer.assert_next_wait();
6084 }
6085
6086 #[tokio::test]
6087 async fn test_select_with_non_matching_prefix() {
6088 let origin = Origin::random().produce();
6089
6090 let limited_producer = origin
6092 .scope(&["allowed/path".into()])
6093 .expect("should create limited producer");
6094
6095 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
6097
6098 assert!(limited_producer.scope(&["other/path".into()]).is_none());
6100 }
6101
6102 #[tokio::test]
6105 async fn test_with_root_trailing_slash_consumer() {
6106 let origin = Origin::random().produce();
6107
6108 let prefix = "some_prefix/".to_string();
6110 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6111
6112 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6113 settle().await;
6114 consumer.assert_next_some("test");
6115 }
6116
6117 #[tokio::test]
6119 async fn test_with_root_trailing_slash_producer() {
6120 let origin = Origin::random().produce();
6121
6122 let prefix = "some_prefix/".to_string();
6124 let rooted = origin.with_root(prefix).unwrap();
6125
6126 let _b = rooted.create_broadcast("test", announce()).unwrap();
6127 settle().await;
6128
6129 let mut consumer = rooted.consume().announced();
6130 consumer.assert_next_some("test");
6131 }
6132
6133 #[tokio::test]
6135 async fn test_with_root_trailing_slash_unannounce() {
6136 tokio::time::pause();
6137
6138 let origin = Origin::random().produce();
6139
6140 let prefix = "some_prefix/".to_string();
6141 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6142
6143 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6144 settle().await;
6145 consumer.assert_next_some("test");
6146
6147 b.finish();
6149 settle().await;
6150
6151 consumer.assert_next_none("test");
6153 }
6154
6155 #[tokio::test]
6156 async fn test_select_maintains_access_with_wider_prefix() {
6157 let origin = Origin::random().produce();
6158
6159 let demo_producer = origin.with_root("demo").expect("should create demo root");
6161 let user_producer = demo_producer
6162 .scope(&["worm-node".into(), "foobar".into()])
6163 .expect("should create user producer");
6164
6165 let _broadcast1 = user_producer
6167 .create_broadcast("worm-node/data", announce())
6168 .expect("publish allowed");
6169 let _broadcast2 = user_producer
6170 .create_broadcast("foobar", announce())
6171 .expect("publish allowed");
6172 settle().await;
6173
6174 let mut consumer = user_producer
6176 .consume()
6177 .scope(&["".into()])
6178 .expect("scope with empty prefix should not fail when user has specific permissions")
6179 .announced();
6180
6181 let a1 = consumer.try_next().expect("expected first announcement");
6183 let a2 = consumer.try_next().expect("expected second announcement");
6184 consumer.assert_next_wait();
6185
6186 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6187 paths.sort();
6188 assert_eq!(paths, ["foobar", "worm-node/data"]);
6189
6190 let mut narrow_consumer = user_producer
6192 .consume()
6193 .scope(&["worm-node".into()])
6194 .expect("should be able to narrow scope to worm-node")
6195 .announced();
6196
6197 narrow_consumer.assert_next_some("worm-node/data");
6198 narrow_consumer.assert_next_wait(); }
6200
6201 #[tokio::test]
6202 async fn test_duplicate_prefixes_deduped() {
6203 let origin = Origin::random().produce();
6204
6205 let producer = origin
6207 .scope(&["demo".into(), "demo".into()])
6208 .expect("should create producer");
6209
6210 let _broadcast = producer
6211 .create_broadcast("demo/stream", announce())
6212 .expect("publish allowed");
6213 settle().await;
6214
6215 let mut consumer = producer.consume().announced();
6216 consumer.assert_next_some("demo/stream");
6217 consumer.assert_next_wait();
6218 }
6219
6220 #[tokio::test]
6221 async fn test_overlapping_prefixes_deduped() {
6222 let origin = Origin::random().produce();
6223
6224 let producer = origin
6226 .scope(&["demo".into(), "demo/foo".into()])
6227 .expect("should create producer");
6228
6229 let _broadcast = producer
6231 .create_broadcast("demo/bar/stream", announce())
6232 .expect("publish allowed");
6233 settle().await;
6234
6235 let mut consumer = producer.consume().announced();
6236 consumer.assert_next_some("demo/bar/stream");
6237 consumer.assert_next_wait();
6238 }
6239
6240 #[tokio::test]
6241 async fn test_overlapping_prefixes_no_duplicate_announcements() {
6242 let origin = Origin::random().produce();
6243
6244 let producer = origin
6246 .scope(&["demo".into(), "demo/foo".into()])
6247 .expect("should create producer");
6248
6249 let _broadcast = producer
6250 .create_broadcast("demo/foo/stream", announce())
6251 .expect("publish allowed");
6252 settle().await;
6253
6254 let mut consumer = producer.consume().announced();
6255 consumer.assert_next_some("demo/foo/stream");
6257 consumer.assert_next_wait();
6258 }
6259
6260 #[tokio::test]
6261 async fn test_allowed_returns_deduped_prefixes() {
6262 let origin = Origin::random().produce();
6263
6264 let producer = origin
6265 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
6266 .expect("should create producer");
6267
6268 let allowed: Vec<_> = producer.allowed().collect();
6269 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
6270 }
6271
6272 #[tokio::test]
6273 async fn test_announced_broadcast_already_announced() {
6274 let origin = Origin::random().produce();
6275
6276 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6277 settle().await;
6278
6279 let consumer = origin.consume();
6280 let result = consumer.announced_broadcast("test").await.expect("should find it");
6281 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6282 }
6283
6284 #[tokio::test]
6285 async fn test_announced_broadcast_delayed() {
6286 tokio::time::pause();
6287
6288 let origin = Origin::random().produce();
6289
6290 let consumer = origin.consume();
6291
6292 let wait = tokio::spawn({
6294 let consumer = consumer.clone();
6295 async move { consumer.announced_broadcast("test").await }
6296 });
6297
6298 tokio::task::yield_now().await;
6300
6301 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6302 settle().await;
6303
6304 let result = wait.await.unwrap().expect("should find it");
6305 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6306 }
6307
6308 #[tokio::test]
6309 async fn test_announced_broadcast_ignores_unrelated_paths() {
6310 tokio::time::pause();
6311
6312 let origin = Origin::random().produce();
6313
6314 let consumer = origin.consume();
6315
6316 let wait = tokio::spawn({
6317 let consumer = consumer.clone();
6318 async move { consumer.announced_broadcast("target").await }
6319 });
6320
6321 tokio::task::yield_now().await;
6322
6323 let _other = origin.create_broadcast("other", announce()).unwrap();
6325 settle().await;
6326 tokio::task::yield_now().await;
6327 assert!(!wait.is_finished(), "must not resolve on unrelated path");
6328
6329 let _target = origin.create_broadcast("target", announce()).unwrap();
6330 settle().await;
6331 let result = wait.await.unwrap().expect("should find target");
6332 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
6333 }
6334
6335 #[tokio::test]
6336 async fn test_announced_broadcast_skips_nested_paths() {
6337 tokio::time::pause();
6338
6339 let origin = Origin::random().produce();
6340
6341 let consumer = origin.consume();
6342
6343 let wait = tokio::spawn({
6344 let consumer = consumer.clone();
6345 async move { consumer.announced_broadcast("foo").await }
6346 });
6347
6348 tokio::task::yield_now().await;
6349
6350 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
6352 settle().await;
6353 tokio::task::yield_now().await;
6354 assert!(!wait.is_finished(), "must not resolve on a nested path");
6355
6356 let _exact = origin.create_broadcast("foo", announce()).unwrap();
6357 settle().await;
6358 let result = wait.await.unwrap().expect("should find foo exactly");
6359 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
6360 }
6361
6362 #[tokio::test]
6363 async fn test_announced_broadcast_disallowed() {
6364 let origin = Origin::random().produce();
6365 let limited = origin
6366 .consume()
6367 .scope(&["allowed".into()])
6368 .expect("should create limited");
6369
6370 assert!(limited.announced_broadcast("notallowed").await.is_none());
6372 }
6373
6374 #[tokio::test]
6375 async fn test_announced_broadcast_scope_too_narrow() {
6376 let origin = Origin::random().produce();
6379 let limited = origin
6380 .consume()
6381 .scope(&["foo/specific".into()])
6382 .expect("should create limited");
6383
6384 let result = limited
6386 .announced_broadcast("foo")
6387 .now_or_never()
6388 .expect("must not block");
6389 assert!(result.is_none());
6390 }
6391
6392 #[tokio::test]
6396 async fn test_coalesce_announce_then_unannounce() {
6397 tokio::time::pause();
6399
6400 let origin = Origin::random().produce();
6401 let mut announced = origin.consume().announced();
6402
6403 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6404 settle().await;
6405 broadcast.finish();
6406
6407 settle().await;
6408
6409 announced.assert_next_wait();
6410 }
6411
6412 #[tokio::test]
6413 async fn test_coalesce_announce_unannounce_announce() {
6414 tokio::time::pause();
6417
6418 let origin = Origin::random().produce();
6419 let mut announced = origin.consume().announced();
6420
6421 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6422 settle().await;
6423 broadcast1.finish();
6424 settle().await;
6425 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6426 settle().await;
6427
6428 announced.assert_next_some("test");
6429 announced.assert_next_wait();
6430 }
6431
6432 #[tokio::test]
6433 async fn test_coalesce_unannounce_announce_preserved() {
6434 tokio::time::pause();
6437
6438 let origin = Origin::random().produce();
6439 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6440 settle().await;
6441
6442 let mut announced = origin.consume().announced();
6443 announced.assert_next_some("test");
6444
6445 broadcast1.finish();
6447 settle().await;
6448
6449 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6450 settle().await;
6451
6452 announced.assert_next_none("test");
6454 announced.assert_next_some("test");
6455 announced.assert_next_wait();
6456 }
6457
6458 #[tokio::test]
6459 async fn test_coalesce_unannounce_announce_unannounce() {
6460 tokio::time::pause();
6463
6464 let origin = Origin::random().produce();
6465 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6466 settle().await;
6467
6468 let mut announced = origin.consume().announced();
6469 announced.assert_next_some("test");
6470
6471 broadcast1.finish();
6472 settle().await;
6473
6474 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6475 settle().await;
6476 broadcast2.finish();
6477 settle().await;
6478
6479 announced.assert_next_none("test");
6480 announced.assert_next_wait();
6481 }
6482
6483 #[tokio::test]
6484 async fn test_coalesce_churn_bounded() {
6485 tokio::time::pause();
6490
6491 let origin = Origin::random().produce();
6492 let mut announced = origin.consume().announced();
6493
6494 for _ in 0..1000 {
6495 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6496 settle().await;
6497 broadcast.finish();
6498 }
6499 settle().await;
6500
6501 let mut collected = Vec::new();
6502 while let Some(update) = announced.try_next() {
6503 collected.push(update);
6504 }
6505 assert!(
6506 collected.len() <= 1,
6507 "expected at most one pending update, got {}",
6508 collected.len()
6509 );
6510 assert!(
6511 collected.iter().all(|a| a.path == Path::new("test")),
6512 "unexpected path in pending updates",
6513 );
6514 }
6515
6516 #[tokio::test]
6520 async fn test_consumer_clone_is_side_effect_free() {
6521 let origin = Origin::random().produce();
6522
6523 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
6524 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
6525 settle().await;
6526
6527 let consumer = origin.consume();
6528 let mut announced = consumer.announced();
6529
6530 for _ in 0..16 {
6533 let cloned = consumer.clone();
6534 assert!(cloned.get_broadcast("test1").is_some());
6535 assert!(cloned.get_broadcast("test2").is_some());
6536 }
6537
6538 let a1 = announced.try_next().expect("first announcement");
6541 let a2 = announced.try_next().expect("second announcement");
6542 announced.assert_next_wait();
6543
6544 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6545 paths.sort();
6546 assert_eq!(paths, ["test1", "test2"]);
6547
6548 let mut fresh = consumer.announced();
6550 let b1 = fresh.try_next().expect("backlog: first");
6551 let b2 = fresh.try_next().expect("backlog: second");
6552 fresh.assert_next_wait();
6553
6554 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
6555 paths.sort();
6556 assert_eq!(paths, ["test1", "test2"]);
6557 }
6558
6559 #[tokio::test]
6561 async fn dynamic_request_unroutable_without_handler() {
6562 let origin = Origin::random().produce();
6563 let consumer = origin.consume();
6564 assert!(matches!(
6565 consumer.request_broadcast("missing").await,
6566 Err(Error::Unroutable)
6567 ));
6568 }
6569
6570 #[tokio::test(start_paused = true)]
6573 async fn dynamic_request_served_not_announced() {
6574 let origin = Origin::random().produce();
6575 let mut dynamic = origin.dynamic();
6576 let consumer = origin.consume();
6577
6578 let mut announced = origin.consume().announced();
6580 announced.assert_next_wait();
6581
6582 let served = broadcast::Info::new().produce();
6583 let request_fut = consumer.request_broadcast("fallback");
6586
6587 let mut served_dynamic = served.dynamic();
6589
6590 let request = dynamic.requested_broadcast().await.unwrap();
6591 assert_eq!(request.path(), &Path::new("fallback"));
6592 request.accept(&served);
6593
6594 let broadcast = request_fut.await.unwrap();
6595 assert!(broadcast.is_clone(&served.consume()));
6596
6597 let track_fut = broadcast.track("video").unwrap().subscribe(None);
6599 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
6600 let mut track = track_fut.await.unwrap();
6601 producer.append_group().unwrap();
6602 track.assert_group();
6603
6604 announced.assert_next_wait();
6606 }
6607
6608 #[tokio::test(start_paused = true)]
6610 async fn dynamic_request_coalesces() {
6611 let origin = Origin::random().produce();
6612 let mut dynamic = origin.dynamic();
6613 let consumer = origin.consume();
6614
6615 let f1 = consumer.request_broadcast("dup");
6617 let f2 = consumer.request_broadcast("dup");
6618
6619 let request = dynamic.requested_broadcast().await.unwrap();
6621 assert_eq!(request.path(), &Path::new("dup"));
6622 assert!(
6623 dynamic.requested_broadcast().now_or_never().is_none(),
6624 "a coalesced request must not be served twice"
6625 );
6626
6627 let served = broadcast::Info::new().produce();
6629 request.accept(&served);
6630 assert!(f1.await.unwrap().is_clone(&served.consume()));
6631 assert!(f2.await.unwrap().is_clone(&served.consume()));
6632 }
6633
6634 #[tokio::test(start_paused = true)]
6637 async fn dynamic_request_dedups_served() {
6638 let origin = Origin::random().produce();
6639 let mut dynamic = origin.dynamic();
6640 let consumer = origin.consume();
6641
6642 let request_fut = consumer.request_broadcast("fallback");
6643 let request = dynamic.requested_broadcast().await.unwrap();
6644 let served = broadcast::Info::new().produce();
6645 request.accept(&served);
6646 let first = request_fut.await.unwrap();
6647 assert!(first.is_clone(&served.consume()));
6648
6649 let second = consumer.request_broadcast("fallback").await.unwrap();
6651 assert!(second.is_clone(&served.consume()));
6652
6653 assert!(
6655 dynamic.requested_broadcast().now_or_never().is_none(),
6656 "a still-live served broadcast must not be re-requested from the handler"
6657 );
6658 }
6659
6660 #[tokio::test(start_paused = true)]
6662 async fn dynamic_request_reserves_after_close() {
6663 let origin = Origin::random().produce();
6664 let mut dynamic = origin.dynamic();
6665 let consumer = origin.consume();
6666
6667 let request_fut = consumer.request_broadcast("fallback");
6668 let request = dynamic.requested_broadcast().await.unwrap();
6669 let served = broadcast::Info::new().produce();
6670 request.accept(&served);
6671 request_fut.await.unwrap();
6672
6673 drop(served);
6675
6676 let request_fut = consumer.request_broadcast("fallback");
6678 let request = dynamic.requested_broadcast().await.unwrap();
6679 assert_eq!(request.path(), &Path::new("fallback"));
6680 let served = broadcast::Info::new().produce();
6681 request.accept(&served);
6682 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6683 }
6684
6685 #[tokio::test(start_paused = true)]
6688 async fn dynamic_request_served_cache_bounded() {
6689 let origin = Origin::random().produce();
6690 let mut dynamic = origin.dynamic();
6691 let consumer = origin.consume();
6692
6693 for i in 0..100 {
6694 let path = format!("one-shot/{i}");
6695 let request_fut = consumer.request_broadcast(&path);
6696 let request = dynamic.requested_broadcast().await.unwrap();
6697 let served = broadcast::Info::new().produce();
6698 request.accept(&served);
6699 request_fut.await.unwrap();
6700 drop(served);
6702 }
6703
6704 assert!(
6707 origin.dynamic.read().served.len() <= 4,
6708 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
6709 origin.dynamic.read().served.len()
6710 );
6711 }
6712
6713 #[tokio::test(start_paused = true)]
6716 async fn dynamic_request_coalesces_after_handoff() {
6717 let origin = Origin::random().produce();
6718 let mut dynamic = origin.dynamic();
6719 let consumer = origin.consume();
6720
6721 let f1 = consumer.request_broadcast("fallback");
6722 let request = dynamic.requested_broadcast().await.unwrap();
6724
6725 let f2 = consumer.request_broadcast("fallback");
6727 assert!(
6728 dynamic.requested_broadcast().now_or_never().is_none(),
6729 "a repeat request during hand-off must coalesce, not re-queue"
6730 );
6731
6732 let served = broadcast::Info::new().produce();
6734 request.accept(&served);
6735 assert!(f1.await.unwrap().is_clone(&served.consume()));
6736 assert!(f2.await.unwrap().is_clone(&served.consume()));
6737 }
6738
6739 #[tokio::test(start_paused = true)]
6741 async fn dynamic_request_dropped_after_handoff() {
6742 let origin = Origin::random().produce();
6743 let mut dynamic = origin.dynamic();
6744 let consumer = origin.consume();
6745
6746 let f1 = consumer.request_broadcast("fallback");
6747 let request = dynamic.requested_broadcast().await.unwrap();
6748 let f2 = consumer.request_broadcast("fallback");
6749
6750 drop(request);
6752 assert!(matches!(f1.await, Err(Error::Unroutable)));
6753 assert!(matches!(f2.await, Err(Error::Unroutable)));
6754 }
6755
6756 #[tokio::test(start_paused = true)]
6758 async fn dynamic_request_rejected() {
6759 let origin = Origin::random().produce();
6760 let mut dynamic = origin.dynamic();
6761 let consumer = origin.consume();
6762
6763 let request_fut = consumer.request_broadcast("fallback");
6764
6765 let request = dynamic.requested_broadcast().await.unwrap();
6766 request.reject(Error::Cancel);
6767
6768 assert!(matches!(request_fut.await, Err(Error::Cancel)));
6769 }
6770
6771 #[tokio::test(start_paused = true)]
6775 async fn dynamic_request_rerequest_after_reject() {
6776 let origin = Origin::random().produce();
6777 let mut dynamic = origin.dynamic();
6778 let consumer = origin.consume();
6779
6780 let f1 = consumer.request_broadcast("fallback");
6781 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
6782 assert!(matches!(f1.await, Err(Error::Unroutable)));
6783
6784 let served = broadcast::Info::new().produce();
6785 let f2 = consumer.request_broadcast("fallback");
6787 let request = dynamic.requested_broadcast().await.unwrap();
6788 assert_eq!(request.path(), &Path::new("fallback"));
6789 request.accept(&served);
6790 assert!(f2.await.unwrap().is_clone(&served.consume()));
6791 }
6792
6793 #[tokio::test(start_paused = true)]
6796 async fn dynamic_request_handler_dropped() {
6797 let origin = Origin::random().produce();
6798 let dynamic = origin.dynamic();
6799 let consumer = origin.consume();
6800
6801 let request_fut = consumer.request_broadcast("fallback");
6802 drop(dynamic);
6803 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
6804
6805 assert!(matches!(
6807 consumer.request_broadcast("again").await,
6808 Err(Error::Unroutable)
6809 ));
6810 }
6811
6812 #[tokio::test(start_paused = true)]
6816 async fn dynamic_request_accept_after_handler_dropped() {
6817 let origin = Origin::random().produce();
6818 let mut dynamic = origin.dynamic();
6819 let consumer = origin.consume();
6820
6821 let request_fut = consumer.request_broadcast("fallback");
6822
6823 let request = dynamic.requested_broadcast().await.unwrap();
6825 drop(dynamic);
6826
6827 let served = broadcast::Info::new().produce();
6828 request.accept(&served);
6830 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6831 }
6832
6833 #[tokio::test(start_paused = true)]
6835 async fn dynamic_request_prefers_announced() {
6836 let origin = Origin::random().produce();
6837 let mut dynamic = origin.dynamic();
6838 let consumer = origin.consume();
6839
6840 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
6841 settle().await;
6842
6843 let got = consumer.request_broadcast("live").await.unwrap();
6844 assert!(
6845 got.is_clone(&consumer.get_broadcast("live").unwrap()),
6846 "should return the published broadcast"
6847 );
6848 assert!(
6849 dynamic.requested_broadcast().now_or_never().is_none(),
6850 "a published path must not queue a fallback request"
6851 );
6852 }
6853
6854 #[tokio::test(start_paused = true)]
6856 async fn dynamic_clone_keeps_alive() {
6857 let origin = Origin::random().produce();
6858 let dynamic = origin.dynamic();
6859 let consumer = origin.consume();
6860
6861 drop(dynamic.clone());
6862
6863 let request_fut = consumer.request_broadcast("fallback");
6866 assert!(
6867 request_fut.now_or_never().is_none(),
6868 "request should stay pending until served"
6869 );
6870 }
6871
6872 fn wedge_watchdog<F>(name: &str, secs: u64, scenario: F)
6878 where
6879 F: std::future::Future<Output = ()> + Send + 'static,
6880 {
6881 let (done_tx, done_rx) = std::sync::mpsc::channel::<()>();
6882 let handle = std::thread::spawn(move || {
6883 let rt = ::tokio::runtime::Builder::new_current_thread()
6884 .enable_time()
6885 .build()
6886 .unwrap();
6887 rt.block_on(scenario);
6888 let _ = done_tx.send(());
6889 });
6890 match done_rx.recv_timeout(std::time::Duration::from_secs(secs)) {
6891 Ok(()) => {
6892 let _ = handle.join();
6893 }
6894 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
6895 let err = handle.join().unwrap_err();
6897 std::panic::resume_unwind(err);
6898 }
6899 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
6900 panic!("{name}: scenario wedged; a task is spinning inside a single poll")
6901 }
6902 }
6903 }
6904
6905 #[test]
6918 fn test_active_corpse_does_not_livelock_takeover() {
6919 wedge_watchdog("active-corpse", 20, async {
6920 let origin = Origin::random().produce();
6921 let consumer = origin.consume();
6922 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
6923
6924 let source_a = origin
6927 .create_broadcast("test", announce().with_hops(hops.clone()))
6928 .unwrap();
6929 let mut dynamic_a = source_a.dynamic();
6930 settle().await;
6931 settle().await;
6932
6933 let broadcast = consumer.request_broadcast("test").await.unwrap();
6934 let subscribing = broadcast.track("video").unwrap().subscribe(None);
6935 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
6936 settle().await;
6937 let mut sub = subscribing.await.unwrap();
6938 producer_a.append_group().unwrap();
6939 sub.assert_group();
6940
6941 let source_b = origin
6945 .create_broadcast("test", announce().with_hops(hops.clone()))
6946 .unwrap();
6947 let dynamic_b = source_b.dynamic();
6948 settle().await;
6949
6950 drop(dynamic_b);
6958 source_b.abort(Error::Dropped).unwrap();
6959 settle().await;
6960
6961 producer_a.append_group().unwrap();
6964 sub.assert_group();
6965 sub.assert_not_closed();
6966 });
6967 }
6968
6969 #[test]
6977 fn test_route_churn_never_wedges() {
6978 for seed in 1..=8u64 {
6979 wedge_watchdog(&format!("churn seed {seed}"), 30, churn_scenario(seed));
6980 }
6981 }
6982
6983 async fn churn_scenario(seed: u64) {
6984 let mut rng = seed.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
6986 let mut next = move || {
6987 rng = rng.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
6988 rng >> 33
6989 };
6990
6991 let origin = Origin::random().produce();
6992 let consumer = origin.consume();
6993 let names: Vec<Arc<str>> = (0..8).map(|i| Arc::from(format!("t{i}"))).collect();
6994
6995 let mut subs: Vec<track::Subscriber> = Vec::new();
6996 let mut pending_subs: Vec<kio::Pending<track::Subscribing>> = Vec::new();
6997
6998 struct Source {
6999 producer: Option<broadcast::Producer>,
7000 server: ::tokio::task::JoinHandle<()>,
7001 }
7002 let mut sources: Vec<Source> = Vec::new();
7003
7004 for step in 0..400u64 {
7005 match next() % 10 {
7006 0 | 1 => {
7009 if sources.len() >= 3 {
7010 continue;
7011 }
7012 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
7013 let route = announce().with_hops(hops).with_cost(next() % 4);
7014 let Ok(source) = origin.create_broadcast("test", route) else {
7015 continue;
7016 };
7017 let mut dynamic = source.dynamic();
7018 let behavior = next();
7019 let server = ::tokio::spawn(async move {
7020 let mut round = 0u64;
7021 let mut kept: Vec<track::Producer> = Vec::new();
7022 while let Ok(request) = dynamic.requested_track().await {
7023 round += 1;
7024 match (behavior >> (round % 16)) % 4 {
7025 0 => drop(request), 1 => {
7027 let mut producer = request.accept(None);
7028 let _ = producer.create_group(group::Info { sequence: round });
7029 let _ = producer.abort(Error::Dropped);
7030 }
7031 2 => {
7032 let mut producer = request.accept(None);
7033 let _ = producer.create_group(group::Info { sequence: round });
7034 kept.push(producer);
7035 }
7036 _ => {
7037 let mut producer = request.accept(None);
7038 let _ = producer.finish();
7039 }
7040 }
7041 }
7042 });
7043 sources.push(Source {
7044 producer: Some(source),
7045 server,
7046 });
7047 }
7048 2 | 3 => {
7050 if sources.is_empty() {
7051 continue;
7052 }
7053 let i = (next() as usize) % sources.len();
7054 let mut source = sources.swap_remove(i);
7055 if next() % 2 == 0
7056 && let Some(producer) = source.producer.take()
7057 {
7058 let _ = producer.abort(Error::Dropped);
7059 }
7060 source.server.abort();
7061 }
7062 4..=6 => {
7064 if subs.len() + pending_subs.len() >= 24 {
7065 continue;
7066 }
7067 let Some(broadcast) = consumer.get_broadcast("test") else {
7068 continue;
7069 };
7070 let name = &names[(next() as usize) % names.len()];
7071 if let Ok(track) = broadcast.track(name.as_ref()) {
7072 pending_subs.push(track.subscribe(None));
7073 }
7074 }
7075 7 => {
7077 if subs.is_empty() {
7078 continue;
7079 }
7080 let i = (next() as usize) % subs.len();
7081 subs.swap_remove(i);
7082 }
7083 _ => {
7085 for sub in pending_subs.drain(..) {
7086 match ::tokio::time::timeout(std::time::Duration::from_millis(5), sub).await {
7087 Ok(Ok(sub)) => subs.push(sub),
7088 Ok(Err(_)) => {}
7089 Err(_) => {}
7091 }
7092 }
7093 for sub in subs.iter_mut() {
7094 while let Some(Ok(Some(_))) = sub.recv_group().now_or_never() {}
7095 }
7096 }
7097 }
7098 if step % 16 == 0 {
7099 settle().await;
7100 }
7101 for _ in 0..(next() % 3) {
7103 ::tokio::task::yield_now().await;
7104 }
7105 }
7106
7107 for source in sources {
7108 source.server.abort();
7109 }
7110 settle().await;
7111 }
7112}