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 let mut may_take_over = true;
1575
1576 'attach: loop {
1577 let leaf = if rest.is_empty() {
1582 node.clone()
1583 } else {
1584 node.lock().leaf(&rest)
1585 };
1586
1587 let (state, broadcast, id) = match attach_source(&ctx, &leaf, &source, route.clone(), may_take_over) {
1588 Attach::Ready(state, broadcast, id) => (state, broadcast, id),
1589 Attach::Parked(incumbent) => {
1590 tracing::debug!(
1591 broadcast = %full,
1592 "path already live with a different publisher; parking this source until it ends",
1593 );
1594 let update = kio::wait(|waiter| {
1598 if let Poll::Ready(update) = source.poll_route_changed(waiter) {
1599 return Poll::Ready(Some(update));
1600 }
1601 match incumbent.poll(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending }) {
1604 Poll::Ready(_) => Poll::Ready(None),
1605 Poll::Pending => Poll::Pending,
1606 }
1607 })
1608 .await;
1609 match update {
1610 Some(Ok(update)) => {
1612 sync_announce(&mut announce, update.announce, &ingress);
1613 if update.hops.iter().next().copied() != route.hops.iter().next().copied() {
1618 may_take_over = true;
1619 }
1620 route = update;
1621 }
1622 Some(Err(_)) => return,
1624 None => {}
1626 }
1627 continue 'attach;
1628 }
1629 };
1630 let publisher = route.hops.iter().next().copied();
1631
1632 loop {
1633 let update = kio::wait(|waiter| {
1634 if state
1640 .poll_ref(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending })
1641 .is_ready()
1642 {
1643 return Poll::Ready(None);
1644 }
1645 source.poll_route_changed(waiter).map(Some)
1646 })
1647 .await;
1648 match update {
1649 None => {
1650 may_take_over = false;
1653 continue 'attach;
1654 }
1655 Some(Ok(update)) => {
1656 let announced = update.announce;
1657 if update.hops.iter().next().copied() != publisher {
1668 detach_source(&state, &broadcast, &leaf, id, true);
1669 sync_announce(&mut announce, announced, &ingress);
1670 may_take_over = true;
1674 route = update;
1675 continue 'attach;
1676 }
1677 {
1678 let carrying = broadcast.demand().is_used();
1679 let Ok(mut s) = state.write() else { return };
1680 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1681 return;
1682 };
1683 if entry.route == update {
1684 continue;
1685 }
1686 entry.route = update;
1687 s.reselect(carrying);
1688 }
1689 sync_announce(&mut announce, announced, &ingress);
1691 sync_front(&state, &broadcast, &leaf);
1692 }
1693 Some(Err(_)) => {
1694 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1697 return;
1698 }
1699 }
1700 }
1701 }
1702}
1703
1704enum Attach {
1706 Ready(kio::Producer<FrontState>, broadcast::Producer, u64),
1709 Parked(kio::Producer<FrontState>),
1716}
1717
1718struct AttachContext<'a> {
1720 origin: &'a Info,
1721 node: &'a Lock<OriginNode>,
1722 full: &'a PathOwned,
1724 rest: &'a PathOwned,
1726}
1727
1728fn same_publisher(a: Option<Origin>, b: Option<Origin>) -> bool {
1737 if a == Some(Origin::UNKNOWN) || b == Some(Origin::UNKNOWN) {
1738 return false;
1739 }
1740 a == b
1741}
1742
1743fn attach_source(
1769 ctx: &AttachContext,
1770 leaf: &Lock<OriginNode>,
1771 source: &broadcast::Consumer,
1772 route: broadcast::Route,
1773 may_take_over: bool,
1774) -> Attach {
1775 let publisher = route.hops.iter().next().copied();
1776 let mut leaf_guard = leaf.lock();
1777
1778 if let Some(existing) = &leaf_guard.broadcast {
1781 let mut joined = None;
1782 let carrying = existing.broadcast.demand().is_used();
1783 if let Ok(mut s) = existing.state.write()
1784 && !s.closed
1785 {
1786 if same_publisher(s.publisher, publisher) {
1787 let id = s.next_route;
1788 s.next_route += 1;
1789 s.routes.push(FrontRoute {
1790 id,
1791 route: route.clone(),
1792 source: source.clone(),
1793 });
1794 s.reselect(carrying);
1795 joined = Some(id);
1796 } else if !may_take_over || !route.announce || s.taints_a_reader(&route) {
1797 return Attach::Parked(existing.state.clone());
1798 } else {
1799 s.closed = true;
1807 tracing::warn!(broadcast = %ctx.full, "replacing a live broadcast from a different publisher");
1808 }
1809 }
1810 if let Some(id) = joined {
1811 let state = existing.state.clone();
1812 let broadcast = existing.broadcast.clone();
1813 drop(leaf_guard);
1814 sync_front(&state, &broadcast, leaf);
1815 return Attach::Ready(state, broadcast, id);
1816 }
1817 }
1818
1819 let announce = route.announce;
1821 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
1822 origin: ctx.origin.clone(),
1823 });
1824 let _ = broadcast.clone().set_route(route.clone());
1825 let state = kio::Producer::new(FrontState {
1826 path: ctx.full.clone(),
1827 self_origin: ctx.origin.id,
1828 publisher,
1829 next_route: 1,
1830 excluded: HashMap::new(),
1831 routes: vec![FrontRoute {
1832 id: 0,
1833 route,
1834 source: source.clone(),
1835 }],
1836 active: Some(0),
1837 linger: ctx.origin.linger,
1838 closed: false,
1839 });
1840
1841 if let Some(stale) = leaf_guard.broadcast.take()
1845 && stale.announced
1846 {
1847 leaf_guard.notify.lock().unannounce(&stale.path);
1848 }
1849 let entry = OriginBroadcast {
1850 path: ctx.full.clone(),
1851 broadcast: broadcast.clone(),
1852 state: state.clone(),
1853 announced: announce,
1854 };
1855 if entry.announced {
1856 leaf_guard
1857 .notify
1858 .lock()
1859 .announce(ctx.full, &broadcast.consume(), &state);
1860 }
1861 leaf_guard.broadcast = Some(entry);
1862 drop(leaf_guard);
1863
1864 web_async::spawn(run_front(
1865 state.clone(),
1866 broadcast.clone(),
1867 ctx.node.clone(),
1868 ctx.rest.clone(),
1869 ));
1870
1871 Attach::Ready(state, broadcast, 0)
1872}
1873
1874async fn run_front(
1877 state: kio::Producer<FrontState>,
1878 mut broadcast: broadcast::Producer,
1879 node: Lock<OriginNode>,
1880 rest: PathOwned,
1881) {
1882 enum Step {
1883 Serve(Arc<str>, super::resume::Producer),
1884 Changed,
1886 Expired,
1888 Closed,
1889 }
1890
1891 let linger = state.read().linger;
1892 let mut deadline = kio::time::Deadline::new();
1897
1898 loop {
1899 let empty = {
1900 let s = state.read();
1901 !s.closed && s.routes.is_empty()
1902 };
1903 deadline.set(match (empty, deadline.deadline()) {
1904 (true, None) => web_async::time::Instant::now().checked_add(linger),
1907 (true, at) => at,
1908 (false, _) => None,
1909 });
1910
1911 let step = {
1912 kio::wait(|waiter| {
1913 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1914 return Poll::Ready(Step::Serve(name, resume));
1915 }
1916 match state.poll(waiter, |s| {
1919 if s.closed || s.routes.is_empty() != empty {
1920 Poll::Ready(())
1921 } else {
1922 Poll::Pending
1923 }
1924 }) {
1925 Poll::Ready(Ok(guard)) => {
1926 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1927 }
1928 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1929 Poll::Pending => {}
1930 }
1931 deadline.poll(waiter).map(|_| Step::Expired)
1932 })
1933 .await
1934 };
1935
1936 match step {
1937 Step::Serve(name, resume) => {
1938 web_async::spawn(serve_track(state.clone(), name, resume));
1941 }
1942 Step::Changed => {}
1943 Step::Expired => {
1944 let close = {
1948 let Ok(mut s) = state.write() else { break };
1949 if !s.closed && s.routes.is_empty() {
1950 s.closed = true;
1951 true
1952 } else {
1953 false
1954 }
1955 };
1956 if close {
1957 break;
1958 }
1959 }
1960 Step::Closed => break,
1961 }
1962 }
1963
1964 broadcast.abort_spliced(Error::Dropped);
1966
1967 broadcast.finish();
1969
1970 node.lock().remove(&state, &rest);
1973}
1974
1975async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1988 enum Step {
1989 Closed,
1990 Splice(u64, broadcast::Consumer),
1991 Complete,
1992 Failed(Error),
1993 NoRoute,
1998 Idle,
2000 Demand,
2002 }
2003
2004 let mut serving: Option<(u64, track::Consumer)> = None;
2006 let mut spliced_edge: Option<u64> = None;
2012 let mut refused: HashSet<u64> = HashSet::new();
2017 let mut refusal: Option<Error> = None;
2018 let mut dead: HashSet<u64> = HashSet::new();
2023 let mut idle_since: Option<web_async::time::Instant> = None;
2025 let mut deadline = kio::time::Deadline::new();
2026
2027 loop {
2028 let serving_id = serving.as_ref().map(|(id, _)| *id);
2029
2030 {
2037 let s = state.read();
2038 refused.retain(|id| s.routes.iter().any(|r| r.id == *id));
2039 dead.retain(|id| s.routes.iter().any(|r| r.id == *id));
2040 let exhausted = !s.routes.is_empty()
2041 && s.serve_route(|id| refused.contains(&id) || dead.contains(&id))
2042 .is_none();
2043 if exhausted && dead.is_empty() {
2044 drop(s);
2045 let err = refusal.take().unwrap_or(Error::NotFound);
2046 tracing::debug!(name = %name, %err, "every source refused track; aborting");
2047 let _ = resume.abort(err);
2048 return;
2049 }
2050 }
2051
2052 let used = resume.is_used();
2064 idle_since = match (resume.is_spliced(), used) {
2065 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
2066 _ => None,
2067 };
2068 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
2069
2070 let step = {
2071 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2072 kio::wait(|waiter| {
2073 match state.poll(waiter, |s| {
2079 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
2080 if s.closed
2081 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
2082 {
2083 Poll::Ready(())
2084 } else {
2085 Poll::Pending
2086 }
2087 }) {
2088 Poll::Ready(Ok(guard)) => {
2089 if guard.closed {
2090 return Poll::Ready(Step::Closed);
2091 }
2092 let Some(next) = guard.serve_route(skip) else {
2093 return Poll::Ready(Step::NoRoute);
2094 };
2095 let source = guard
2096 .routes
2097 .iter()
2098 .find(|r| r.id == next)
2099 .expect("servable source in table")
2100 .source
2101 .clone();
2102 return Poll::Ready(Step::Splice(next, source));
2103 }
2104 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
2105 Poll::Pending => {}
2106 }
2107
2108 let edge = match used {
2113 true => resume.poll_unused(waiter),
2114 false => resume.poll_used(waiter),
2115 };
2116 if edge.is_ready() {
2117 return Poll::Ready(Step::Demand);
2118 }
2119
2120 if let Some((_, track)) = &serving
2123 && let Poll::Ready(result) = track.poll_complete(waiter)
2124 {
2125 return Poll::Ready(match result {
2126 Ok(()) => Step::Complete,
2127 Err(err) => Step::Failed(err),
2128 });
2129 }
2130
2131 deadline.poll(waiter).map(|_| Step::Idle)
2132 })
2133 .await
2134 };
2135
2136 match step {
2137 Step::Closed => return,
2139 Step::Complete => {
2140 let _ = resume.finish();
2141 return;
2142 }
2143 Step::Failed(err) => {
2144 if resume.latest() == spliced_edge
2152 && let Some(id) = serving_id
2153 {
2154 let closing = state
2155 .read()
2156 .routes
2157 .iter()
2158 .find(|r| r.id == id)
2159 .is_some_and(|r| r.source.is_closing());
2160 if closing {
2161 dead.insert(id);
2162 } else {
2163 refused.insert(id);
2164 refusal = Some(err);
2165 }
2166 }
2167 serving = None;
2168 }
2169 Step::Demand => {}
2171 Step::NoRoute => serving = None,
2177 Step::Idle => {
2178 if resume.release().is_err() {
2183 return;
2185 }
2186 serving = None;
2187 }
2188 Step::Splice(id, source) => {
2189 let attempt = match source.track(&name) {
2193 Ok(track) => {
2194 let query = track.info().into_inner();
2197 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2198 let info = kio::wait(|waiter| {
2199 if let Poll::Ready(result) = query.poll(waiter) {
2200 return Poll::Ready(Some(result));
2201 }
2202 match state.poll(waiter, |s| {
2203 if s.closed || s.serve_route(skip) != Some(id) {
2204 Poll::Ready(())
2205 } else {
2206 Poll::Pending
2207 }
2208 }) {
2209 Poll::Ready(_) => Poll::Ready(None),
2210 Poll::Pending => Poll::Pending,
2211 }
2212 })
2213 .await;
2214 match info {
2215 None => continue,
2217 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2220 Poll::Ready(Err(err)) => Err(err),
2221 _ => Ok(track),
2222 },
2223 Some(Err(err)) => Err(err),
2224 }
2225 }
2226 Err(err) => Err(err),
2227 };
2228
2229 match attempt {
2230 Ok(track) => {
2231 if let Err(err) = resume.takeover(&track) {
2232 let _ = resume.abort(err);
2237 return;
2238 }
2239 spliced_edge = resume.latest();
2251 serving = Some((id, track));
2252 }
2253 Err(_) if source.is_closing() => {
2257 dead.insert(id);
2258 serving = None;
2259 }
2260 Err(err) => {
2266 tracing::debug!(name = %name, source = id, %err, "source refused track");
2267 refused.insert(id);
2268 refusal = Some(err);
2269 serving = None;
2270 }
2271 }
2272 }
2273 }
2274 }
2275}
2276
2277#[derive(Default)]
2283struct OriginDynamicState {
2284 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2287
2288 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2294}
2295
2296#[derive(Default)]
2303struct PendingBroadcast {
2304 resolved: Option<Result<broadcast::Consumer, Error>>,
2305}
2306
2307pub struct Dynamic {
2318 info: Origin,
2319 root: PathOwned,
2320 state: kio::Shared<OriginDynamicState>,
2321}
2322
2323impl Clone for Dynamic {
2324 fn clone(&self) -> Self {
2325 self.state.lock().requests.add_handler();
2329
2330 Self {
2331 info: self.info,
2332 root: self.root.clone(),
2333 state: self.state.clone(),
2334 }
2335 }
2336}
2337
2338impl Dynamic {
2339 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2340 state.lock().requests.add_handler();
2341
2342 Self { info, root, state }
2343 }
2344
2345 pub fn info(&self) -> &Origin {
2347 &self.info
2348 }
2349
2350 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2352 let mut state = ready!(self.state.poll(waiter, |state| {
2353 if state.requests.has_queued() {
2354 Poll::Ready(())
2355 } else {
2356 Poll::Pending
2357 }
2358 }));
2359
2360 let path = state.requests.pop().expect("predicate guaranteed a request");
2361 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2367 Poll::Ready(Ok(Request {
2368 path,
2369 producer,
2370 state: self.state.clone(),
2371 }))
2372 }
2373
2374 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2377 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2378 }
2379
2380 pub fn root(&self) -> &Path<'_> {
2382 &self.root
2383 }
2384}
2385
2386impl Drop for Dynamic {
2387 fn drop(&mut self) {
2388 let mut state = self.state.lock();
2391 if state.requests.remove_handler() {
2392 state.requests.drain_queued();
2396 }
2397 }
2398}
2399
2400pub struct Request {
2407 path: PathOwned,
2409
2410 producer: kio::Producer<PendingBroadcast>,
2413
2414 state: kio::Shared<OriginDynamicState>,
2416}
2417
2418impl Request {
2419 pub fn path(&self) -> &Path<'_> {
2421 &self.path
2422 }
2423
2424 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2430 let broadcast = broadcast.consume();
2431
2432 let resolved = {
2438 let mut state = self.state.lock();
2439 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2440 state
2441 .requests
2442 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2443 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2444 };
2445
2446 if let Ok(mut pending) = self.producer.write() {
2447 pending.resolved = Some(Ok(resolved));
2448 }
2449 }
2451
2452 pub fn reject(self, err: Error) {
2454 self.state
2455 .lock()
2456 .requests
2457 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2458 if let Ok(mut state) = self.producer.write() {
2459 state.resolved = Some(Err(err));
2460 }
2461 }
2462}
2463
2464impl Drop for Request {
2465 fn drop(&mut self) {
2466 self.state
2474 .lock()
2475 .requests
2476 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2477 }
2478}
2479
2480pub struct Requesting {
2487 inner: RequestState,
2488 stats: stats::Scope,
2491}
2492
2493enum RequestState {
2494 Ready(broadcast::Consumer),
2496 Failed(Error),
2499 Pending(kio::Consumer<PendingBroadcast>),
2501}
2502
2503impl Requesting {
2504 fn ready(broadcast: broadcast::Consumer) -> Self {
2505 Self {
2506 inner: RequestState::Ready(broadcast),
2507 stats: stats::Scope::default(),
2508 }
2509 }
2510
2511 fn failed(error: Error) -> Self {
2512 Self {
2513 inner: RequestState::Failed(error),
2514 stats: stats::Scope::default(),
2515 }
2516 }
2517
2518 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2519 Self {
2520 inner: RequestState::Pending(consumer),
2521 stats: stats::Scope::default(),
2522 }
2523 }
2524
2525 fn with_stats(mut self, scope: stats::Scope) -> Self {
2526 self.stats = scope;
2527 self
2528 }
2529
2530 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2532 match &self.inner {
2533 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2534 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2535 RequestState::Pending(consumer) => Poll::Ready(
2536 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2537 Some(result) => Poll::Ready(result.clone()),
2538 None => Poll::Pending,
2539 })) {
2540 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2541 Err(_closed) => Err(Error::Unroutable),
2543 },
2544 ),
2545 }
2546 }
2547}
2548
2549impl kio::Pollable for Requesting {
2550 type Output = Result<broadcast::Consumer, Error>;
2551
2552 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2553 self.poll_ok(waiter)
2554 }
2555}
2556
2557pub trait Consume<T> {
2565 fn consume(&self) -> T;
2567}
2568
2569impl<T, U: Consume<T>> Consume<T> for &U {
2570 fn consume(&self) -> T {
2571 (**self).consume()
2572 }
2573}
2574
2575impl Consume<Consumer> for Producer {
2576 fn consume(&self) -> Consumer {
2577 Consumer::new(
2581 self.info,
2582 self.root.clone(),
2583 self.nodes.clone(),
2584 self.dynamic.clone(),
2585 stats::Session::default(),
2586 )
2587 }
2588}
2589
2590impl Consume<Consumer> for Consumer {
2591 fn consume(&self) -> Consumer {
2592 self.clone()
2593 }
2594}
2595
2596impl Consume<broadcast::Consumer> for broadcast::Producer {
2597 fn consume(&self) -> broadcast::Consumer {
2598 self.consume()
2600 }
2601}
2602
2603impl Consume<broadcast::Consumer> for broadcast::Consumer {
2604 fn consume(&self) -> broadcast::Consumer {
2605 self.clone()
2606 }
2607}
2608
2609impl Consume<track::Consumer> for track::Producer {
2610 fn consume(&self) -> track::Consumer {
2611 self.consume()
2612 }
2613}
2614
2615impl Consume<track::Consumer> for track::Consumer {
2616 fn consume(&self) -> track::Consumer {
2617 self.clone()
2618 }
2619}
2620
2621#[derive(Clone)]
2627pub struct Consumer {
2628 info: Origin,
2630 nodes: OriginNodes,
2631
2632 root: PathOwned,
2634
2635 dynamic: kio::Shared<OriginDynamicState>,
2638
2639 stats: stats::Session,
2643
2644 exclude: Option<Origin>,
2648}
2649
2650impl std::ops::Deref for Consumer {
2651 type Target = Origin;
2652
2653 fn deref(&self) -> &Self::Target {
2654 &self.info
2655 }
2656}
2657
2658impl Consumer {
2659 fn new(
2660 info: Origin,
2661 root: PathOwned,
2662 nodes: OriginNodes,
2663 dynamic: kio::Shared<OriginDynamicState>,
2664 stats: stats::Session,
2665 ) -> Self {
2666 Self {
2667 info,
2668 nodes,
2669 root,
2670 dynamic,
2671 stats,
2672 exclude: None,
2673 }
2674 }
2675
2676 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2681 self.exclude = Some(peer);
2682 self
2683 }
2684
2685 pub fn with_stats(mut self, session: stats::Session) -> Self {
2689 self.stats = session;
2690 self
2691 }
2692
2693 fn untagged(&self) -> Self {
2697 Self {
2698 stats: stats::Session::default(),
2699 ..self.clone()
2700 }
2701 }
2702
2703 pub(crate) fn empty(&self) -> Self {
2708 Self {
2709 info: self.info,
2710 nodes: OriginNodes { nodes: Vec::new() },
2711 root: self.root.clone(),
2712 dynamic: self.dynamic.clone(),
2713 stats: self.stats.clone(),
2714 exclude: self.exclude,
2715 }
2716 }
2717
2718 pub fn announced(&self) -> AnnounceConsumer {
2725 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone(), self.exclude)
2726 }
2727
2728 pub fn consume(&self) -> Self {
2730 self.clone()
2731 }
2732
2733 fn resolve(&self, path: impl AsPath) -> Resolved {
2742 let path = path.as_path();
2743 let Some((root, rest)) = self.nodes.get(&path) else {
2744 return Resolved::Missing;
2745 };
2746 let state = root.lock();
2747 state.resolve_broadcast(&rest, self.exclude)
2748 }
2749
2750 #[cfg(test)]
2752 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2753 match self.resolve(path) {
2754 Resolved::Found(broadcast) => Some(broadcast),
2755 Resolved::Excluded | Resolved::Missing => None,
2756 }
2757 }
2758
2759 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2771 let path = path.as_path();
2772
2773 let consumer = self.scope(std::slice::from_ref(&path))?;
2775
2776 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2780 return None;
2781 }
2782
2783 let mut announced = consumer.untagged().announced();
2787 let scope = self.stats.egress(self.root.join(&path).to_owned());
2788 loop {
2789 let OriginAnnounce {
2790 path: announced_path,
2791 broadcast,
2792 } = announced.next().await?;
2793 if announced_path.as_path() == path
2795 && let Some(broadcast) = broadcast
2796 {
2797 return Some(broadcast.with_stats(scope));
2798 }
2799 }
2800 }
2801
2802 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2808 let prefixes = PathPrefixes::new(prefixes);
2809 Some(Consumer {
2810 info: self.info,
2811 root: self.root.clone(),
2812 nodes: self.nodes.select(&prefixes)?,
2813 dynamic: self.dynamic.clone(),
2814 stats: self.stats.clone(),
2815 exclude: self.exclude,
2816 })
2817 }
2818
2819 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2838 let path = path.as_path();
2839
2840 let absolute = self.root.join(&path).to_owned();
2844 let scope = self.stats.egress(&absolute);
2845
2846 match self.resolve(&path) {
2852 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
2853 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
2854 Resolved::Missing => {}
2855 }
2856
2857 let mut state = self.dynamic.lock();
2858
2859 if let Some(weak) = state.served.get(&absolute) {
2863 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2864 }
2865
2866 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2869 producer.consume()
2870 } else {
2871 let producer = kio::Producer::<PendingBroadcast>::default();
2872 let consumer = producer.consume();
2873 if state.requests.insert(absolute, producer).is_err() {
2874 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2875 }
2876 consumer
2877 };
2878
2879 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2880 }
2881
2882 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2887 let prefix = prefix.as_path();
2888
2889 Some(Self {
2890 info: self.info,
2891 root: self.root.join(&prefix).to_owned(),
2892 nodes: self.nodes.root(&prefix)?,
2893 dynamic: self.dynamic.clone(),
2894 stats: self.stats.clone(),
2895 exclude: self.exclude,
2896 })
2897 }
2898
2899 pub fn root(&self) -> &Path<'_> {
2901 &self.root
2902 }
2903
2904 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2907 self.nodes.nodes.iter().map(|(root, _)| root)
2908 }
2909
2910 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2912 self.root.join(path)
2913 }
2914}
2915
2916#[derive(Clone)]
2921pub struct AnnounceProducer {
2922 nodes: OriginNodes,
2923 root: PathOwned,
2924}
2925
2926impl AnnounceProducer {
2927 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2928 Self { nodes, root }
2929 }
2930
2931 pub fn consume(&self) -> AnnounceConsumer {
2937 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default(), None)
2940 }
2941
2942 pub fn root(&self) -> &Path<'_> {
2944 &self.root
2945 }
2946}
2947
2948pub struct AnnounceConsumer {
2953 id: ConsumerId,
2954 nodes: OriginNodes,
2955 root: PathOwned,
2956
2957 state: kio::Producer<OriginConsumerState>,
2960
2961 stats: stats::Session,
2964
2965 guards: HashMap<PathOwned, stats::Announce>,
2969}
2970
2971impl AnnounceConsumer {
2972 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session, exclude: Option<Origin>) -> Self {
2973 let state = kio::Producer::<OriginConsumerState>::default();
2974 let id = ConsumerId::new();
2975
2976 for (_, node) in &nodes.nodes {
2977 let notify = AnnounceConsumerNotify {
2978 root: root.clone(),
2979 state: state.clone(),
2980 exclude,
2981 };
2982 node.lock().consume(id, notify);
2983 }
2984
2985 Self {
2986 id,
2987 nodes,
2988 root,
2989 state,
2990 stats,
2991 guards: HashMap::new(),
2992 }
2993 }
2994
2995 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
3001 let OriginAnnounce { path, broadcast } = update;
3002 let absolute = self.root.join(&path).to_owned();
3003 match broadcast {
3004 Some(broadcast) => {
3005 let scope = self.stats.egress(&absolute);
3006 self.guards.entry(absolute).or_insert_with(|| scope.announce());
3007 OriginAnnounce {
3008 path,
3009 broadcast: Some(broadcast.with_stats(scope)),
3010 }
3011 }
3012 None => {
3013 self.guards.remove(&absolute);
3014 OriginAnnounce { path, broadcast: None }
3015 }
3016 }
3017 }
3018
3019 pub async fn next(&mut self) -> Option<OriginAnnounce> {
3026 kio::wait(|waiter| self.poll_next(waiter)).await
3027 }
3028
3029 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
3035 let update = {
3036 let mut state = match ready!(self.state.poll(waiter, |state| {
3037 if state.pending.is_empty() {
3038 Poll::Pending
3039 } else {
3040 Poll::Ready(())
3041 }
3042 })) {
3043 Ok(state) => state,
3044 Err(_) => return Poll::Ready(None),
3046 };
3047 state.take().expect("predicate guaranteed an update")
3048 };
3049 Poll::Ready(Some(self.attribute(update)))
3050 }
3051
3052 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
3057 let update = self.state.write().ok()?.take()?;
3058 Some(self.attribute(update))
3059 }
3060
3061 pub fn is_closed(&self) -> bool {
3063 self.state.write().is_err()
3064 }
3065
3066 pub fn root(&self) -> &Path<'_> {
3068 &self.root
3069 }
3070
3071 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
3073 self.root.join(path)
3074 }
3075}
3076
3077impl Drop for AnnounceConsumer {
3078 fn drop(&mut self) {
3079 for (_, root) in &self.nodes.nodes {
3080 root.lock().unconsume(self.id);
3081 }
3082 }
3083}
3084
3085#[cfg(test)]
3086use futures::FutureExt;
3087
3088#[cfg(test)]
3089#[allow(missing_docs)] impl AnnounceConsumer {
3091 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3092 let expected = expected.as_path();
3093 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3094 assert_eq!(announce.path, expected, "wrong path");
3095 let announced = announce.broadcast.expect("should be an active announce");
3096 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3097 }
3098
3099 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3103 let expected = expected.as_path();
3104 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3105 assert_eq!(announce.path, expected, "wrong path");
3106 announce.broadcast.expect("should be an active announce")
3107 }
3108
3109 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3110 let expected = expected.as_path();
3111 let announce = self.try_next().expect("no next");
3112 assert_eq!(announce.path, expected, "wrong path");
3113 let announced = announce.broadcast.expect("should be an active announce");
3114 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3115 }
3116
3117 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3119 let expected = expected.as_path();
3120 let announce = self.try_next().expect("no next");
3121 assert_eq!(announce.path, expected, "wrong path");
3122 announce.broadcast.expect("should be an active announce")
3123 }
3124
3125 pub fn assert_next_none(&mut self, expected: impl AsPath) {
3126 let expected = expected.as_path();
3127 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3128 assert_eq!(announce.path, expected, "wrong path");
3129 assert!(announce.broadcast.is_none(), "should be unannounced");
3130 }
3131
3132 pub fn assert_next_wait(&mut self) {
3133 if let Some(res) = self.next().now_or_never() {
3134 panic!("next should block: got {:?}", res.map(|a| a.path));
3135 }
3136 }
3137
3138 }
3147
3148#[cfg(test)]
3149mod tests {
3150 use crate::coding::Decode;
3151 use crate::group;
3152
3153 use super::*;
3154
3155 fn announce() -> broadcast::Route {
3157 broadcast::Route::new().with_announce(true)
3158 }
3159
3160 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
3166 let name = Path::new(name);
3167 let peer_key = fnv_key(&name, [peer]);
3168 (100u64..)
3169 .map(|id| Origin::new(id).unwrap())
3170 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
3171 .unwrap()
3172 }
3173
3174 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
3177 let source = broadcast::Info::new().produce().consume();
3178 FrontState {
3179 path: Path::new("test").to_owned(),
3180 self_origin,
3181 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
3182 next_route: routes.len() as u64,
3183 excluded: HashMap::new(),
3184 routes: routes
3185 .into_iter()
3186 .enumerate()
3187 .map(|(id, route)| FrontRoute {
3188 id: id as u64,
3189 route,
3190 source: source.clone(),
3191 })
3192 .collect(),
3193 active: Some(0),
3194 linger: Duration::ZERO,
3195 closed: false,
3196 }
3197 }
3198
3199 fn sibling_route(peer: Origin) -> broadcast::Route {
3202 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3203 announce().with_hops(hops)
3204 }
3205
3206 fn upstream_route(cost: u64) -> broadcast::Route {
3208 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3209 announce().with_hops(hops).with_cost(cost)
3210 }
3211
3212 #[test]
3216 fn test_carrying_gate_keys() {
3217 let peer = Origin::new(3).unwrap();
3218
3219 let mut lost = front_state(
3221 origin_keyed("test", peer, false),
3222 vec![upstream_route(10), sibling_route(peer)],
3223 );
3224 lost.reselect(true);
3225 assert_eq!(
3226 lost.active,
3227 Some(0),
3228 "carrying front re-parented onto a higher-keyed peer"
3229 );
3230 lost.reselect(false);
3231 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3232
3233 let mut won = front_state(
3235 origin_keyed("test", peer, true),
3236 vec![upstream_route(10), sibling_route(peer)],
3237 );
3238 won.reselect(true);
3239 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3240 }
3241
3242 #[test]
3247 fn test_carrying_gate_symmetric_race() {
3248 let a = Origin::new(1).unwrap();
3249 let b = Origin::new(2).unwrap();
3250
3251 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3252 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3253 a_view.reselect(true);
3254 b_view.reselect(true);
3255
3256 let a_moved = a_view.active == Some(1);
3257 let b_moved = b_view.active == Some(1);
3258 assert!(
3259 a_moved != b_moved,
3260 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3261 );
3262 }
3263
3264 #[test]
3269 fn test_carrying_switches_to_benign_routes() {
3270 let peer = Origin::new(3).unwrap();
3271 let lost = origin_keyed("test", peer, false);
3272
3273 let mut forwarder = sibling_route(peer).with_cost(4);
3275 forwarder.advertised = 4;
3276 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3277 state.reselect(true);
3278 assert_eq!(
3279 state.active,
3280 Some(1),
3281 "a cheaper forwarder path must win while carrying"
3282 );
3283
3284 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3286 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3287 state.reselect(true);
3288 assert_eq!(
3289 state.active,
3290 Some(1),
3291 "a direct publisher route must win while carrying"
3292 );
3293
3294 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3299 state.reselect(true);
3300 assert_eq!(
3301 state.active,
3302 Some(1),
3303 "a reconnect on an identical chain must win while carrying"
3304 );
3305 }
3306
3307 #[test]
3310 fn test_carrying_gate_ignores_unannounced_incumbent() {
3311 let peer = Origin::new(3).unwrap();
3312 let unannounced = upstream_route(10).with_announce(false);
3313 let mut state = front_state(
3314 origin_keyed("test", peer, false),
3315 vec![unannounced, sibling_route(peer)],
3316 );
3317 state.reselect(true);
3318 assert_eq!(
3319 state.active,
3320 Some(1),
3321 "an unannounced incumbent must always be displaced"
3322 );
3323 }
3324
3325 #[test]
3330 fn test_reflection_through_an_exposed_peer_cannot_take_over() {
3331 let peer = Origin::new(42).unwrap();
3332 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3333 let reflected = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3334
3335 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3336
3337 assert!(!state.taints_a_reader(&reflected));
3339
3340 *state.excluded.entry(peer).or_default() += 1;
3343 assert!(state.taints_a_reader(&reflected));
3344 }
3345
3346 #[test]
3349 fn test_rival_publisher_through_another_peer_still_takes_over() {
3350 let peer = Origin::new(42).unwrap();
3351 let elsewhere = Origin::new(43).unwrap();
3352 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3353 let rival = announce().with_hops(OriginList::try_from(vec![Origin::UNKNOWN, elsewhere]).unwrap());
3354
3355 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3356 *state.excluded.entry(peer).or_default() += 1;
3357
3358 assert!(!state.taints_a_reader(&rival), "only the peer we feed is a reflection");
3359 }
3360
3361 #[test]
3366 fn test_opaque_peer_understates_its_depth() {
3367 let us = Origin::new(1).unwrap();
3368
3369 let direct = || {
3371 announce()
3372 .with_hops(OriginList::try_from(vec![Origin::new(7).unwrap(), Origin::new(8).unwrap()]).unwrap())
3373 .with_cost(2)
3374 };
3375 let opaque = |cost| {
3378 announce()
3379 .with_hops(OriginList::try_from(vec![Origin::new(42).unwrap()]).unwrap())
3380 .with_cost(cost)
3381 };
3382
3383 let mut state = front_state(us, vec![direct(), opaque(1)]);
3385 state.reselect(false);
3386 assert_eq!(
3387 state.active,
3388 Some(1),
3389 "an unpriced opaque link out-ranks a shorter real path"
3390 );
3391
3392 let mut state = front_state(us, vec![direct(), opaque(16)]);
3394 state.reselect(false);
3395 assert_eq!(
3396 state.active,
3397 Some(0),
3398 "pricing the opaque link restores the intended order"
3399 );
3400 }
3401
3402 async fn settle() {
3405 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3406 }
3407
3408 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3411 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3412 .await
3413 .expect("timed out waiting for a track request")
3414 .expect("source closed");
3415 assert_eq!(request.name(), name, "unexpected track dispatched");
3416 request.accept(None)
3417 }
3418
3419 async fn accept_tracks(dynamic: &mut broadcast::Dynamic, count: usize) -> HashMap<String, track::Producer> {
3423 let mut accepted = HashMap::new();
3424 for _ in 0..count {
3425 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3426 .await
3427 .expect("timed out waiting for a track request")
3428 .expect("source closed");
3429 let name = request.name().to_string();
3430 accepted.insert(name, request.accept(None));
3431 }
3432 accepted
3433 }
3434
3435 #[tokio::test]
3439 async fn test_stats_tagged_end_to_end() {
3440 use crate::Timestamp;
3441 use crate::stats::{Config, Registry, Tier};
3442 use bytes::Bytes;
3443
3444 tokio::time::pause();
3445
3446 let registry = Registry::new(Config::new());
3447 let ctx = registry.tier(Tier::default()).session("acme");
3448
3449 let origin = Origin::random().produce();
3450 let ingress = origin.clone().with_stats(ctx.clone());
3451 let egress = origin.consume().with_stats(ctx.clone());
3452
3453 let mut announced = egress.announced();
3456
3457 let source = ingress.create_broadcast("demo", announce()).unwrap();
3459 let mut dynamic = source.dynamic();
3460 settle().await;
3461 settle().await;
3462
3463 let update = announced.next().await.unwrap();
3465 assert_eq!(update.path.as_str(), "demo");
3466 let broadcast = update.broadcast.unwrap();
3467
3468 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3470 let mut producer = accept_track(&mut dynamic, "video").await;
3471 settle().await;
3472 let mut sub = subscribing.await.unwrap();
3473
3474 let mut group = producer.append_group().unwrap();
3476 group
3477 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3478 .unwrap();
3479 group
3480 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3481 .unwrap();
3482 group.finish().unwrap();
3483
3484 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3486 let mut frames = 0;
3487 while let Some(frame) = group_c.read_frame().await.unwrap() {
3488 assert_eq!(frame.payload.len(), 5);
3489 frames += 1;
3490 }
3491 assert_eq!(frames, 2);
3492 settle().await;
3493
3494 let report = registry.report();
3495 let entry = report
3496 .traffic
3497 .iter()
3498 .find(|e| e.path.as_str() == "demo")
3499 .expect("demo tracked");
3500 let path_len = "demo".len() as u64;
3501
3502 let egress = &entry.publisher;
3504 assert_eq!(egress.announced, 1, "one egress announce");
3505 assert_eq!(egress.announced_bytes, path_len);
3506 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3507 assert_eq!(egress.broadcasts, 1, "one viewer");
3508 assert_eq!(egress.groups, 1);
3509 assert_eq!(egress.frames, 2);
3510 assert_eq!(egress.bytes, 10);
3511 assert_eq!(egress.fetches, 0);
3512
3513 let ingress = &entry.subscriber;
3515 assert_eq!(ingress.announced, 1, "one ingress announce");
3516 assert_eq!(ingress.announced_bytes, path_len);
3517 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3518 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3519 assert_eq!(ingress.groups, 1);
3520 assert_eq!(ingress.frames, 2);
3521 assert_eq!(ingress.bytes, 10);
3522
3523 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3525 let _ = fetched;
3526 settle().await;
3527 let report = registry.report();
3528 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3529 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3530 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3531 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3532 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3536 }
3537
3538 #[tokio::test]
3543 async fn test_stats_read_frame_counts_once() {
3544 use crate::Timestamp;
3545 use crate::stats::{Config, Registry, Tier};
3546 use bytes::Bytes;
3547
3548 tokio::time::pause();
3549
3550 let registry = Registry::new(Config::new());
3551 let ctx = registry.tier(Tier::default()).session("acme");
3552
3553 let origin = Origin::random().produce();
3554 let ingress = origin.clone().with_stats(ctx.clone());
3555 let egress = origin.consume().with_stats(ctx.clone());
3556
3557 let mut announced = egress.announced();
3558 let source = ingress.create_broadcast("demo", announce()).unwrap();
3559 let mut dynamic = source.dynamic();
3560 settle().await;
3561 settle().await;
3562
3563 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3564 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3565 let mut producer = accept_track(&mut dynamic, "video").await;
3566 settle().await;
3567 let mut sub = subscribing.await.unwrap();
3568
3569 producer
3571 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3572 .unwrap();
3573
3574 let frame = sub.read_frame().await.unwrap().expect("frame");
3575 assert_eq!(frame.payload.len(), 5);
3576 settle().await;
3577
3578 let report = registry.report();
3579 let entry = report
3580 .traffic
3581 .iter()
3582 .find(|e| e.path.as_str() == "demo")
3583 .expect("demo tracked");
3584 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3585 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3586 assert_eq!(
3587 entry.publisher.bytes, 5,
3588 "payload counted once, not zero and not doubled"
3589 );
3590 }
3591
3592 #[tokio::test]
3596 async fn test_stats_datagrams_counted_both_sides() {
3597 use crate::Timestamp;
3598 use crate::stats::{Config, Registry, Tier};
3599
3600 tokio::time::pause();
3601
3602 let registry = Registry::new(Config::new());
3603 let ctx = registry.tier(Tier::default()).session("acme");
3604
3605 let origin = Origin::random().produce();
3606 let ingress = origin.clone().with_stats(ctx.clone());
3607 let egress = origin.consume().with_stats(ctx.clone());
3608
3609 let mut announced = egress.announced();
3610 let source = ingress.create_broadcast("demo", announce()).unwrap();
3611 let mut dynamic = source.dynamic();
3612 settle().await;
3613 settle().await;
3614
3615 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3616 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3617 let mut producer = accept_track(&mut dynamic, "video").await;
3618 settle().await;
3619 let mut sub = subscribing.await.unwrap();
3620
3621 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3622 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3623 assert_eq!(&datagram.payload[..], b"hello");
3624 settle().await;
3625
3626 let report = registry.report();
3627 let entry = report
3628 .traffic
3629 .iter()
3630 .find(|e| e.path.as_str() == "demo")
3631 .expect("demo tracked");
3632
3633 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3634 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3635 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3636 assert_eq!(traffic.frames, 1, "{side}: one frame");
3637 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3638 }
3639 }
3640
3641 #[test]
3642 fn origin_rejects_reserved_ids() {
3643 assert!(Origin::new(0).is_err());
3644 assert!(Origin::new(1u64 << 62).is_err());
3645 assert_eq!(Origin::new(1).unwrap().id(), 1);
3646
3647 let mut zero = [0u8].as_slice();
3648 assert_eq!(
3649 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3650 Origin::UNKNOWN
3651 );
3652 }
3653
3654 #[test]
3655 fn origin_list_push_fails_at_limit() {
3656 let mut list = OriginList::new();
3657 for _ in 0..MAX_HOPS {
3658 list.push(Origin::random()).unwrap();
3659 }
3660 assert_eq!(list.len(), MAX_HOPS);
3661 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3662 }
3663
3664 #[test]
3665 fn origin_list_replace_first() {
3666 let mut list = OriginList::new();
3667 for _ in 0..3 {
3668 list.push(Origin::UNKNOWN).unwrap();
3669 }
3670
3671 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3673 assert_eq!(
3674 list.as_slice(),
3675 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3676 );
3677
3678 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3680 assert_eq!(list.len(), 3);
3681 }
3682
3683 #[test]
3684 fn origin_list_try_from_vec_enforces_limit() {
3685 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3686 assert!(OriginList::try_from(under).is_ok());
3687
3688 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3689 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3690 }
3691
3692 #[tokio::test]
3693 async fn test_announce() {
3694 tokio::time::pause();
3695
3696 let origin = Origin::random().produce();
3697
3698 let mut consumer1 = origin.consume().announced();
3699 consumer1.assert_next_wait();
3700
3701 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3703 settle().await;
3704
3705 consumer1.assert_next_some("test1");
3706 consumer1.assert_next_wait();
3707
3708 let mut consumer2 = origin.consume().announced();
3711
3712 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3714 settle().await;
3715
3716 consumer1.assert_next_some("test2");
3717 consumer1.assert_next_wait();
3718
3719 consumer2.assert_next_some("test1");
3720 consumer2.assert_next_some("test2");
3721 consumer2.assert_next_wait();
3722
3723 broadcast1.finish();
3725 settle().await;
3726
3727 consumer1.assert_next_none("test1");
3729 consumer2.assert_next_none("test1");
3730 consumer1.assert_next_wait();
3731 consumer2.assert_next_wait();
3732
3733 let mut consumer3 = origin.consume().announced();
3735 consumer3.assert_next_some("test2");
3736 consumer3.assert_next_wait();
3737
3738 broadcast2.finish();
3739 settle().await;
3740
3741 consumer1.assert_next_none("test2");
3742 consumer2.assert_next_none("test2");
3743 consumer3.assert_next_none("test2");
3744 }
3745
3746 #[tokio::test]
3750 async fn test_duplicate() {
3751 tokio::time::pause();
3752
3753 let origin = Origin::random().produce();
3754 let consumer = origin.consume();
3755 let mut announced = consumer.announced();
3756
3757 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3758 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3759 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3760 settle().await;
3761 assert!(consumer.get_broadcast("test").is_some());
3762
3763 announced.assert_next_some("test");
3764 announced.assert_next_wait();
3765
3766 broadcast2.finish();
3768 settle().await;
3769 assert!(consumer.get_broadcast("test").is_some());
3770 announced.assert_next_wait();
3771
3772 broadcast1.finish();
3774 settle().await;
3775 assert!(consumer.get_broadcast("test").is_some());
3776 announced.assert_next_wait();
3777
3778 broadcast3.finish();
3780 settle().await;
3781 assert!(consumer.get_broadcast("test").is_none());
3782
3783 announced.assert_next_none("test");
3784 announced.assert_next_wait();
3785 }
3786
3787 #[tokio::test]
3790 async fn test_route_failover() {
3791 tokio::time::pause();
3792
3793 let origin = Origin::random().produce();
3794 let consumer = origin.consume();
3795 let mut announced = consumer.announced();
3796
3797 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3800 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3801
3802 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3804 let mut dynamic_a = source_a.dynamic();
3805 settle().await;
3806 settle().await;
3807 let broadcast = consumer.request_broadcast("test").await.unwrap();
3808 announced.assert_next_some("test");
3809
3810 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3812 let mut dynamic_b = source_b.dynamic();
3813 settle().await;
3814 settle().await;
3815 announced.assert_next_wait();
3816
3817 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3819 let mut producer = accept_track(&mut dynamic_a, "video").await;
3820 settle().await;
3821 dynamic_b.assert_no_request();
3822
3823 let mut sub = subscribing.await.unwrap();
3824 sub.assert_no_group();
3827 assert_eq!(producer.subscription().unwrap().group_start, None);
3828
3829 producer.append_group().unwrap();
3830 producer.append_group().unwrap();
3831 assert_eq!(sub.assert_group().sequence, 0);
3832 assert_eq!(sub.assert_group().sequence, 1);
3833
3834 producer.abort(Error::Dropped).unwrap();
3838 source_a.abort(Error::Dropped).unwrap();
3839 drop(dynamic_a);
3840 settle().await;
3841 announced.assert_next_wait();
3842
3843 let mut producer = accept_track(&mut dynamic_b, "video").await;
3847 settle().await;
3848 sub.assert_no_group();
3849 assert_eq!(producer.subscription().unwrap().group_start, None);
3850 producer.create_group(group::Info { sequence: 1 }).unwrap();
3851 producer.create_group(group::Info { sequence: 2 }).unwrap();
3852 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3853 sub.assert_not_closed();
3854 }
3855
3856 #[tokio::test]
3863 async fn test_route_failover_restores_every_track() {
3864 tokio::time::pause();
3865
3866 let origin = Origin::random().produce();
3867 let consumer = origin.consume();
3868
3869 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3871 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3872
3873 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3874 let mut dynamic_a = source_a.dynamic();
3875 settle().await;
3876 settle().await;
3877 let broadcast = consumer.request_broadcast("test").await.unwrap();
3878
3879 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3881 let mut dynamic_b = source_b.dynamic();
3882 settle().await;
3883 settle().await;
3884
3885 const TRACKS: [(&str, u64); 3] = [("video", 3), ("audio", 1), ("data", 2)];
3889
3890 let subscribing: Vec<_> = TRACKS
3891 .iter()
3892 .map(|(name, _)| broadcast.track(name).unwrap().subscribe(None))
3893 .collect();
3894 let mut producers_a = accept_tracks(&mut dynamic_a, TRACKS.len()).await;
3895 settle().await;
3896
3897 let mut subs = Vec::new();
3898 for ((name, groups), subscribing) in TRACKS.iter().zip(subscribing) {
3899 let mut sub = subscribing.await.unwrap();
3900 let producer = producers_a
3901 .get_mut(*name)
3902 .unwrap_or_else(|| panic!("{name} was never dispatched"));
3903 for expected in 0..*groups {
3904 producer.append_group().unwrap();
3905 assert_eq!(sub.assert_group().sequence, expected, "{name} did not start");
3906 }
3907 subs.push((*name, *groups, sub));
3908 }
3909
3910 for (_, producer) in producers_a.drain() {
3912 producer.abort(Error::Dropped).unwrap();
3913 }
3914 source_a.abort(Error::Dropped).unwrap();
3915 drop(dynamic_a);
3916 settle().await;
3917
3918 let mut producers_b = accept_tracks(&mut dynamic_b, TRACKS.len()).await;
3921 settle().await;
3922
3923 for (_, _, sub) in subs.iter_mut() {
3926 sub.assert_no_group();
3927 }
3928 settle().await;
3929
3930 for (name, groups, sub) in subs.iter_mut() {
3931 let producer = producers_b
3932 .get_mut(*name)
3933 .unwrap_or_else(|| panic!("{name} was never re-dispatched to the standby"));
3934 assert_eq!(
3937 producer
3938 .subscription()
3939 .unwrap_or_else(|| panic!("{name} resumed without a subscription"))
3940 .group_start,
3941 None,
3942 "{name} must keep the subscriber's live-edge demand"
3943 );
3944 let boundary = *groups;
3945
3946 producer.create_group(group::Info { sequence: boundary - 1 }).unwrap();
3948 producer.create_group(group::Info { sequence: boundary }).unwrap();
3949 assert_eq!(sub.assert_group().sequence, boundary, "{name} did not resume");
3950 sub.assert_not_closed();
3951 }
3952 }
3953
3954 #[tokio::test]
3957 async fn test_broadcast_route_watch() {
3958 let mut producer = broadcast::Info::new().produce();
3959 let mut consumer = producer.consume();
3960
3961 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3963
3964 producer.set_route(broadcast::Route::default()).unwrap();
3966 assert!(consumer.route_changed().now_or_never().is_none());
3967
3968 let mut hops = OriginList::new();
3969 hops.push(Origin::new(7).unwrap()).unwrap();
3970 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3971 producer.set_route(route.clone()).unwrap();
3972 assert_eq!(consumer.route_changed().await.unwrap(), route);
3973
3974 let mut fresh = producer.consume();
3976 assert_eq!(fresh.route_changed().await.unwrap(), route);
3977
3978 drop(producer);
3979 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3980 }
3981
3982 #[tokio::test]
3986 async fn test_route_cost_update() {
3987 tokio::time::pause();
3988
3989 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3993 let consumer = origin.consume();
3994 let mut announced = consumer.announced();
3995
3996 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3999 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4000
4001 let mut source_a = origin
4003 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4004 .unwrap();
4005 let mut dynamic_a = source_a.dynamic();
4006 settle().await;
4007 let broadcast = consumer.request_broadcast("test").await.unwrap();
4008 announced.assert_next_some("test");
4009
4010 let mut watch = broadcast.clone();
4011 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
4012
4013 let mut source_b = origin
4014 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4015 .unwrap();
4016 let mut dynamic_b = source_b.dynamic();
4017 settle().await;
4018 assert!(
4019 watch.route_changed().now_or_never().is_none(),
4020 "a losing standby must not change the advertised route"
4021 );
4022
4023 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4025 let mut producer = accept_track(&mut dynamic_a, "video").await;
4026 settle().await;
4027 let mut sub = subscribing.await.unwrap();
4028 producer.append_group().unwrap();
4029 assert_eq!(sub.assert_group().sequence, 0);
4030
4031 source_a
4034 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
4035 .unwrap();
4036 settle().await;
4037 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
4038 announced.assert_next_wait();
4039
4040 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
4041 settle().await;
4042 sub.assert_no_group();
4045 assert_eq!(producer_b.subscription().unwrap().group_start, None);
4046 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
4047 assert_eq!(sub.assert_group().sequence, 1);
4048 sub.assert_not_closed();
4049
4050 source_b
4052 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
4053 .unwrap();
4054 settle().await;
4055 let advertised = watch.route_changed().await.unwrap();
4056 assert_eq!(advertised.hops, hops_b);
4057 assert_eq!(advertised.cost, 5);
4058 announced.assert_next_wait();
4059 }
4060
4061 #[tokio::test]
4064 async fn test_completed_track_survives_route_churn() {
4065 tokio::time::pause();
4066
4067 let origin = Origin::random().produce();
4068 let consumer = origin.consume();
4069
4070 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4072 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4073
4074 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4075 let mut dynamic_a = source_a.dynamic();
4076 settle().await;
4077 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4078 let mut dynamic_b = source_b.dynamic();
4079 settle().await;
4080 settle().await;
4081 let broadcast = consumer.request_broadcast("test").await.unwrap();
4082
4083 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4085 let mut producer = accept_track(&mut dynamic_a, "video").await;
4086 settle().await;
4087 let mut sub = subscribing.await.unwrap();
4088 producer.append_group().unwrap();
4089 assert_eq!(sub.assert_group().sequence, 0);
4090 producer.finish().unwrap();
4091 drop(producer);
4092 settle().await;
4093 sub.assert_closed();
4094
4095 source_a.abort(Error::Dropped).unwrap();
4097 drop(dynamic_a);
4098 settle().await;
4099 dynamic_b.assert_no_request();
4100
4101 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
4103 late.assert_closed();
4104 }
4105
4106 #[tokio::test]
4110 async fn test_refused_track_aborts_instantly() {
4111 tokio::time::pause();
4112
4113 let origin = Origin::random().produce();
4114 let consumer = origin.consume();
4115
4116 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4117 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4118 let mut dynamic = source.dynamic();
4119 settle().await;
4120 settle().await;
4121 let broadcast = consumer.request_broadcast("test").await.unwrap();
4122
4123 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4124 let request = dynamic.requested_track().await.unwrap();
4125 request.reject(Error::NotFound);
4126 settle().await;
4127
4128 assert!(matches!(subscribing.await, Err(Error::NotFound)));
4130 dynamic.assert_no_request();
4131 }
4132
4133 #[tokio::test]
4138 async fn test_stale_rejection_does_not_abort_a_handover() {
4139 tokio::time::pause();
4140
4141 let origin = Origin::random().produce();
4142 let consumer = origin.consume();
4143
4144 let publisher = Origin::new(1).unwrap();
4145 let peer = Origin::new(5).unwrap();
4146 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4147 let local = OriginList::try_from(vec![publisher]).unwrap();
4148
4149 let source_remote = origin
4151 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4152 .unwrap();
4153 let mut dynamic_remote = source_remote.dynamic();
4154 settle().await;
4155 settle().await;
4156 let broadcast = consumer.request_broadcast("test").await.unwrap();
4157 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4158 let request_remote = dynamic_remote.requested_track().await.unwrap();
4159
4160 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4163 let mut dynamic_local = source_local.dynamic();
4164 request_remote.reject(Error::NotFound);
4165 settle().await;
4166
4167 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4169 settle().await;
4170 let mut sub = subscribing
4171 .await
4172 .expect("the handover must win over the stale rejection");
4173 producer_local.append_group().unwrap();
4174 assert_eq!(sub.assert_group().sequence, 0);
4175 sub.assert_not_closed();
4176 }
4177
4178 #[tokio::test]
4182 async fn test_route_handover() {
4183 tokio::time::pause();
4184
4185 let origin = Origin::random().produce();
4186 let consumer = origin.consume();
4187 let mut announced = consumer.announced();
4188
4189 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4191 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4192
4193 let source_a = origin
4194 .create_broadcast("test", announce().with_hops(hops_long))
4195 .unwrap();
4196 let mut dynamic_a = source_a.dynamic();
4197 settle().await;
4198 settle().await;
4199 let broadcast = consumer.request_broadcast("test").await.unwrap();
4200 announced.assert_next_some("test");
4201
4202 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4203 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
4204 settle().await;
4205 let mut sub = subscribing.await.unwrap();
4206 producer_a.append_group().unwrap();
4207 producer_a.append_group().unwrap();
4208 assert_eq!(sub.assert_group().sequence, 0);
4209 assert_eq!(sub.assert_group().sequence, 1);
4210
4211 let source_b = origin
4214 .create_broadcast("test", announce().with_hops(hops_short))
4215 .unwrap();
4216 let mut dynamic_b = source_b.dynamic();
4217 settle().await;
4218 settle().await;
4219 announced.assert_next_wait();
4220
4221 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
4222 settle().await;
4223
4224 sub.assert_no_group();
4228 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
4229 assert_eq!(producer_b.subscription().unwrap().group_start, None);
4230
4231 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
4233 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
4234 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
4235 assert_eq!(sub.assert_group().sequence, 2);
4236 assert_eq!(sub.assert_group().sequence, 3);
4237 sub.assert_no_group();
4238 sub.assert_not_closed();
4239 }
4240
4241 #[tokio::test(start_paused = true)]
4244 async fn test_route_unannounce_immediate() {
4245 let origin = Origin::random().produce();
4246 let consumer = origin.consume();
4247 let mut announced = consumer.announced();
4248
4249 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4250 let mut source = origin
4251 .create_broadcast("test", announce().with_hops(hops.clone()))
4252 .unwrap();
4253 settle().await;
4254 let broadcast = consumer.request_broadcast("test").await.unwrap();
4255 announced.assert_next_some("test");
4256
4257 source.finish();
4260 settle().await;
4261 announced.assert_next_none("test");
4262
4263 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4265 settle().await;
4266 let fresh = consumer.request_broadcast("test").await.unwrap();
4267 announced.assert_next_some("test");
4268 assert!(
4269 !fresh.is_clone(&broadcast),
4270 "re-create must not splice the old broadcast"
4271 );
4272 }
4273
4274 #[tokio::test(start_paused = true)]
4279 async fn test_route_detach_immediate() {
4280 let origin = Origin::random().produce();
4281 let consumer = origin.consume();
4282 let mut announced = consumer.announced();
4283
4284 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4285 let source = origin
4286 .create_broadcast("test", announce().with_hops(hops.clone()))
4287 .unwrap();
4288 let mut dynamic = source.dynamic();
4289 settle().await;
4290 settle().await;
4291 let broadcast = consumer.request_broadcast("test").await.unwrap();
4292 announced.assert_next_some("test");
4293
4294 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4295 let producer = accept_track(&mut dynamic, "video").await;
4296 settle().await;
4297 let mut sub = subscribing.await.unwrap();
4298
4299 drop(producer);
4301 source.abort(Error::Dropped).unwrap();
4302 drop(dynamic);
4303
4304 settle().await;
4305 announced.assert_next_none("test");
4306 sub.assert_error();
4307
4308 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4311 settle().await;
4312 settle().await;
4313 let fresh = consumer.request_broadcast("test").await.unwrap();
4314 announced.assert_next_some("test");
4315 assert!(
4316 !fresh.is_clone(&broadcast),
4317 "re-create must not splice the old broadcast"
4318 );
4319 }
4320
4321 #[tokio::test(start_paused = true)]
4326 async fn test_idle_track_releases_without_respinning() {
4327 let origin = Info::new(Origin::random()).produce();
4328 let consumer = origin.consume();
4329
4330 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4331 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4332 let mut dynamic = source.dynamic();
4333 settle().await;
4334 let broadcast = consumer.request_broadcast("test").await.unwrap();
4335
4336 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4337 let producer = accept_track(&mut dynamic, "video").await;
4338 settle().await;
4339 let sub = subscribing.await.unwrap();
4340
4341 drop(sub);
4344 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
4345 settle().await;
4346 assert!(
4347 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
4348 "the copy must stay spliced inside the linger",
4349 );
4350
4351 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4354 settle().await;
4355 assert!(
4356 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4357 "an idle copy must be released after the linger",
4358 );
4359
4360 for _ in 0..3 {
4364 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4365 settle().await;
4366 assert!(
4367 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4368 "an unread copy must stay released, not be re-spliced",
4369 );
4370 }
4371 assert!(
4372 dynamic.requested_track().now_or_never().is_none(),
4373 "an unread track must not be re-requested",
4374 );
4375 drop(producer);
4376
4377 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4379 let mut producer = accept_track(&mut dynamic, "video").await;
4380 settle().await;
4381 let mut sub = subscribing.await.unwrap();
4382 producer.append_group().unwrap();
4383 assert_eq!(sub.assert_group().sequence, 0);
4384 }
4385
4386 #[tokio::test(start_paused = true)]
4390 async fn test_back_to_back_fetches_reuse_the_track() {
4391 let origin = Info::new(Origin::random()).produce();
4392 let consumer = origin.consume();
4393
4394 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4395 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4396 let mut dynamic = source.dynamic();
4397 settle().await;
4398 let broadcast = consumer.request_broadcast("test").await.unwrap();
4399
4400 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4402 let mut producer = accept_track(&mut dynamic, "video").await;
4403 producer.append_group().unwrap().finish().unwrap();
4404 settle().await;
4405 let first = fetching.await.expect("first fetch");
4406 drop(first);
4407
4408 settle().await;
4410 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4411 settle().await;
4412 assert!(
4413 dynamic.requested_track().now_or_never().is_none(),
4414 "a fetch inside the linger must reuse the track, not re-request it",
4415 );
4416 drop(fetching.await.expect("second fetch"));
4417
4418 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4420 settle().await;
4421 assert!(
4422 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4423 "the copy must be released once the fetches stop",
4424 );
4425 drop(producer);
4426
4427 settle().await;
4429 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4430 let mut producer = accept_track(&mut dynamic, "video").await;
4431 producer.append_group().unwrap().finish().unwrap();
4432 settle().await;
4433 fetching.await.expect("fetch after the linger");
4434 }
4435
4436 #[tokio::test(start_paused = true)]
4440 async fn test_linger_reconnect_splices() {
4441 let origin = Info::new(Origin::random())
4442 .with_linger(Duration::from_secs(5))
4443 .produce();
4444 let consumer = origin.consume();
4445 let mut announced = consumer.announced();
4446
4447 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4448 let source = origin
4449 .create_broadcast("test", announce().with_hops(hops.clone()))
4450 .unwrap();
4451 let mut dynamic = source.dynamic();
4452 settle().await;
4453 settle().await;
4454 let broadcast = consumer.request_broadcast("test").await.unwrap();
4455 announced.assert_next_some("test");
4456
4457 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4458 let mut producer = accept_track(&mut dynamic, "video").await;
4459 settle().await;
4460 let mut sub = subscribing.await.unwrap();
4461
4462 producer.append_group().unwrap();
4463 producer.append_group().unwrap();
4464 assert_eq!(sub.assert_group().sequence, 0);
4465 assert_eq!(sub.assert_group().sequence, 1);
4466
4467 drop(producer);
4470 source.abort(Error::Dropped).unwrap();
4471 drop(dynamic);
4472 settle().await;
4473
4474 announced.assert_next_wait();
4476 sub.assert_no_group();
4477 sub.assert_not_closed();
4478
4479 let during = consumer.request_broadcast("test").await.unwrap();
4481 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4482
4483 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4486 let mut dynamic = source.dynamic();
4487 settle().await;
4488 settle().await;
4489 announced.assert_next_wait();
4490 let again = consumer.request_broadcast("test").await.unwrap();
4491 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4492
4493 let mut producer = accept_track(&mut dynamic, "video").await;
4497 settle().await;
4498 sub.assert_no_group();
4499 assert_eq!(producer.subscription().unwrap().group_start, None);
4500 producer.create_group(group::Info { sequence: 2 }).unwrap();
4501 assert_eq!(sub.assert_group().sequence, 2);
4502 sub.assert_not_closed();
4503 }
4504
4505 #[tokio::test(start_paused = true)]
4508 async fn test_linger_expiry_closes() {
4509 let origin = Info::new(Origin::random())
4510 .with_linger(Duration::from_secs(5))
4511 .produce();
4512 let consumer = origin.consume();
4513 let mut announced = consumer.announced();
4514
4515 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4516 let source = origin
4517 .create_broadcast("test", announce().with_hops(hops.clone()))
4518 .unwrap();
4519 let mut dynamic = source.dynamic();
4520 settle().await;
4521 settle().await;
4522 let broadcast = consumer.request_broadcast("test").await.unwrap();
4523 announced.assert_next_some("test");
4524
4525 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4526 let producer = accept_track(&mut dynamic, "video").await;
4527 settle().await;
4528 let mut sub = subscribing.await.unwrap();
4529
4530 drop(producer);
4531 source.abort(Error::Dropped).unwrap();
4532 drop(dynamic);
4533 settle().await;
4534 announced.assert_next_wait();
4535
4536 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4538 settle().await;
4539 announced.assert_next_none("test");
4540 sub.assert_error();
4541
4542 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4544 settle().await;
4545 settle().await;
4546 let fresh = consumer.request_broadcast("test").await.unwrap();
4547 announced.assert_next_some("test");
4548 assert!(
4549 !fresh.is_clone(&broadcast),
4550 "a late re-create must not splice the expired broadcast"
4551 );
4552 }
4553
4554 #[tokio::test(start_paused = true)]
4558 async fn test_linger_forever() {
4559 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4560 let consumer = origin.consume();
4561 let mut announced = consumer.announced();
4562
4563 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4564 let source = origin
4565 .create_broadcast("test", announce().with_hops(hops.clone()))
4566 .unwrap();
4567 settle().await;
4568 let broadcast = consumer.request_broadcast("test").await.unwrap();
4569 announced.assert_next_some("test");
4570
4571 source.abort(Error::Dropped).unwrap();
4572 settle().await;
4573
4574 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4576 announced.assert_next_wait();
4577 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4578 settle().await;
4579 settle().await;
4580 let again = consumer.request_broadcast("test").await.unwrap();
4581 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4582 drop(source);
4583 }
4584
4585 #[tokio::test(start_paused = true)]
4599 async fn test_linger_parks_a_live_subscription() {
4600 let origin = Info::new(Origin::random())
4601 .with_linger(Duration::from_secs(5))
4602 .produce();
4603 let consumer = origin.consume();
4604
4605 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4606 let source = origin
4607 .create_broadcast("test", announce().with_hops(hops.clone()))
4608 .unwrap();
4609 let mut dynamic = source.dynamic();
4610 settle().await;
4611 settle().await;
4612 let broadcast = consumer.request_broadcast("test").await.unwrap();
4613
4614 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4617 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4618 settle().await;
4619 let mut sub = subscribing.await.unwrap();
4620 producer.append_group().unwrap();
4621 assert_eq!(sub.assert_group().sequence, 0);
4622
4623 source.abort(Error::Dropped).unwrap();
4628 settle().await;
4629 settle().await;
4630 sub.assert_not_closed();
4631
4632 tokio::time::sleep(Duration::from_secs(4)).await;
4635 settle().await;
4636 sub.assert_not_closed();
4637
4638 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4640 let mut dynamic = source.dynamic();
4641 settle().await;
4642 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4643 settle().await;
4644 producer.create_group(group::Info { sequence: 1 }).unwrap();
4645 assert_eq!(sub.assert_group().sequence, 1);
4646 sub.assert_not_closed();
4647 }
4648
4649 #[tokio::test(start_paused = true)]
4660 async fn test_idle_release_survives_the_route_leaving() {
4661 let origin = Info::new(Origin::random())
4662 .with_linger(Duration::from_secs(600))
4663 .produce();
4664 let consumer = origin.consume();
4665
4666 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4667 let source = origin
4668 .create_broadcast("test", announce().with_hops(hops.clone()))
4669 .unwrap();
4670 let mut dynamic = source.dynamic();
4671 settle().await;
4672 settle().await;
4673 let broadcast = consumer.request_broadcast("test").await.unwrap();
4674
4675 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4676 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4677 settle().await;
4678 let mut sub = subscribing.await.unwrap();
4679 producer.append_group().unwrap();
4680 assert_eq!(sub.assert_group().sequence, 0);
4681
4682 source.abort(Error::Dropped).unwrap();
4685 settle().await;
4686 settle().await;
4687 drop(sub);
4688 drop(producer);
4689 drop(dynamic);
4690 settle().await;
4691 tokio::time::sleep(TRACK_IDLE_LINGER + Duration::from_secs(1)).await;
4692 settle().await;
4693
4694 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4697 let mut dynamic = source.dynamic();
4698 settle().await;
4699 settle().await;
4700 let broadcast = consumer.request_broadcast("test").await.unwrap();
4701 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4702 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4703 settle().await;
4704 let mut sub = subscribing.await.unwrap();
4705 producer.append_group().unwrap();
4706 assert_eq!(
4707 sub.assert_group().sequence,
4708 0,
4709 "the reconnect's first group must not be filtered by a stale boundary"
4710 );
4711 }
4712
4713 #[tokio::test(start_paused = true)]
4716 async fn test_linger_skipped_on_finish() {
4717 let origin = Info::new(Origin::random())
4718 .with_linger(Duration::from_secs(5))
4719 .produce();
4720 let consumer = origin.consume();
4721 let mut announced = consumer.announced();
4722
4723 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4724 let mut source = origin
4725 .create_broadcast("test", announce().with_hops(hops.clone()))
4726 .unwrap();
4727 settle().await;
4728 let broadcast = consumer.request_broadcast("test").await.unwrap();
4729 announced.assert_next_some("test");
4730
4731 source.finish();
4734 settle().await;
4735 announced.assert_next_none("test");
4736
4737 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4739 settle().await;
4740 let fresh = consumer.request_broadcast("test").await.unwrap();
4741 announced.assert_next_some("test");
4742 assert!(
4743 !fresh.is_clone(&broadcast),
4744 "a finish must not leave a lingering broadcast to splice into"
4745 );
4746 }
4747
4748 #[tokio::test]
4751 async fn test_announce_toggle() {
4752 tokio::time::pause();
4753
4754 let origin = Origin::random().produce();
4755 let consumer = origin.consume();
4756 let mut announced = consumer.announced();
4757
4758 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4759 settle().await;
4760
4761 announced.assert_next_wait();
4763 let broadcast = consumer
4764 .get_broadcast("test")
4765 .expect("offline broadcast is still routable");
4766 assert!(!broadcast.route().announce);
4767
4768 let requested = consumer.request_broadcast("test").await.unwrap();
4770 assert!(requested.is_clone(&broadcast));
4771
4772 source.set_route(announce()).unwrap();
4774 settle().await;
4775 let face = announced.assert_next_some("test");
4776 assert!(face.is_clone(&broadcast));
4777
4778 let mut fresh = origin.consume().announced();
4780 fresh.assert_next_some("test");
4781 fresh.assert_next_wait();
4782
4783 source.set_route(broadcast::Route::new()).unwrap();
4785 settle().await;
4786 announced.assert_next_none("test");
4787 assert!(consumer.get_broadcast("test").is_some());
4788 let mut fresh = origin.consume().announced();
4789 fresh.assert_next_wait();
4790
4791 source.finish();
4792 settle().await;
4793 assert!(consumer.get_broadcast("test").is_none());
4794 }
4795
4796 #[tokio::test]
4799 async fn test_announce_beats_offline() {
4800 tokio::time::pause();
4801
4802 let origin = Origin::random().produce();
4803 let consumer = origin.consume();
4804 let mut announced = consumer.announced();
4805
4806 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4808 settle().await;
4809 announced.assert_next_wait();
4810
4811 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
4814 settle().await;
4815 announced.assert_next_some("test");
4816 let face = consumer.get_broadcast("test").unwrap();
4817 assert!(face.route().announce);
4818 assert_eq!(face.route().cost, 10);
4819
4820 announced_source.finish();
4823 settle().await;
4824 announced.assert_next_none("test");
4825 assert!(consumer.get_broadcast("test").is_some());
4826 }
4827
4828 #[tokio::test]
4831 async fn test_better_source_no_churn() {
4832 tokio::time::pause();
4833
4834 let origin = Origin::random().produce();
4835 let mut announced = origin.consume().announced();
4836
4837 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4840 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4841 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4842 settle().await;
4843 let face = announced.assert_next_some("test");
4844
4845 let _b = origin
4846 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4847 .unwrap();
4848 settle().await;
4849 announced.assert_next_wait();
4850 let current = origin.consume().get_broadcast("test").unwrap();
4851 assert!(current.is_clone(&face), "the broadcast identity must not change");
4852 assert_eq!(current.route().hops, hops_b);
4854 }
4855
4856 #[tokio::test]
4862 async fn test_publisher_mismatch_replaces() {
4863 tokio::time::pause();
4864
4865 let origin = Origin::random().produce();
4866 let consumer = origin.consume();
4867 let mut announced = consumer.announced();
4868
4869 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4870 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4871
4872 let mut source_a = origin
4873 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4874 .unwrap();
4875 settle().await;
4876 let face_a = announced.assert_next_some("test");
4877
4878 let _source_b = origin
4881 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4882 .unwrap();
4883 settle().await;
4884 settle().await;
4885 announced.assert_next_none("test");
4886 let face_b = announced.assert_next_some("test");
4887 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
4888 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4889 assert!(face_a.is_closed(), "the displaced front must close");
4893
4894 source_a.finish();
4896 settle().await;
4897 settle().await;
4898 announced.assert_next_wait();
4899 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4900 }
4901
4902 #[tokio::test]
4905 async fn test_displaced_publisher_reclaims_path_when_replacement_leaves() {
4906 tokio::time::pause();
4907
4908 let origin = Origin::random().produce();
4909 let consumer = origin.consume();
4910 let mut announced = consumer.announced();
4911
4912 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4913 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4914
4915 let mut source_a = origin
4917 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4918 .unwrap();
4919 settle().await;
4920 announced.assert_next_some("test");
4921
4922 let mut source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4925 settle().await;
4926 settle().await;
4927 announced.assert_next_none("test");
4928 announced.assert_next_some("test");
4929 announced.assert_next_wait();
4932
4933 source_b.finish();
4935 settle().await;
4936 settle().await;
4937
4938 assert!(
4939 consumer.get_broadcast("test").is_some(),
4940 "the still-live publisher A should reclaim the path once its replacement leaves"
4941 );
4942 let recovered = consumer.request_broadcast("test").await.unwrap();
4943 assert_eq!(recovered.route().hops, hops_a);
4944
4945 announced.assert_next_none("test");
4946 announced.assert_next_some("test");
4947 announced.assert_next_wait();
4948
4949 source_a.finish();
4950 }
4951
4952 #[tokio::test]
4956 async fn test_reclaimed_publisher_can_still_replace_itself() {
4957 tokio::time::pause();
4958
4959 let origin = Origin::random().produce();
4960 let consumer = origin.consume();
4961
4962 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4963 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4964 let hops_c = OriginList::try_from(vec![Origin::new(3).unwrap()]).unwrap();
4965
4966 let mut source_a1 = origin
4968 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4969 .unwrap();
4970 let mut source_a2 = origin
4971 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4972 .unwrap();
4973 settle().await;
4974 settle().await;
4975 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_a);
4976
4977 let mut source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4979 settle().await;
4980 settle().await;
4981
4982 source_b.finish();
4984 settle().await;
4985 settle().await;
4986 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_a);
4987
4988 source_a1.set_route(announce().with_hops(hops_c.clone())).unwrap();
4992 settle().await;
4993 settle().await;
4994 assert_eq!(
4995 consumer.get_broadcast("test").unwrap().route().hops,
4996 hops_c,
4997 "the new publisher must take the path over, not stand by behind the old front"
4998 );
4999
5000 source_a1.finish();
5001 source_a2.finish();
5002 }
5003
5004 #[tokio::test]
5007 async fn test_repricing_does_not_earn_a_takeover() {
5008 tokio::time::pause();
5009
5010 let origin = Origin::random().produce();
5011 let consumer = origin.consume();
5012 let mut announced = consumer.announced();
5013
5014 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5015 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5016
5017 let mut source_a = origin
5018 .create_broadcast("test", announce().with_hops(hops_a.clone()).with_cost(5))
5019 .unwrap();
5020 settle().await;
5021 announced.assert_next_some("test");
5022
5023 let mut source_b = origin
5025 .create_broadcast("test", announce().with_hops(hops_b.clone()))
5026 .unwrap();
5027 settle().await;
5028 settle().await;
5029 announced.assert_next_none("test");
5030 announced.assert_next_some("test");
5031 announced.assert_next_wait();
5032
5033 source_a
5036 .set_route(announce().with_hops(hops_a.clone()).with_cost(9))
5037 .unwrap();
5038 settle().await;
5039 settle().await;
5040 assert_eq!(
5041 consumer.get_broadcast("test").unwrap().route().hops,
5042 hops_b,
5043 "a repricing must not take the path back from the live front"
5044 );
5045 announced.assert_next_wait();
5046
5047 source_a.finish();
5048 source_b.finish();
5049 }
5050
5051 #[tokio::test]
5056 async fn test_reconnect_wins_over_stale_route() {
5057 tokio::time::pause();
5058
5059 let origin = Origin::random().produce();
5060 let consumer = origin.consume();
5061
5062 let publisher = Origin::new(1).unwrap();
5063 let hops = OriginList::try_from(vec![publisher]).unwrap();
5064
5065 let stale = origin
5068 .create_broadcast("test", announce().with_hops(hops.clone()))
5069 .unwrap();
5070 let mut stale_dynamic = stale.dynamic();
5071 settle().await;
5072
5073 let fresh = origin
5075 .create_broadcast("test", announce().with_hops(hops.clone()))
5076 .unwrap();
5077 let mut fresh_dynamic = fresh.dynamic();
5078 settle().await;
5079 settle().await;
5080
5081 let broadcast = consumer.request_broadcast("test").await.unwrap();
5083 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5084 settle().await;
5085 let _producer = accept_track(&mut fresh_dynamic, "video").await;
5086 settle().await;
5087 subscribing.await.unwrap();
5088 stale_dynamic.assert_no_request();
5089 }
5090
5091 #[tokio::test]
5097 async fn test_carrying_reconnect_switches_immediately() {
5098 tokio::time::pause();
5099
5100 let origin = Origin::random().produce();
5101 let consumer = origin.consume();
5102
5103 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5104
5105 let stale = origin
5106 .create_broadcast("test", announce().with_hops(hops.clone()))
5107 .unwrap();
5108 let mut stale_dynamic = stale.dynamic();
5109 settle().await;
5110
5111 let broadcast = consumer.request_broadcast("test").await.unwrap();
5113 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5114 settle().await;
5115 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
5116 settle().await;
5117 let _subscription = subscribing.await.unwrap();
5119
5120 let fresh = origin
5122 .create_broadcast("test", announce().with_hops(hops.clone()))
5123 .unwrap();
5124 let mut fresh_dynamic = fresh.dynamic();
5125 settle().await;
5126 settle().await;
5127
5128 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
5131 }
5132
5133 #[tokio::test]
5143 async fn test_offline_mismatch_never_evicts_a_live_front() {
5144 tokio::time::pause();
5145
5146 let origin = Origin::random().produce();
5147 let consumer = origin.consume();
5148 let mut announced = consumer.announced();
5149
5150 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5151 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5152
5153 let mut live = origin
5154 .create_broadcast("test", announce().with_hops(hops_live.clone()))
5155 .unwrap();
5156 let mut live_dynamic = live.dynamic();
5157 settle().await;
5158 let face = announced.assert_next_some("test");
5159
5160 let broadcast = consumer.request_broadcast("test").await.unwrap();
5162 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5163 settle().await;
5164 let _producer = accept_track(&mut live_dynamic, "video").await;
5165 settle().await;
5166 subscribing.await.unwrap();
5167
5168 let cache = origin
5171 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
5172 .unwrap();
5173 settle().await;
5174 settle().await;
5175 announced.assert_next_wait();
5176 assert!(!face.is_closed(), "the live front must survive");
5177 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
5178
5179 live.finish();
5182 settle().await;
5183 settle().await;
5184 announced.assert_next_none("test");
5185 let taken = consumer
5186 .get_broadcast("test")
5187 .expect("the parked source must take over");
5188 assert_eq!(taken.route().hops, hops_cache);
5189 announced.assert_next_wait();
5191 drop(cache);
5192 }
5193
5194 #[tokio::test]
5200 async fn test_dispatch_excludes_requester() {
5201 tokio::time::pause();
5202
5203 let origin = Origin::random().produce();
5204 let consumer = origin.consume();
5205
5206 let peer = Origin::new(5).unwrap();
5207 let publisher = Origin::new(1).unwrap();
5208 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
5210 let clean = OriginList::try_from(vec![publisher]).unwrap();
5211
5212 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5213 let mut dynamic_a = source_a.dynamic();
5214 settle().await;
5215 let source_b = origin
5216 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
5217 .unwrap();
5218 let mut dynamic_b = source_b.dynamic();
5219 settle().await;
5220 settle().await;
5221
5222 let shared = consumer.request_broadcast("test").await.unwrap();
5225 let subscribing = shared.track("video").unwrap().subscribe(None);
5226 let _producer_a = accept_track(&mut dynamic_a, "video").await;
5227 settle().await;
5228 subscribing.await.unwrap();
5229
5230 let scoped = consumer.clone().excluding(peer);
5235 let pinned = scoped.request_broadcast("test").await.unwrap();
5236 let subscribing = pinned.track("video").unwrap().subscribe(None);
5237 let _producer_b = accept_track(&mut dynamic_b, "video").await;
5238 settle().await;
5239 subscribing.await.unwrap();
5240 dynamic_a.assert_no_request();
5241 }
5242
5243 #[tokio::test]
5248 async fn test_unknown_publishers_do_not_splice() {
5249 tokio::time::pause();
5250
5251 let origin = Origin::random().produce();
5252 let consumer = origin.consume();
5253 let mut announced = consumer.announced();
5254
5255 let unknown_a = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
5256 let unknown_b = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
5257
5258 let mut source_a = origin
5259 .create_broadcast("test", announce().with_hops(unknown_a.clone()))
5260 .unwrap();
5261 settle().await;
5262 settle().await;
5263 announced.assert_next_some("test");
5264
5265 let source_b = origin
5269 .create_broadcast("test", announce().with_hops(unknown_b))
5270 .unwrap();
5271 settle().await;
5272 settle().await;
5273 announced.assert_next_none("test");
5274 let live = announced.assert_next_some("test");
5275
5276 source_a
5279 .set_route(announce().with_hops(unknown_a).with_cost(9))
5280 .unwrap();
5281 settle().await;
5282 settle().await;
5283 assert!(
5284 consumer.get_broadcast("test").unwrap().is_clone(&live),
5285 "UNKNOWN-to-UNKNOWN repricing must not replace the live front"
5286 );
5287 announced.assert_next_wait();
5288
5289 drop(source_a);
5290 drop(source_b);
5291 }
5292
5293 #[tokio::test]
5296 async fn test_known_publishers_still_splice() {
5297 tokio::time::pause();
5298
5299 let origin = Origin::random().produce();
5300 let consumer = origin.consume();
5301 let mut announced = consumer.announced();
5302
5303 let publisher = Origin::new(1).unwrap();
5304 let hops_a = OriginList::try_from(vec![publisher]).unwrap();
5305 let hops_b = OriginList::try_from(vec![publisher, Origin::new(3).unwrap()]).unwrap();
5306
5307 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
5308 settle().await;
5309 settle().await;
5310 announced.assert_next_some("test");
5311
5312 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
5314 settle().await;
5315 settle().await;
5316 announced.assert_next_wait();
5317
5318 drop(source_a);
5319 drop(source_b);
5320 }
5321
5322 #[tokio::test]
5328 async fn test_standby_join_splices_live_subscriber() {
5329 tokio::time::pause();
5330
5331 let origin = Origin::random().produce();
5332 let consumer = origin.consume();
5333
5334 let publisher = Origin::new(1).unwrap();
5335 let peer = Origin::new(5).unwrap();
5336 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5337 let local = OriginList::try_from(vec![publisher]).unwrap();
5338
5339 let source_remote = origin
5341 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5342 .unwrap();
5343 let mut dynamic_remote = source_remote.dynamic();
5344 settle().await;
5345 settle().await;
5346 let broadcast = consumer.request_broadcast("test").await.unwrap();
5347 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5348 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
5349 settle().await;
5350 let mut sub = subscribing.await.unwrap();
5351 producer_remote.append_group().unwrap();
5352 assert_eq!(sub.assert_group().sequence, 0);
5353
5354 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5358 let mut dynamic_local = source_local.dynamic();
5359 settle().await;
5360 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
5361 settle().await;
5362 sub.assert_no_group();
5363 assert_eq!(producer_local.subscription().unwrap().group_start, None);
5364 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
5365 assert_eq!(sub.assert_group().sequence, 1);
5366 sub.assert_not_closed();
5367 }
5368
5369 #[tokio::test]
5376 async fn test_standby_with_a_partial_track_list_splits_per_track() {
5377 tokio::time::pause();
5378
5379 let origin = Origin::random().produce();
5380 let consumer = origin.consume();
5381
5382 let publisher = Origin::new(1).unwrap();
5383 let peer = Origin::new(5).unwrap();
5384 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5385 let local = OriginList::try_from(vec![publisher]).unwrap();
5386
5387 let source_remote = origin
5389 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5390 .unwrap();
5391 let mut dynamic_remote = source_remote.dynamic();
5392 settle().await;
5393 settle().await;
5394 let broadcast = consumer.request_broadcast("test").await.unwrap();
5395
5396 let subscribing_video = broadcast.track("video").unwrap().subscribe(None);
5397 let subscribing_audio = broadcast.track("audio").unwrap().subscribe(None);
5398 let mut producers_remote = accept_tracks(&mut dynamic_remote, 2).await;
5399 settle().await;
5400
5401 let mut sub_video = subscribing_video.await.unwrap();
5402 let mut sub_audio = subscribing_audio.await.unwrap();
5403 for name in ["video", "audio"] {
5404 producers_remote.get_mut(name).unwrap().append_group().unwrap();
5405 }
5406 assert_eq!(sub_video.assert_group().sequence, 0);
5407 assert_eq!(sub_audio.assert_group().sequence, 0);
5408
5409 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5411 let mut dynamic_local = source_local.dynamic();
5412 settle().await;
5413
5414 let mut producer_local = None;
5415 for _ in 0..2 {
5416 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic_local.requested_track())
5417 .await
5418 .expect("timed out waiting for a track request")
5419 .expect("source closed");
5420 match request.name() {
5421 "video" => producer_local = Some(request.accept(None)),
5422 "audio" => request.reject(Error::NotFound),
5423 other => panic!("unexpected track dispatched: {other}"),
5424 }
5425 }
5426 settle().await;
5427 let mut producer_local = producer_local.expect("the standby was never asked for video");
5428
5429 sub_video.assert_no_group();
5432 assert_eq!(producer_local.subscription().unwrap().group_start, None);
5433 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
5434 assert_eq!(
5435 sub_video.assert_group().sequence,
5436 1,
5437 "video did not move to the standby"
5438 );
5439
5440 producers_remote.get_mut("audio").unwrap().append_group().unwrap();
5442 assert_eq!(
5443 sub_audio.assert_group().sequence,
5444 1,
5445 "audio did not stay on the incumbent"
5446 );
5447 sub_video.assert_not_closed();
5448 sub_audio.assert_not_closed();
5449 }
5450
5451 #[tokio::test]
5458 async fn test_standby_missing_track_keeps_incumbent() {
5459 tokio::time::pause();
5460
5461 let origin = Origin::random().produce();
5462 let consumer = origin.consume();
5463
5464 let publisher = Origin::new(1).unwrap();
5465 let peer = Origin::new(5).unwrap();
5466 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5467 let local = OriginList::try_from(vec![publisher]).unwrap();
5468
5469 let source_remote = origin
5471 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5472 .unwrap();
5473 let mut dynamic_remote = source_remote.dynamic();
5474 settle().await;
5475 settle().await;
5476 let broadcast = consumer.request_broadcast("test").await.unwrap();
5477 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5478 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
5479 settle().await;
5480 let mut sub = subscribing.await.unwrap();
5481 producer_remote.append_group().unwrap();
5482 assert_eq!(sub.assert_group().sequence, 0);
5483
5484 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5487 let mut dynamic_local = source_local.dynamic();
5488 settle().await;
5489 let request = dynamic_local.requested_track().await.unwrap();
5490 assert_eq!(request.name(), "audio");
5491 request.reject(Error::NotFound);
5492 settle().await;
5493
5494 producer_remote.append_group().unwrap();
5496 assert_eq!(sub.assert_group().sequence, 1);
5497 sub.assert_not_closed();
5498
5499 source_remote.abort(Error::Dropped).unwrap();
5502 settle().await;
5503 settle().await;
5504 sub.assert_closed();
5505 dynamic_local.assert_no_request();
5506
5507 let retry = broadcast.track("audio").unwrap().subscribe(None);
5509 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
5510 settle().await;
5511 let mut sub = retry.await.expect("a fresh request must reach the standby");
5512 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
5513 assert_eq!(sub.assert_group().sequence, 2);
5514 }
5515
5516 #[tokio::test]
5521 async fn test_unservable_track_retried_by_a_later_request() {
5522 tokio::time::pause();
5523
5524 let origin = Origin::random().produce();
5525 let consumer = origin.consume();
5526
5527 let source = origin.create_broadcast("test", announce()).unwrap();
5528 let mut dynamic = source.dynamic();
5529 settle().await;
5530 settle().await;
5531 let broadcast = consumer.request_broadcast("test").await.unwrap();
5532
5533 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5535 let request = dynamic.requested_track().await.unwrap();
5536 request.reject(Error::NotFound);
5537 settle().await;
5538 assert!(matches!(subscribing.await, Err(Error::NotFound)));
5539
5540 let retry = broadcast.track("audio").unwrap().subscribe(None);
5542 let mut producer = accept_track(&mut dynamic, "audio").await;
5543 settle().await;
5544 let mut sub = retry.await.expect("a fresh request must reach the source");
5545 producer.append_group().unwrap();
5546 assert_eq!(sub.assert_group().sequence, 0);
5547 }
5548
5549 #[tokio::test]
5555 async fn test_track_dying_without_progress_aborts() {
5556 tokio::time::pause();
5557
5558 let origin = Origin::random().produce();
5559 let consumer = origin.consume();
5560
5561 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5562 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5563 let mut dynamic = source.dynamic();
5564 settle().await;
5565 settle().await;
5566 let broadcast = consumer.request_broadcast("test").await.unwrap();
5567
5568 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5569 let producer = accept_track(&mut dynamic, "video").await;
5570 settle().await;
5571 let mut sub = subscribing.await.unwrap();
5572
5573 drop(producer);
5576 settle().await;
5577 sub.assert_closed();
5578 dynamic.assert_no_request();
5579
5580 let retry = broadcast.track("video").unwrap().subscribe(None);
5582 let mut producer = accept_track(&mut dynamic, "video").await;
5583 settle().await;
5584 let mut sub = retry.await.expect("a fresh request must reach the source");
5585 producer.append_group().unwrap();
5586 assert_eq!(sub.assert_group().sequence, 0);
5587 }
5588
5589 #[tokio::test]
5594 async fn test_delivered_copy_death_survives_unrelated_wakes() {
5595 tokio::time::pause();
5596
5597 let origin = Origin::random().produce();
5598 let consumer = origin.consume();
5599
5600 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5601 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5602 let mut dynamic = source.dynamic();
5603 settle().await;
5604 settle().await;
5605 let broadcast = consumer.request_broadcast("test").await.unwrap();
5606
5607 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5608 let mut producer = accept_track(&mut dynamic, "video").await;
5609 settle().await;
5610 let mut sub = subscribing.await.unwrap();
5611 producer.append_group().unwrap();
5612 assert_eq!(sub.assert_group().sequence, 0);
5613
5614 drop(sub);
5617 settle().await;
5618 let resubscribing = broadcast.track("video").unwrap().subscribe(None);
5619 settle().await;
5620 let mut sub = resubscribing.await.unwrap();
5621 assert_eq!(sub.assert_group().sequence, 0, "cached group re-served");
5622
5623 drop(producer);
5626 let mut producer = accept_track(&mut dynamic, "video").await;
5627 settle().await;
5628 producer.create_group(group::Info { sequence: 1 }).unwrap();
5629 assert_eq!(sub.assert_group().sequence, 1);
5630 sub.assert_not_closed();
5631 }
5632
5633 #[tokio::test]
5639 async fn test_per_track_fallback_respects_exclusion() {
5640 tokio::time::pause();
5641
5642 let origin = Origin::random().produce();
5643 let consumer = origin.consume();
5644
5645 let publisher = Origin::new(1).unwrap();
5646 let peer = Origin::new(5).unwrap();
5647 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5648 let local = OriginList::try_from(vec![publisher]).unwrap();
5649
5650 let source_tainted = origin
5652 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5653 .unwrap();
5654 let mut dynamic_tainted = source_tainted.dynamic();
5655 settle().await;
5656 settle().await;
5657
5658 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5661 let mut dynamic_clean = source_clean.dynamic();
5662 settle().await;
5663
5664 let scoped = consumer.clone().excluding(peer);
5665 let broadcast = scoped.request_broadcast("test").await.unwrap();
5666 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5667 settle().await;
5668
5669 let request = dynamic_clean.requested_track().await.unwrap();
5673 request.reject(Error::NotFound);
5674 settle().await;
5675 dynamic_tainted.assert_no_request();
5676 }
5677
5678 #[tokio::test]
5682 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
5683 tokio::time::pause();
5684
5685 let origin = Origin::random().produce();
5686 let consumer = origin.consume();
5687
5688 let publisher = Origin::new(1).unwrap();
5689 let peer = Origin::new(5).unwrap();
5690 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5691 let local = OriginList::try_from(vec![publisher]).unwrap();
5692
5693 let source_tainted = origin
5694 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5695 .unwrap();
5696 let mut dynamic_tainted = source_tainted.dynamic();
5697 settle().await;
5698 settle().await;
5699 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5700 let mut dynamic_clean = source_clean.dynamic();
5701 settle().await;
5702
5703 let scoped = consumer.clone().excluding(peer);
5704 let broadcast = scoped.request_broadcast("test").await.unwrap();
5705 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5706 let _clean = accept_track(&mut dynamic_clean, "video").await;
5707 settle().await;
5708
5709 source_clean.abort(Error::Dropped).unwrap();
5711 settle().await;
5712 settle().await;
5713 dynamic_tainted.assert_no_request();
5714
5715 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5718 }
5719
5720 #[tokio::test]
5725 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
5726 tokio::time::pause();
5727
5728 let origin = Origin::random().produce();
5729 let consumer = origin.consume();
5730
5731 let publisher = Origin::new(1).unwrap();
5732 let peer = Origin::new(5).unwrap();
5733 let local = OriginList::try_from(vec![publisher]).unwrap();
5734 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5735
5736 let source_clean = origin
5740 .create_broadcast("test", announce().with_hops(local).with_cost(5))
5741 .unwrap();
5742 let mut dynamic_clean = source_clean.dynamic();
5743 settle().await;
5744 settle().await;
5745 let scoped = consumer.clone().excluding(peer);
5746 let broadcast = scoped.request_broadcast("test").await.unwrap();
5747 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5748 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
5749 settle().await;
5750 let mut sub = subscribing.await.unwrap();
5751 producer_clean.append_group().unwrap();
5752 assert_eq!(sub.assert_group().sequence, 0);
5753
5754 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
5761 tainted.advertised = 1;
5762 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
5763 let mut dynamic_tainted = source_tainted.dynamic();
5764 settle().await;
5765 settle().await;
5766 dynamic_tainted.assert_no_request();
5767 producer_clean.append_group().unwrap();
5768 assert_eq!(sub.assert_group().sequence, 1);
5769 sub.assert_not_closed();
5770
5771 drop(sub);
5774 drop(broadcast);
5775 drop(scoped);
5776 settle().await;
5777 let mut bumped = announce().with_hops(via_peer).with_cost(1);
5778 bumped.advertised = 1;
5779 source_tainted.set_route(bumped).unwrap();
5780 settle().await;
5781 let plain = consumer.request_broadcast("test").await.unwrap();
5782 let _plain_track = plain.track("video").unwrap().subscribe(None);
5783 settle().await;
5784 settle().await;
5785 assert!(
5786 dynamic_tainted.requested_track().now_or_never().is_some(),
5787 "the front must be free to use the route again once the peer is gone"
5788 );
5789 }
5790
5791 #[tokio::test]
5795 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
5796 tokio::time::pause();
5797
5798 let origin = Origin::random().produce();
5799 let consumer = origin.consume();
5800 let mut dynamic = origin.dynamic();
5801
5802 let peer = Origin::new(5).unwrap();
5803 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5804 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5805 settle().await;
5806 settle().await;
5807
5808 let scoped = consumer.clone().excluding(peer);
5809 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5810 assert!(
5811 dynamic.requested_broadcast().now_or_never().is_none(),
5812 "the dynamic handler was asked to route around the exclusion"
5813 );
5814
5815 let _pending = scoped.request_broadcast("other");
5817 settle().await;
5818 assert!(
5819 dynamic.requested_broadcast().now_or_never().is_some(),
5820 "a genuinely missing path must still fall back"
5821 );
5822 }
5823
5824 #[tokio::test]
5827 async fn test_dispatch_all_tainted_unroutable() {
5828 tokio::time::pause();
5829
5830 let origin = Origin::random().produce();
5831 let consumer = origin.consume();
5832
5833 let peer = Origin::new(5).unwrap();
5834 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5835 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5836 settle().await;
5837 settle().await;
5838
5839 let scoped = consumer.clone().excluding(peer);
5840 match scoped.request_broadcast("test").await {
5841 Err(Error::Unroutable) => {}
5842 Err(err) => panic!("expected Unroutable, got {err:?}"),
5843 Ok(_) => panic!("expected Unroutable, got a broadcast"),
5844 }
5845
5846 consumer.request_broadcast("test").await.unwrap();
5848 }
5849
5850 #[tokio::test]
5851 async fn test_duplicate_reverse() {
5852 tokio::time::pause();
5853
5854 let origin = Origin::random().produce();
5855
5856 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5857 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5858 settle().await;
5859 assert!(origin.consume().get_broadcast("test").is_some());
5860
5861 broadcast2.finish();
5863 settle().await;
5864 assert!(origin.consume().get_broadcast("test").is_some());
5865
5866 broadcast1.finish();
5867 settle().await;
5868 assert!(origin.consume().get_broadcast("test").is_none());
5869 }
5870
5871 #[tokio::test]
5872 async fn test_deterministic_tiebreak() {
5873 tokio::time::pause();
5874
5875 fn hops(ids: &[u64]) -> OriginList {
5876 OriginList::try_from(
5877 ids.iter()
5878 .copied()
5879 .map(|id| Origin::new(id).unwrap())
5880 .collect::<Vec<_>>(),
5881 )
5882 .unwrap()
5883 }
5884
5885 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
5888 let origin = Origin::random().produce();
5889 let _a = origin
5890 .create_broadcast("test", announce().with_hops(hops(first)))
5891 .unwrap();
5892 let _b = origin
5893 .create_broadcast("test", announce().with_hops(hops(second)))
5894 .unwrap();
5895 settle().await;
5896 origin.consume().get_broadcast("test").unwrap().route().hops
5897 }
5898
5899 let forward = winner(&[5, 20], &[5, 40]).await;
5903 let reverse = winner(&[5, 40], &[5, 20]).await;
5904 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
5905
5906 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
5908 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
5909 }
5910
5911 #[tokio::test]
5916 async fn test_many_announces() {
5917 let origin = Origin::random().produce();
5918
5919 let mut consumer = origin.consume().announced();
5920 let mut broadcasts = Vec::new();
5922 for i in 0..256 {
5923 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5924 settle().await;
5925 }
5926
5927 for i in 0..256 {
5928 consumer.assert_next_some(format!("test{i:03}"));
5929 }
5930 consumer.assert_next_wait();
5931 }
5932
5933 #[tokio::test]
5934 async fn test_many_announces_try() {
5935 let origin = Origin::random().produce();
5936
5937 let mut consumer = origin.consume().announced();
5938 let mut broadcasts = Vec::new();
5940 for i in 0..256 {
5941 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5942 settle().await;
5943 }
5944
5945 for i in 0..256 {
5946 consumer.assert_try_next_some(format!("test{i:03}"));
5947 }
5948 }
5949
5950 #[tokio::test]
5951 async fn test_with_root_basic() {
5952 let origin = Origin::random().produce();
5953
5954 let foo_producer = origin.with_root("foo").expect("should create root");
5956 assert_eq!(foo_producer.root().as_str(), "foo");
5957
5958 let mut consumer = origin.consume().announced();
5959
5960 let _broadcast = foo_producer
5962 .create_broadcast("bar/baz", announce())
5963 .expect("publish allowed");
5964 settle().await;
5965 consumer.assert_next_some("foo/bar/baz");
5967
5968 let mut foo_consumer = foo_producer.consume().announced();
5970 foo_consumer.assert_next_some("bar/baz");
5971 }
5972
5973 #[tokio::test]
5974 async fn test_with_root_nested() {
5975 let origin = Origin::random().produce();
5976
5977 let foo_producer = origin.with_root("foo").expect("should create foo root");
5979 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
5980 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
5981
5982 let mut consumer = origin.consume().announced();
5983
5984 let _broadcast = foo_bar_producer
5986 .create_broadcast("baz", announce())
5987 .expect("publish allowed");
5988 settle().await;
5989 consumer.assert_next_some("foo/bar/baz");
5991
5992 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
5994 foo_bar_consumer.assert_next_some("baz");
5995 }
5996
5997 #[tokio::test]
5998 async fn test_publish_scope_allows() {
5999 let origin = Origin::random().produce();
6000
6001 let limited_producer = origin
6003 .scope(&["allowed/path1".into(), "allowed/path2".into()])
6004 .expect("should create limited producer");
6005
6006 let _broadcast = limited_producer
6008 .create_broadcast("allowed/path1", announce())
6009 .expect("publish allowed");
6010 let _keep2 = limited_producer
6011 .create_broadcast("allowed/path1/nested", announce())
6012 .expect("publish allowed");
6013 let _keep3 = limited_producer
6014 .create_broadcast("allowed/path2", announce())
6015 .expect("publish allowed");
6016 settle().await;
6017
6018 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
6020 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
6022 }
6023
6024 #[tokio::test]
6025 async fn test_publish_max_parts() {
6026 let origin = Origin::random().produce();
6027
6028 let at_limit = (0..Path::MAX_PARTS)
6029 .map(|i| i.to_string())
6030 .collect::<Vec<_>>()
6031 .join("/");
6032 let _broadcast = origin
6033 .create_broadcast(at_limit.as_str(), announce())
6034 .expect("publish allowed");
6035 settle().await;
6036
6037 let too_deep = format!("{at_limit}/extra");
6038 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
6039
6040 let rooted = origin.with_root("root").expect("wildcard allows any root");
6042 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
6043 }
6044
6045 #[tokio::test]
6046 async fn test_publish_scope_empty() {
6047 let origin = Origin::random().produce();
6048
6049 assert!(origin.scope(&[]).is_none());
6051 }
6052
6053 #[tokio::test]
6054 async fn test_consume_scope_filters() {
6055 let origin = Origin::random().produce();
6056
6057 let mut consumer = origin.consume().announced();
6058
6059 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
6061 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
6062 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
6063 settle().await;
6064
6065 let mut limited_consumer = origin
6067 .consume()
6068 .scope(&["allowed".into()])
6069 .expect("should create limited consumer")
6070 .announced();
6071
6072 limited_consumer.assert_next_some("allowed");
6074 limited_consumer.assert_next_some("allowed/nested");
6075 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
6079 consumer.assert_next_some("allowed/nested");
6080 consumer.assert_next_some("notallowed");
6081 }
6082
6083 #[tokio::test]
6084 async fn test_consume_scope_multiple_prefixes() {
6085 let origin = Origin::random().produce();
6086
6087 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
6088 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
6089 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
6090 settle().await;
6091
6092 let mut limited_consumer = origin
6094 .consume()
6095 .scope(&["foo".into(), "bar".into()])
6096 .expect("should create limited consumer")
6097 .announced();
6098
6099 limited_consumer.assert_next_some("bar/test");
6101 limited_consumer.assert_next_some("foo/test");
6102 limited_consumer.assert_next_wait(); }
6104
6105 #[tokio::test]
6106 async fn test_with_root_and_publish_scope() {
6107 let origin = Origin::random().produce();
6108
6109 let foo_producer = origin.with_root("foo").expect("should create foo root");
6111
6112 let limited_producer = foo_producer
6114 .scope(&["bar".into(), "goop/pee".into()])
6115 .expect("should create limited producer");
6116
6117 let mut consumer = origin.consume().announced();
6118
6119 let _broadcast = limited_producer
6121 .create_broadcast("bar", announce())
6122 .expect("publish allowed");
6123 let _keep2 = limited_producer
6124 .create_broadcast("bar/nested", announce())
6125 .expect("publish allowed");
6126 let _keep3 = limited_producer
6127 .create_broadcast("goop/pee", announce())
6128 .expect("publish allowed");
6129 let _keep4 = limited_producer
6130 .create_broadcast("goop/pee/nested", announce())
6131 .expect("publish allowed");
6132 settle().await;
6133
6134 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
6136 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
6138
6139 consumer.assert_next_some("foo/bar");
6141 consumer.assert_next_some("foo/bar/nested");
6142 consumer.assert_next_some("foo/goop/pee");
6143 consumer.assert_next_some("foo/goop/pee/nested");
6144 }
6145
6146 #[tokio::test]
6147 async fn test_with_root_and_consume_scope() {
6148 let origin = Origin::random().produce();
6149
6150 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
6152 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
6153 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
6154 settle().await;
6155
6156 let foo_producer = origin.with_root("foo").expect("should create foo root");
6158
6159 let mut limited_consumer = foo_producer
6161 .consume()
6162 .scope(&["bar".into(), "goop/pee".into()])
6163 .expect("should create limited consumer")
6164 .announced();
6165
6166 limited_consumer.assert_next_some("bar/test");
6168 limited_consumer.assert_next_some("goop/pee/test");
6169 limited_consumer.assert_next_wait(); }
6171
6172 #[tokio::test]
6173 async fn test_with_root_unauthorized() {
6174 let origin = Origin::random().produce();
6175
6176 let limited_producer = origin
6178 .scope(&["allowed".into()])
6179 .expect("should create limited producer");
6180
6181 assert!(limited_producer.with_root("notallowed").is_none());
6183
6184 let allowed_root = limited_producer
6186 .with_root("allowed")
6187 .expect("should create allowed root");
6188 assert_eq!(allowed_root.root().as_str(), "allowed");
6189 }
6190
6191 #[tokio::test]
6192 async fn test_wildcard_permission() {
6193 let origin = Origin::random().produce();
6194
6195 let root_producer = origin.clone();
6197
6198 let _broadcast = root_producer
6200 .create_broadcast("any/path", announce())
6201 .expect("publish allowed");
6202 let _keep2 = root_producer
6203 .create_broadcast("other/path", announce())
6204 .expect("publish allowed");
6205 settle().await;
6206
6207 let foo_producer = root_producer.with_root("foo").expect("should create any root");
6209 assert_eq!(foo_producer.root().as_str(), "foo");
6210 }
6211
6212 #[tokio::test]
6213 async fn test_consume_broadcast_with_permissions() {
6214 let origin = Origin::random().produce();
6215
6216 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
6217 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
6218 settle().await;
6219
6220 let limited_consumer = origin
6222 .consume()
6223 .scope(&["allowed".into()])
6224 .expect("should create limited consumer");
6225
6226 let result = limited_consumer.get_broadcast("allowed/test");
6228 assert!(result.is_some());
6229 assert!(
6230 result
6231 .unwrap()
6232 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
6233 );
6234
6235 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
6237
6238 let consumer = origin.consume();
6240 assert!(consumer.get_broadcast("allowed/test").is_some());
6241 assert!(consumer.get_broadcast("notallowed/test").is_some());
6242 }
6243
6244 #[tokio::test]
6245 async fn test_nested_paths_with_permissions() {
6246 let origin = Origin::random().produce();
6247
6248 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
6250
6251 let _broadcast = limited_producer
6253 .create_broadcast("a/b/c", announce())
6254 .expect("publish allowed");
6255 let _keep2 = limited_producer
6256 .create_broadcast("a/b/c/d", announce())
6257 .expect("publish allowed");
6258 let _keep3 = limited_producer
6259 .create_broadcast("a/b/c/d/e", announce())
6260 .expect("publish allowed");
6261 settle().await;
6262
6263 assert!(limited_producer.create_broadcast("a", announce()).is_err());
6265 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
6266 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
6267 }
6268
6269 #[tokio::test]
6270 async fn test_multiple_consumers_with_different_permissions() {
6271 let origin = Origin::random().produce();
6272
6273 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
6275 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
6276 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
6277 settle().await;
6278
6279 let mut foo_consumer = origin
6281 .consume()
6282 .scope(&["foo".into()])
6283 .expect("should create foo consumer")
6284 .announced();
6285
6286 let mut bar_consumer = origin
6287 .consume()
6288 .scope(&["bar".into()])
6289 .expect("should create bar consumer")
6290 .announced();
6291
6292 let mut foobar_consumer = origin
6293 .consume()
6294 .scope(&["foo".into(), "bar".into()])
6295 .expect("should create foobar consumer")
6296 .announced();
6297
6298 foo_consumer.assert_next_some("foo/test");
6300 foo_consumer.assert_next_wait();
6301
6302 bar_consumer.assert_next_some("bar/test");
6303 bar_consumer.assert_next_wait();
6304
6305 foobar_consumer.assert_next_some("bar/test");
6306 foobar_consumer.assert_next_some("foo/test");
6307 foobar_consumer.assert_next_wait();
6308 }
6309
6310 #[tokio::test]
6311 async fn test_select_with_empty_prefix() {
6312 let origin = Origin::random().produce();
6313
6314 let demo_producer = origin.with_root("demo").expect("should create demo root");
6316 let limited_producer = demo_producer
6317 .scope(&["worm-node".into(), "foobar".into()])
6318 .expect("should create limited producer");
6319
6320 let _broadcast1 = limited_producer
6322 .create_broadcast("worm-node/test", announce())
6323 .expect("publish allowed");
6324 let _broadcast2 = limited_producer
6325 .create_broadcast("foobar/test", announce())
6326 .expect("publish allowed");
6327 settle().await;
6328
6329 let mut consumer = limited_producer
6331 .consume()
6332 .scope(&["".into()])
6333 .expect("should create consumer with empty prefix")
6334 .announced();
6335
6336 let a1 = consumer.try_next().expect("expected first announcement");
6338 let a2 = consumer.try_next().expect("expected second announcement");
6339 consumer.assert_next_wait();
6340
6341 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6342 paths.sort();
6343 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
6344 }
6345
6346 #[tokio::test]
6347 async fn test_select_narrowing_scope() {
6348 let origin = Origin::random().produce();
6349
6350 let demo_producer = origin.with_root("demo").expect("should create demo root");
6352 let limited_producer = demo_producer
6353 .scope(&["worm-node".into(), "foobar".into()])
6354 .expect("should create limited producer");
6355
6356 let _broadcast1 = limited_producer
6358 .create_broadcast("worm-node", announce())
6359 .expect("publish allowed");
6360 let _broadcast2 = limited_producer
6361 .create_broadcast("worm-node/foo", announce())
6362 .expect("publish allowed");
6363 let _broadcast3 = limited_producer
6364 .create_broadcast("foobar/bar", announce())
6365 .expect("publish allowed");
6366 settle().await;
6367
6368 let mut worm_consumer = limited_producer
6370 .consume()
6371 .scope(&["worm-node".into()])
6372 .expect("should create worm-node consumer")
6373 .announced();
6374
6375 worm_consumer.assert_next_some("worm-node");
6377 worm_consumer.assert_next_some("worm-node/foo");
6378 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
6382 .consume()
6383 .scope(&["worm-node/foo".into()])
6384 .expect("should create worm-node/foo consumer")
6385 .announced();
6386
6387 foo_consumer.assert_next_some("worm-node/foo");
6388 foo_consumer.assert_next_wait(); }
6390
6391 #[tokio::test]
6392 async fn test_select_multiple_roots_with_empty_prefix() {
6393 let origin = Origin::random().produce();
6394
6395 let limited_producer = origin
6397 .scope(&["app1".into(), "app2".into(), "shared".into()])
6398 .expect("should create limited producer");
6399
6400 let _broadcast1 = limited_producer
6402 .create_broadcast("app1/data", announce())
6403 .expect("publish allowed");
6404 let _broadcast2 = limited_producer
6405 .create_broadcast("app2/config", announce())
6406 .expect("publish allowed");
6407 let _broadcast3 = limited_producer
6408 .create_broadcast("shared/resource", announce())
6409 .expect("publish allowed");
6410 settle().await;
6411
6412 let mut consumer = limited_producer
6414 .consume()
6415 .scope(&["".into()])
6416 .expect("should create consumer with empty prefix")
6417 .announced();
6418
6419 consumer.assert_next_some("app1/data");
6421 consumer.assert_next_some("app2/config");
6422 consumer.assert_next_some("shared/resource");
6423 consumer.assert_next_wait();
6424 }
6425
6426 #[tokio::test]
6427 async fn test_publish_scope_with_empty_prefix() {
6428 let origin = Origin::random().produce();
6429
6430 let limited_producer = origin
6432 .scope(&["services/api".into(), "services/web".into()])
6433 .expect("should create limited producer");
6434
6435 let same_producer = limited_producer
6437 .scope(&["".into()])
6438 .expect("should create producer with empty prefix");
6439
6440 let _broadcast = same_producer
6442 .create_broadcast("services/api", announce())
6443 .expect("publish allowed");
6444 let _keep2 = same_producer
6445 .create_broadcast("services/web", announce())
6446 .expect("publish allowed");
6447 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
6448 assert!(same_producer.create_broadcast("other", announce()).is_err());
6449 }
6450
6451 #[tokio::test]
6452 async fn test_select_narrowing_to_deeper_path() {
6453 let origin = Origin::random().produce();
6454
6455 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
6457
6458 let _broadcast1 = limited_producer
6460 .create_broadcast("org/team1/project1", announce())
6461 .expect("publish allowed");
6462 let _broadcast2 = limited_producer
6463 .create_broadcast("org/team1/project2", announce())
6464 .expect("publish allowed");
6465 let _broadcast3 = limited_producer
6466 .create_broadcast("org/team2/project1", announce())
6467 .expect("publish allowed");
6468 settle().await;
6469
6470 let mut team2_consumer = limited_producer
6472 .consume()
6473 .scope(&["org/team2".into()])
6474 .expect("should create team2 consumer")
6475 .announced();
6476
6477 team2_consumer.assert_next_some("org/team2/project1");
6478 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
6482 .consume()
6483 .scope(&["org/team1/project1".into()])
6484 .expect("should create project1 consumer")
6485 .announced();
6486
6487 project1_consumer.assert_next_some("org/team1/project1");
6489 project1_consumer.assert_next_wait();
6490 }
6491
6492 #[tokio::test]
6493 async fn test_select_with_non_matching_prefix() {
6494 let origin = Origin::random().produce();
6495
6496 let limited_producer = origin
6498 .scope(&["allowed/path".into()])
6499 .expect("should create limited producer");
6500
6501 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
6503
6504 assert!(limited_producer.scope(&["other/path".into()]).is_none());
6506 }
6507
6508 #[tokio::test]
6511 async fn test_with_root_trailing_slash_consumer() {
6512 let origin = Origin::random().produce();
6513
6514 let prefix = "some_prefix/".to_string();
6516 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6517
6518 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6519 settle().await;
6520 consumer.assert_next_some("test");
6521 }
6522
6523 #[tokio::test]
6525 async fn test_with_root_trailing_slash_producer() {
6526 let origin = Origin::random().produce();
6527
6528 let prefix = "some_prefix/".to_string();
6530 let rooted = origin.with_root(prefix).unwrap();
6531
6532 let _b = rooted.create_broadcast("test", announce()).unwrap();
6533 settle().await;
6534
6535 let mut consumer = rooted.consume().announced();
6536 consumer.assert_next_some("test");
6537 }
6538
6539 #[tokio::test]
6541 async fn test_with_root_trailing_slash_unannounce() {
6542 tokio::time::pause();
6543
6544 let origin = Origin::random().produce();
6545
6546 let prefix = "some_prefix/".to_string();
6547 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6548
6549 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6550 settle().await;
6551 consumer.assert_next_some("test");
6552
6553 b.finish();
6555 settle().await;
6556
6557 consumer.assert_next_none("test");
6559 }
6560
6561 #[tokio::test]
6562 async fn test_select_maintains_access_with_wider_prefix() {
6563 let origin = Origin::random().produce();
6564
6565 let demo_producer = origin.with_root("demo").expect("should create demo root");
6567 let user_producer = demo_producer
6568 .scope(&["worm-node".into(), "foobar".into()])
6569 .expect("should create user producer");
6570
6571 let _broadcast1 = user_producer
6573 .create_broadcast("worm-node/data", announce())
6574 .expect("publish allowed");
6575 let _broadcast2 = user_producer
6576 .create_broadcast("foobar", announce())
6577 .expect("publish allowed");
6578 settle().await;
6579
6580 let mut consumer = user_producer
6582 .consume()
6583 .scope(&["".into()])
6584 .expect("scope with empty prefix should not fail when user has specific permissions")
6585 .announced();
6586
6587 let a1 = consumer.try_next().expect("expected first announcement");
6589 let a2 = consumer.try_next().expect("expected second announcement");
6590 consumer.assert_next_wait();
6591
6592 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6593 paths.sort();
6594 assert_eq!(paths, ["foobar", "worm-node/data"]);
6595
6596 let mut narrow_consumer = user_producer
6598 .consume()
6599 .scope(&["worm-node".into()])
6600 .expect("should be able to narrow scope to worm-node")
6601 .announced();
6602
6603 narrow_consumer.assert_next_some("worm-node/data");
6604 narrow_consumer.assert_next_wait(); }
6606
6607 #[tokio::test]
6608 async fn test_duplicate_prefixes_deduped() {
6609 let origin = Origin::random().produce();
6610
6611 let producer = origin
6613 .scope(&["demo".into(), "demo".into()])
6614 .expect("should create producer");
6615
6616 let _broadcast = producer
6617 .create_broadcast("demo/stream", announce())
6618 .expect("publish allowed");
6619 settle().await;
6620
6621 let mut consumer = producer.consume().announced();
6622 consumer.assert_next_some("demo/stream");
6623 consumer.assert_next_wait();
6624 }
6625
6626 #[tokio::test]
6627 async fn test_overlapping_prefixes_deduped() {
6628 let origin = Origin::random().produce();
6629
6630 let producer = origin
6632 .scope(&["demo".into(), "demo/foo".into()])
6633 .expect("should create producer");
6634
6635 let _broadcast = producer
6637 .create_broadcast("demo/bar/stream", announce())
6638 .expect("publish allowed");
6639 settle().await;
6640
6641 let mut consumer = producer.consume().announced();
6642 consumer.assert_next_some("demo/bar/stream");
6643 consumer.assert_next_wait();
6644 }
6645
6646 #[tokio::test]
6647 async fn test_overlapping_prefixes_no_duplicate_announcements() {
6648 let origin = Origin::random().produce();
6649
6650 let producer = origin
6652 .scope(&["demo".into(), "demo/foo".into()])
6653 .expect("should create producer");
6654
6655 let _broadcast = producer
6656 .create_broadcast("demo/foo/stream", announce())
6657 .expect("publish allowed");
6658 settle().await;
6659
6660 let mut consumer = producer.consume().announced();
6661 consumer.assert_next_some("demo/foo/stream");
6663 consumer.assert_next_wait();
6664 }
6665
6666 #[tokio::test]
6667 async fn test_allowed_returns_deduped_prefixes() {
6668 let origin = Origin::random().produce();
6669
6670 let producer = origin
6671 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
6672 .expect("should create producer");
6673
6674 let allowed: Vec<_> = producer.allowed().collect();
6675 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
6676 }
6677
6678 #[tokio::test]
6679 async fn test_announced_broadcast_already_announced() {
6680 let origin = Origin::random().produce();
6681
6682 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6683 settle().await;
6684
6685 let consumer = origin.consume();
6686 let result = consumer.announced_broadcast("test").await.expect("should find it");
6687 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6688 }
6689
6690 #[tokio::test]
6691 async fn test_announced_broadcast_delayed() {
6692 tokio::time::pause();
6693
6694 let origin = Origin::random().produce();
6695
6696 let consumer = origin.consume();
6697
6698 let wait = tokio::spawn({
6700 let consumer = consumer.clone();
6701 async move { consumer.announced_broadcast("test").await }
6702 });
6703
6704 tokio::task::yield_now().await;
6706
6707 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6708 settle().await;
6709
6710 let result = wait.await.unwrap().expect("should find it");
6711 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6712 }
6713
6714 #[tokio::test]
6715 async fn test_announced_broadcast_ignores_unrelated_paths() {
6716 tokio::time::pause();
6717
6718 let origin = Origin::random().produce();
6719
6720 let consumer = origin.consume();
6721
6722 let wait = tokio::spawn({
6723 let consumer = consumer.clone();
6724 async move { consumer.announced_broadcast("target").await }
6725 });
6726
6727 tokio::task::yield_now().await;
6728
6729 let _other = origin.create_broadcast("other", announce()).unwrap();
6731 settle().await;
6732 tokio::task::yield_now().await;
6733 assert!(!wait.is_finished(), "must not resolve on unrelated path");
6734
6735 let _target = origin.create_broadcast("target", announce()).unwrap();
6736 settle().await;
6737 let result = wait.await.unwrap().expect("should find target");
6738 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
6739 }
6740
6741 #[tokio::test]
6742 async fn test_announced_broadcast_skips_nested_paths() {
6743 tokio::time::pause();
6744
6745 let origin = Origin::random().produce();
6746
6747 let consumer = origin.consume();
6748
6749 let wait = tokio::spawn({
6750 let consumer = consumer.clone();
6751 async move { consumer.announced_broadcast("foo").await }
6752 });
6753
6754 tokio::task::yield_now().await;
6755
6756 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
6758 settle().await;
6759 tokio::task::yield_now().await;
6760 assert!(!wait.is_finished(), "must not resolve on a nested path");
6761
6762 let _exact = origin.create_broadcast("foo", announce()).unwrap();
6763 settle().await;
6764 let result = wait.await.unwrap().expect("should find foo exactly");
6765 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
6766 }
6767
6768 #[tokio::test]
6769 async fn test_announced_broadcast_disallowed() {
6770 let origin = Origin::random().produce();
6771 let limited = origin
6772 .consume()
6773 .scope(&["allowed".into()])
6774 .expect("should create limited");
6775
6776 assert!(limited.announced_broadcast("notallowed").await.is_none());
6778 }
6779
6780 #[tokio::test]
6781 async fn test_announced_broadcast_scope_too_narrow() {
6782 let origin = Origin::random().produce();
6785 let limited = origin
6786 .consume()
6787 .scope(&["foo/specific".into()])
6788 .expect("should create limited");
6789
6790 let result = limited
6792 .announced_broadcast("foo")
6793 .now_or_never()
6794 .expect("must not block");
6795 assert!(result.is_none());
6796 }
6797
6798 #[tokio::test]
6802 async fn test_coalesce_announce_then_unannounce() {
6803 tokio::time::pause();
6805
6806 let origin = Origin::random().produce();
6807 let mut announced = origin.consume().announced();
6808
6809 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6810 settle().await;
6811 broadcast.finish();
6812
6813 settle().await;
6814
6815 announced.assert_next_wait();
6816 }
6817
6818 #[tokio::test]
6819 async fn test_coalesce_announce_unannounce_announce() {
6820 tokio::time::pause();
6823
6824 let origin = Origin::random().produce();
6825 let mut announced = origin.consume().announced();
6826
6827 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6828 settle().await;
6829 broadcast1.finish();
6830 settle().await;
6831 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6832 settle().await;
6833
6834 announced.assert_next_some("test");
6835 announced.assert_next_wait();
6836 }
6837
6838 #[tokio::test]
6839 async fn test_coalesce_unannounce_announce_preserved() {
6840 tokio::time::pause();
6843
6844 let origin = Origin::random().produce();
6845 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6846 settle().await;
6847
6848 let mut announced = origin.consume().announced();
6849 announced.assert_next_some("test");
6850
6851 broadcast1.finish();
6853 settle().await;
6854
6855 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6856 settle().await;
6857
6858 announced.assert_next_none("test");
6860 announced.assert_next_some("test");
6861 announced.assert_next_wait();
6862 }
6863
6864 #[tokio::test]
6865 async fn test_coalesce_unannounce_announce_unannounce() {
6866 tokio::time::pause();
6869
6870 let origin = Origin::random().produce();
6871 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6872 settle().await;
6873
6874 let mut announced = origin.consume().announced();
6875 announced.assert_next_some("test");
6876
6877 broadcast1.finish();
6878 settle().await;
6879
6880 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6881 settle().await;
6882 broadcast2.finish();
6883 settle().await;
6884
6885 announced.assert_next_none("test");
6886 announced.assert_next_wait();
6887 }
6888
6889 #[tokio::test]
6890 async fn test_coalesce_churn_bounded() {
6891 tokio::time::pause();
6896
6897 let origin = Origin::random().produce();
6898 let mut announced = origin.consume().announced();
6899
6900 for _ in 0..1000 {
6901 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6902 settle().await;
6903 broadcast.finish();
6904 }
6905 settle().await;
6906
6907 let mut collected = Vec::new();
6908 while let Some(update) = announced.try_next() {
6909 collected.push(update);
6910 }
6911 assert!(
6912 collected.len() <= 1,
6913 "expected at most one pending update, got {}",
6914 collected.len()
6915 );
6916 assert!(
6917 collected.iter().all(|a| a.path == Path::new("test")),
6918 "unexpected path in pending updates",
6919 );
6920 }
6921
6922 #[tokio::test]
6926 async fn test_consumer_clone_is_side_effect_free() {
6927 let origin = Origin::random().produce();
6928
6929 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
6930 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
6931 settle().await;
6932
6933 let consumer = origin.consume();
6934 let mut announced = consumer.announced();
6935
6936 for _ in 0..16 {
6939 let cloned = consumer.clone();
6940 assert!(cloned.get_broadcast("test1").is_some());
6941 assert!(cloned.get_broadcast("test2").is_some());
6942 }
6943
6944 let a1 = announced.try_next().expect("first announcement");
6947 let a2 = announced.try_next().expect("second announcement");
6948 announced.assert_next_wait();
6949
6950 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6951 paths.sort();
6952 assert_eq!(paths, ["test1", "test2"]);
6953
6954 let mut fresh = consumer.announced();
6956 let b1 = fresh.try_next().expect("backlog: first");
6957 let b2 = fresh.try_next().expect("backlog: second");
6958 fresh.assert_next_wait();
6959
6960 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
6961 paths.sort();
6962 assert_eq!(paths, ["test1", "test2"]);
6963 }
6964
6965 #[tokio::test]
6967 async fn dynamic_request_unroutable_without_handler() {
6968 let origin = Origin::random().produce();
6969 let consumer = origin.consume();
6970 assert!(matches!(
6971 consumer.request_broadcast("missing").await,
6972 Err(Error::Unroutable)
6973 ));
6974 }
6975
6976 #[tokio::test(start_paused = true)]
6979 async fn dynamic_request_served_not_announced() {
6980 let origin = Origin::random().produce();
6981 let mut dynamic = origin.dynamic();
6982 let consumer = origin.consume();
6983
6984 let mut announced = origin.consume().announced();
6986 announced.assert_next_wait();
6987
6988 let served = broadcast::Info::new().produce();
6989 let request_fut = consumer.request_broadcast("fallback");
6992
6993 let mut served_dynamic = served.dynamic();
6995
6996 let request = dynamic.requested_broadcast().await.unwrap();
6997 assert_eq!(request.path(), &Path::new("fallback"));
6998 request.accept(&served);
6999
7000 let broadcast = request_fut.await.unwrap();
7001 assert!(broadcast.is_clone(&served.consume()));
7002
7003 let track_fut = broadcast.track("video").unwrap().subscribe(None);
7005 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
7006 let mut track = track_fut.await.unwrap();
7007 producer.append_group().unwrap();
7008 track.assert_group();
7009
7010 announced.assert_next_wait();
7012 }
7013
7014 #[tokio::test(start_paused = true)]
7016 async fn dynamic_request_coalesces() {
7017 let origin = Origin::random().produce();
7018 let mut dynamic = origin.dynamic();
7019 let consumer = origin.consume();
7020
7021 let f1 = consumer.request_broadcast("dup");
7023 let f2 = consumer.request_broadcast("dup");
7024
7025 let request = dynamic.requested_broadcast().await.unwrap();
7027 assert_eq!(request.path(), &Path::new("dup"));
7028 assert!(
7029 dynamic.requested_broadcast().now_or_never().is_none(),
7030 "a coalesced request must not be served twice"
7031 );
7032
7033 let served = broadcast::Info::new().produce();
7035 request.accept(&served);
7036 assert!(f1.await.unwrap().is_clone(&served.consume()));
7037 assert!(f2.await.unwrap().is_clone(&served.consume()));
7038 }
7039
7040 #[tokio::test(start_paused = true)]
7043 async fn dynamic_request_dedups_served() {
7044 let origin = Origin::random().produce();
7045 let mut dynamic = origin.dynamic();
7046 let consumer = origin.consume();
7047
7048 let request_fut = consumer.request_broadcast("fallback");
7049 let request = dynamic.requested_broadcast().await.unwrap();
7050 let served = broadcast::Info::new().produce();
7051 request.accept(&served);
7052 let first = request_fut.await.unwrap();
7053 assert!(first.is_clone(&served.consume()));
7054
7055 let second = consumer.request_broadcast("fallback").await.unwrap();
7057 assert!(second.is_clone(&served.consume()));
7058
7059 assert!(
7061 dynamic.requested_broadcast().now_or_never().is_none(),
7062 "a still-live served broadcast must not be re-requested from the handler"
7063 );
7064 }
7065
7066 #[tokio::test(start_paused = true)]
7068 async fn dynamic_request_reserves_after_close() {
7069 let origin = Origin::random().produce();
7070 let mut dynamic = origin.dynamic();
7071 let consumer = origin.consume();
7072
7073 let request_fut = consumer.request_broadcast("fallback");
7074 let request = dynamic.requested_broadcast().await.unwrap();
7075 let served = broadcast::Info::new().produce();
7076 request.accept(&served);
7077 request_fut.await.unwrap();
7078
7079 drop(served);
7081
7082 let request_fut = consumer.request_broadcast("fallback");
7084 let request = dynamic.requested_broadcast().await.unwrap();
7085 assert_eq!(request.path(), &Path::new("fallback"));
7086 let served = broadcast::Info::new().produce();
7087 request.accept(&served);
7088 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
7089 }
7090
7091 #[tokio::test(start_paused = true)]
7094 async fn dynamic_request_served_cache_bounded() {
7095 let origin = Origin::random().produce();
7096 let mut dynamic = origin.dynamic();
7097 let consumer = origin.consume();
7098
7099 for i in 0..100 {
7100 let path = format!("one-shot/{i}");
7101 let request_fut = consumer.request_broadcast(&path);
7102 let request = dynamic.requested_broadcast().await.unwrap();
7103 let served = broadcast::Info::new().produce();
7104 request.accept(&served);
7105 request_fut.await.unwrap();
7106 drop(served);
7108 }
7109
7110 assert!(
7113 origin.dynamic.read().served.len() <= 4,
7114 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
7115 origin.dynamic.read().served.len()
7116 );
7117 }
7118
7119 #[tokio::test(start_paused = true)]
7122 async fn dynamic_request_coalesces_after_handoff() {
7123 let origin = Origin::random().produce();
7124 let mut dynamic = origin.dynamic();
7125 let consumer = origin.consume();
7126
7127 let f1 = consumer.request_broadcast("fallback");
7128 let request = dynamic.requested_broadcast().await.unwrap();
7130
7131 let f2 = consumer.request_broadcast("fallback");
7133 assert!(
7134 dynamic.requested_broadcast().now_or_never().is_none(),
7135 "a repeat request during hand-off must coalesce, not re-queue"
7136 );
7137
7138 let served = broadcast::Info::new().produce();
7140 request.accept(&served);
7141 assert!(f1.await.unwrap().is_clone(&served.consume()));
7142 assert!(f2.await.unwrap().is_clone(&served.consume()));
7143 }
7144
7145 #[tokio::test(start_paused = true)]
7147 async fn dynamic_request_dropped_after_handoff() {
7148 let origin = Origin::random().produce();
7149 let mut dynamic = origin.dynamic();
7150 let consumer = origin.consume();
7151
7152 let f1 = consumer.request_broadcast("fallback");
7153 let request = dynamic.requested_broadcast().await.unwrap();
7154 let f2 = consumer.request_broadcast("fallback");
7155
7156 drop(request);
7158 assert!(matches!(f1.await, Err(Error::Unroutable)));
7159 assert!(matches!(f2.await, Err(Error::Unroutable)));
7160 }
7161
7162 #[tokio::test(start_paused = true)]
7164 async fn dynamic_request_rejected() {
7165 let origin = Origin::random().produce();
7166 let mut dynamic = origin.dynamic();
7167 let consumer = origin.consume();
7168
7169 let request_fut = consumer.request_broadcast("fallback");
7170
7171 let request = dynamic.requested_broadcast().await.unwrap();
7172 request.reject(Error::Cancel);
7173
7174 assert!(matches!(request_fut.await, Err(Error::Cancel)));
7175 }
7176
7177 #[tokio::test(start_paused = true)]
7181 async fn dynamic_request_rerequest_after_reject() {
7182 let origin = Origin::random().produce();
7183 let mut dynamic = origin.dynamic();
7184 let consumer = origin.consume();
7185
7186 let f1 = consumer.request_broadcast("fallback");
7187 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
7188 assert!(matches!(f1.await, Err(Error::Unroutable)));
7189
7190 let served = broadcast::Info::new().produce();
7191 let f2 = consumer.request_broadcast("fallback");
7193 let request = dynamic.requested_broadcast().await.unwrap();
7194 assert_eq!(request.path(), &Path::new("fallback"));
7195 request.accept(&served);
7196 assert!(f2.await.unwrap().is_clone(&served.consume()));
7197 }
7198
7199 #[tokio::test(start_paused = true)]
7202 async fn dynamic_request_handler_dropped() {
7203 let origin = Origin::random().produce();
7204 let dynamic = origin.dynamic();
7205 let consumer = origin.consume();
7206
7207 let request_fut = consumer.request_broadcast("fallback");
7208 drop(dynamic);
7209 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
7210
7211 assert!(matches!(
7213 consumer.request_broadcast("again").await,
7214 Err(Error::Unroutable)
7215 ));
7216 }
7217
7218 #[tokio::test(start_paused = true)]
7222 async fn dynamic_request_accept_after_handler_dropped() {
7223 let origin = Origin::random().produce();
7224 let mut dynamic = origin.dynamic();
7225 let consumer = origin.consume();
7226
7227 let request_fut = consumer.request_broadcast("fallback");
7228
7229 let request = dynamic.requested_broadcast().await.unwrap();
7231 drop(dynamic);
7232
7233 let served = broadcast::Info::new().produce();
7234 request.accept(&served);
7236 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
7237 }
7238
7239 #[tokio::test(start_paused = true)]
7241 async fn dynamic_request_prefers_announced() {
7242 let origin = Origin::random().produce();
7243 let mut dynamic = origin.dynamic();
7244 let consumer = origin.consume();
7245
7246 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
7247 settle().await;
7248
7249 let got = consumer.request_broadcast("live").await.unwrap();
7250 assert!(
7251 got.is_clone(&consumer.get_broadcast("live").unwrap()),
7252 "should return the published broadcast"
7253 );
7254 assert!(
7255 dynamic.requested_broadcast().now_or_never().is_none(),
7256 "a published path must not queue a fallback request"
7257 );
7258 }
7259
7260 #[tokio::test(start_paused = true)]
7262 async fn dynamic_clone_keeps_alive() {
7263 let origin = Origin::random().produce();
7264 let dynamic = origin.dynamic();
7265 let consumer = origin.consume();
7266
7267 drop(dynamic.clone());
7268
7269 let request_fut = consumer.request_broadcast("fallback");
7272 assert!(
7273 request_fut.now_or_never().is_none(),
7274 "request should stay pending until served"
7275 );
7276 }
7277
7278 fn wedge_watchdog<F>(name: &str, secs: u64, scenario: F)
7284 where
7285 F: std::future::Future<Output = ()> + Send + 'static,
7286 {
7287 let (done_tx, done_rx) = std::sync::mpsc::channel::<()>();
7288 let handle = std::thread::spawn(move || {
7289 let rt = ::tokio::runtime::Builder::new_current_thread()
7290 .enable_time()
7291 .build()
7292 .unwrap();
7293 rt.block_on(scenario);
7294 let _ = done_tx.send(());
7295 });
7296 match done_rx.recv_timeout(std::time::Duration::from_secs(secs)) {
7297 Ok(()) => {
7298 let _ = handle.join();
7299 }
7300 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
7301 let err = handle.join().unwrap_err();
7303 std::panic::resume_unwind(err);
7304 }
7305 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
7306 panic!("{name}: scenario wedged; a task is spinning inside a single poll")
7307 }
7308 }
7309 }
7310
7311 #[test]
7324 fn test_active_corpse_does_not_livelock_takeover() {
7325 wedge_watchdog("active-corpse", 20, async {
7326 let origin = Origin::random().produce();
7327 let consumer = origin.consume();
7328 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
7329
7330 let source_a = origin
7333 .create_broadcast("test", announce().with_hops(hops.clone()))
7334 .unwrap();
7335 let mut dynamic_a = source_a.dynamic();
7336 settle().await;
7337 settle().await;
7338
7339 let broadcast = consumer.request_broadcast("test").await.unwrap();
7340 let subscribing = broadcast.track("video").unwrap().subscribe(None);
7341 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
7342 settle().await;
7343 let mut sub = subscribing.await.unwrap();
7344 producer_a.append_group().unwrap();
7345 sub.assert_group();
7346
7347 let source_b = origin
7351 .create_broadcast("test", announce().with_hops(hops.clone()))
7352 .unwrap();
7353 let dynamic_b = source_b.dynamic();
7354 settle().await;
7355
7356 drop(dynamic_b);
7364 source_b.abort(Error::Dropped).unwrap();
7365 settle().await;
7366
7367 producer_a.append_group().unwrap();
7370 sub.assert_group();
7371 sub.assert_not_closed();
7372 });
7373 }
7374
7375 #[test]
7383 fn test_route_churn_never_wedges() {
7384 for seed in 1..=8u64 {
7385 wedge_watchdog(&format!("churn seed {seed}"), 30, churn_scenario(seed));
7386 }
7387 }
7388
7389 async fn churn_scenario(seed: u64) {
7390 let mut rng = seed.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
7392 let mut next = move || {
7393 rng = rng.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
7394 rng >> 33
7395 };
7396
7397 let origin = Origin::random().produce();
7398 let consumer = origin.consume();
7399 let names: Vec<Arc<str>> = (0..8).map(|i| Arc::from(format!("t{i}"))).collect();
7400
7401 let mut subs: Vec<track::Subscriber> = Vec::new();
7402 let mut pending_subs: Vec<kio::Pending<track::Subscribing>> = Vec::new();
7403
7404 struct Source {
7405 producer: Option<broadcast::Producer>,
7406 server: ::tokio::task::JoinHandle<()>,
7407 }
7408 let mut sources: Vec<Source> = Vec::new();
7409
7410 for step in 0..400u64 {
7411 match next() % 10 {
7412 0 | 1 => {
7415 if sources.len() >= 3 {
7416 continue;
7417 }
7418 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
7419 let route = announce().with_hops(hops).with_cost(next() % 4);
7420 let Ok(source) = origin.create_broadcast("test", route) else {
7421 continue;
7422 };
7423 let mut dynamic = source.dynamic();
7424 let behavior = next();
7425 let server = ::tokio::spawn(async move {
7426 let mut round = 0u64;
7427 let mut kept: Vec<track::Producer> = Vec::new();
7428 while let Ok(request) = dynamic.requested_track().await {
7429 round += 1;
7430 match (behavior >> (round % 16)) % 4 {
7431 0 => drop(request), 1 => {
7433 let mut producer = request.accept(None);
7434 let _ = producer.create_group(group::Info { sequence: round });
7435 let _ = producer.abort(Error::Dropped);
7436 }
7437 2 => {
7438 let mut producer = request.accept(None);
7439 let _ = producer.create_group(group::Info { sequence: round });
7440 kept.push(producer);
7441 }
7442 _ => {
7443 let mut producer = request.accept(None);
7444 let _ = producer.finish();
7445 }
7446 }
7447 }
7448 });
7449 sources.push(Source {
7450 producer: Some(source),
7451 server,
7452 });
7453 }
7454 2 | 3 => {
7456 if sources.is_empty() {
7457 continue;
7458 }
7459 let i = (next() as usize) % sources.len();
7460 let mut source = sources.swap_remove(i);
7461 if next() % 2 == 0
7462 && let Some(producer) = source.producer.take()
7463 {
7464 let _ = producer.abort(Error::Dropped);
7465 }
7466 source.server.abort();
7467 }
7468 4..=6 => {
7470 if subs.len() + pending_subs.len() >= 24 {
7471 continue;
7472 }
7473 let Some(broadcast) = consumer.get_broadcast("test") else {
7474 continue;
7475 };
7476 let name = &names[(next() as usize) % names.len()];
7477 if let Ok(track) = broadcast.track(name.as_ref()) {
7478 pending_subs.push(track.subscribe(None));
7479 }
7480 }
7481 7 => {
7483 if subs.is_empty() {
7484 continue;
7485 }
7486 let i = (next() as usize) % subs.len();
7487 subs.swap_remove(i);
7488 }
7489 _ => {
7491 for sub in pending_subs.drain(..) {
7492 match ::tokio::time::timeout(std::time::Duration::from_millis(5), sub).await {
7493 Ok(Ok(sub)) => subs.push(sub),
7494 Ok(Err(_)) => {}
7495 Err(_) => {}
7497 }
7498 }
7499 for sub in subs.iter_mut() {
7500 while let Some(Ok(Some(_))) = sub.recv_group().now_or_never() {}
7501 }
7502 }
7503 }
7504 if step % 16 == 0 {
7505 settle().await;
7506 }
7507 for _ in 0..(next() % 3) {
7509 ::tokio::task::yield_now().await;
7510 }
7511 }
7512
7513 for source in sources {
7514 source.server.abort();
7515 }
7516 settle().await;
7517 }
7518}