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 sources: usize,
670
671 nested: HashMap<String, Lock<OriginNode>>,
673
674 notify: Lock<NotifyNode>,
676}
677
678impl OriginNode {
679 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
680 Self {
681 broadcast: None,
682 sources: 0,
683 nested: HashMap::new(),
684 notify: Lock::new(NotifyNode::new(parent)),
685 }
686 }
687
688 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
689 let (dir, rest) = path.next_part().expect("leaf called with empty path");
690
691 let next = self.entry(dir);
692 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
693 }
694
695 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
696 match self.nested.get(dir) {
697 Some(next) => next.clone(),
698 None => {
699 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
700 self.nested.insert(dir.to_string(), next.clone());
701 next
702 }
703 }
704 }
705
706 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
710 let Some(existing) = &mut self.broadcast else { return };
711 if !existing.state.same_channel(expect) || existing.announced == announce {
712 return;
713 }
714 existing.announced = announce;
715 let path = existing.path.clone();
716 let consumer = existing.broadcast.consume();
717 let state = existing.state.clone();
718 let mut notify = self.notify.lock();
719 if announce {
720 notify.announce(&path, &consumer, &state);
721 } else {
722 notify.unannounce(&path);
723 }
724 }
725
726 fn consume_at(&mut self, id: ConsumerId, notify: AnnounceConsumerNotify, relative: impl AsPath) {
733 let relative = relative.as_path();
734
735 let Some((dir, relative)) = relative.next_part() else {
736 return self.consume(id, notify);
737 };
738
739 let nested = self.entry(dir);
740 nested.lock().consume_at(id, notify, &relative);
741 }
742
743 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
744 self.consume_initial(&mut notify);
745 self.notify.lock().consumers.insert(id, notify);
746 }
747
748 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
749 if let Some(broadcast) = &self.broadcast
752 && broadcast.announced
753 {
754 notify.announce(&broadcast.path, broadcast.broadcast.consume(), &broadcast.state);
755 }
756
757 for nested in self.nested.values() {
759 nested.lock().consume_initial(notify);
760 }
761 }
762
763 fn resolve_broadcast(&self, rest: impl AsPath, exclude: Option<Origin>) -> Resolved {
764 let rest = rest.as_path();
765
766 if let Some((dir, rest)) = rest.next_part() {
767 let Some(node) = self.nested.get(dir) else {
768 return Resolved::Missing;
769 };
770 let node = node.lock();
771 return node.resolve_broadcast(&rest, exclude);
772 }
773
774 let Some(broadcast) = self.broadcast.as_ref() else {
775 return Resolved::Missing;
776 };
777 let Some(origin) = exclude else {
778 return Resolved::Found(broadcast.broadcast.consume());
779 };
780
781 let state = broadcast.state.read();
797 if !state.routes.iter().any(|r| r.route.hops.contains(&origin)) {
798 drop(state);
799 let shared = broadcast.broadcast.consume();
800 return match ExclusionGuard::new(&broadcast.state, origin) {
801 Some(guard) => Resolved::Found(shared.with_exclusion(guard)),
802 None => Resolved::Found(shared),
805 };
806 }
807 match state
808 .dispatch(Some(origin))
809 .and_then(|clean| state.routes.iter().find(|r| r.id == clean))
810 {
811 Some(route) => Resolved::Found(route.source.clone()),
812 None => Resolved::Excluded,
813 }
814 }
815
816 fn detach(&mut self, claim: Claim, relative: impl AsPath) {
824 let relative = relative.as_path();
825
826 let Some((dir, relative)) = relative.next_part() else {
827 match claim {
828 Claim::Consumer(id) => {
829 self.notify.lock().consumers.remove(&id).expect("consumer not found");
830 }
831 Claim::Source => self.sources -= 1,
832 }
833 return;
834 };
835
836 let nested = self.nested.get(dir).expect("claimed node missing").clone();
839 let mut locked = nested.lock();
840 locked.detach(claim, &relative);
841
842 if locked.is_empty() {
843 drop(locked);
844 self.nested.remove(dir);
845 }
846 }
847
848 fn reserve(&mut self, relative: impl AsPath) {
851 let relative = relative.as_path();
852
853 let Some((dir, relative)) = relative.next_part() else {
854 self.sources += 1;
855 return;
856 };
857
858 let nested = self.entry(dir);
859 nested.lock().reserve(&relative);
860 }
861
862 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
866 let relative = relative.as_path();
867
868 if let Some((dir, relative)) = relative.next_part() {
869 let Some(nested) = self.nested.get(dir) else { return };
870 let nested = nested.clone();
871 let mut locked = nested.lock();
872 locked.remove(expect, &relative);
873
874 if locked.is_empty() {
875 drop(locked);
876 self.nested.remove(dir);
877 }
878 } else if let Some(existing) = &self.broadcast
879 && existing.state.same_channel(expect)
880 {
881 let existing = self.broadcast.take().expect("checked above");
882 if existing.announced {
883 self.notify.lock().unannounce(&existing.path);
884 }
885 }
886 }
887
888 fn is_empty(&self) -> bool {
889 self.broadcast.is_none()
890 && self.sources == 0
891 && self.nested.is_empty()
892 && self.notify.lock().consumers.is_empty()
893 }
894
895 #[cfg(test)]
898 fn count(&self) -> usize {
899 1 + self.nested.values().map(|nested| nested.lock().count()).sum::<usize>()
900 }
901}
902
903#[derive(Clone, Copy)]
906enum Claim {
907 Consumer(ConsumerId),
909 Source,
911}
912
913struct SourceReservation {
925 tree: Lock<OriginNode>,
926 path: PathOwned,
927}
928
929impl SourceReservation {
930 fn new(tree: Lock<OriginNode>, path: PathOwned) -> Self {
931 tree.lock().reserve(&path);
932 Self { tree, path }
933 }
934}
935
936impl Drop for SourceReservation {
937 fn drop(&mut self) {
938 self.tree.lock().detach(Claim::Source, &self.path);
939 }
940}
941
942#[derive(Clone)]
950struct OriginNodes {
951 tree: Lock<OriginNode>,
954
955 nodes: Vec<(PathOwned, PathOwned)>,
958}
959
960impl OriginNodes {
961 fn empty() -> Self {
964 Self {
965 tree: Lock::new(OriginNode::new(None)),
966 nodes: Vec::new(),
967 }
968 }
969
970 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
973 let mut roots = Vec::new();
974
975 for (root, absolute) in &self.nodes {
976 for prefix in prefixes {
977 if root.has_prefix(prefix) {
978 roots.push((root.to_owned(), absolute.clone()));
980 continue;
981 }
982
983 if let Some(suffix) = prefix.strip_prefix(root) {
984 roots.push((prefix.to_owned(), absolute.join(&suffix)));
986 }
987 }
988 }
989
990 if roots.is_empty() {
991 None
992 } else {
993 Some(self.with_nodes(roots))
994 }
995 }
996
997 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
998 let new_root = new_root.as_path();
999 let mut roots = Vec::new();
1000
1001 if new_root.is_empty() {
1002 return Some(self.clone());
1003 }
1004
1005 for (root, absolute) in &self.nodes {
1006 if let Some(suffix) = root.strip_prefix(&new_root) {
1007 roots.push((suffix.to_owned(), absolute.clone()));
1009 } else if let Some(suffix) = new_root.strip_prefix(root) {
1010 roots.push(("".into(), absolute.join(&suffix)));
1013 }
1014 }
1015
1016 if roots.is_empty() {
1017 None
1018 } else {
1019 Some(self.with_nodes(roots))
1020 }
1021 }
1022
1023 fn with_nodes(&self, nodes: Vec<(PathOwned, PathOwned)>) -> Self {
1024 Self {
1025 tree: self.tree.clone(),
1026 nodes,
1027 }
1028 }
1029
1030 pub fn get(&self, path: impl AsPath) -> Option<PathOwned> {
1032 let path = path.as_path();
1033
1034 for (root, absolute) in &self.nodes {
1035 if let Some(suffix) = path.strip_prefix(root) {
1036 return Some(absolute.join(&suffix));
1037 }
1038 }
1039
1040 None
1041 }
1042}
1043
1044impl Default for OriginNodes {
1045 fn default() -> Self {
1046 Self {
1047 tree: Lock::new(OriginNode::new(None)),
1048 nodes: vec![("".into(), "".into())],
1049 }
1050 }
1051}
1052
1053#[derive(Clone)]
1055pub struct OriginAnnounce {
1056 pub path: PathOwned,
1058 pub broadcast: Option<broadcast::Consumer>,
1064}
1065
1066#[derive(Clone)]
1068pub struct Producer {
1069 info: Origin,
1073
1074 nodes: OriginNodes,
1077
1078 root: PathOwned,
1080
1081 dynamic: kio::Shared<OriginDynamicState>,
1085
1086 pool: cache::Pool,
1089
1090 cache_duration: Duration,
1093
1094 latency_default: Duration,
1097
1098 linger: Duration,
1101
1102 stats: stats::Session,
1106}
1107
1108impl std::ops::Deref for Producer {
1109 type Target = Origin;
1110
1111 fn deref(&self) -> &Self::Target {
1112 &self.info
1113 }
1114}
1115
1116impl Producer {
1117 pub fn new(info: Info) -> Self {
1121 Self {
1122 info: info.id,
1123 nodes: OriginNodes::default(),
1124 root: PathOwned::default(),
1125 dynamic: kio::Shared::default(),
1126 pool: info.pool,
1127 cache_duration: info.cache_duration,
1128 latency_default: info.latency_default,
1129 linger: info.linger,
1130 stats: stats::Session::default(),
1131 }
1132 }
1133
1134 pub fn with_stats(mut self, session: stats::Session) -> Self {
1138 self.stats = session;
1139 self
1140 }
1141
1142 pub fn with_linger(mut self, linger: Duration) -> Self {
1151 self.linger = linger;
1152 self
1153 }
1154
1155 pub fn info(&self) -> Info {
1158 Info {
1159 id: self.info,
1160 pool: self.pool.clone(),
1161 cache_duration: self.cache_duration,
1162 latency_default: self.latency_default,
1163 linger: self.linger,
1164 }
1165 }
1166
1167 pub(crate) fn latency_default(&self) -> Duration {
1170 self.latency_default
1171 }
1172
1173 pub(crate) fn empty(info: Origin) -> Self {
1178 Self {
1179 info,
1180 nodes: OriginNodes::empty(),
1181 root: PathOwned::default(),
1182 dynamic: kio::Shared::default(),
1183 pool: cache::Pool::default(),
1184 cache_duration: Duration::MAX,
1185 latency_default: track::DEFAULT_LATENCY_MAX,
1186 linger: Duration::ZERO,
1187 stats: stats::Session::default(),
1188 }
1189 }
1190
1191 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
1242 let path = path.as_path();
1243
1244 debug_assert!(
1245 !route.hops.contains(&self.info),
1246 "create_broadcast called with a looping hop chain",
1247 );
1248
1249 let full = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
1254 let tree = self.nodes.tree.clone();
1255
1256 if full.parts().count() > Path::MAX_PARTS {
1260 return Err(BoundsExceeded.into());
1261 }
1262
1263 let ingress = self.stats.ingress(&full);
1267
1268 let mut source = broadcast::Info { origin: self.info() }
1269 .produce()
1270 .with_stats(ingress.clone());
1271 source.set_route(route).expect("fresh producer");
1272
1273 let reservation = SourceReservation::new(tree.clone(), full.clone());
1276 web_async::spawn(run_source(
1277 self.info(),
1278 tree,
1279 full,
1280 source.consume(),
1281 ingress,
1282 reservation,
1283 ));
1284
1285 Ok(source)
1286 }
1287
1288 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
1294 let prefixes = PathPrefixes::new(prefixes);
1295 Some(Producer {
1296 info: self.info,
1297 nodes: self.nodes.select(&prefixes)?,
1298 root: self.root.clone(),
1299 dynamic: self.dynamic.clone(),
1300 pool: self.pool.clone(),
1301 cache_duration: self.cache_duration,
1302 latency_default: self.latency_default,
1303 linger: self.linger,
1304 stats: self.stats.clone(),
1305 })
1306 }
1307
1308 pub fn dynamic(&self) -> Dynamic {
1317 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
1318 }
1319
1320 pub fn consume(&self) -> Consumer {
1325 Consumer::new(
1328 self.info,
1329 self.root.clone(),
1330 self.nodes.clone(),
1331 self.dynamic.clone(),
1332 stats::Session::default(),
1333 )
1334 }
1335
1336 pub fn announces(&self) -> AnnounceProducer {
1342 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
1343 }
1344
1345 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1350 let prefix = prefix.as_path();
1351
1352 Some(Self {
1353 info: self.info,
1354 root: self.root.join(&prefix).to_owned(),
1355 nodes: self.nodes.root(&prefix)?,
1356 dynamic: self.dynamic.clone(),
1357 pool: self.pool.clone(),
1358 cache_duration: self.cache_duration,
1359 latency_default: self.latency_default,
1360 linger: self.linger,
1361 stats: self.stats.clone(),
1362 })
1363 }
1364
1365 pub fn root(&self) -> &Path<'_> {
1367 &self.root
1368 }
1369
1370 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1373 self.nodes.nodes.iter().map(|(root, _)| root)
1374 }
1375
1376 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1378 self.root.join(path)
1379 }
1380
1381 #[cfg(test)]
1383 pub(crate) fn node_count(&self) -> usize {
1384 self.nodes.tree.lock().count()
1385 }
1386}
1387
1388const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1401
1402struct FrontRoute {
1404 id: u64,
1405 route: broadcast::Route,
1408 source: broadcast::Consumer,
1410}
1411
1412struct FrontState {
1414 path: PathOwned,
1416 self_origin: Origin,
1418 publisher: Option<Origin>,
1426 next_route: u64,
1429 routes: Vec<FrontRoute>,
1430 excluded: HashMap<Origin, usize>,
1443 active: Option<u64>,
1445 linger: Duration,
1448 closed: bool,
1453}
1454
1455impl FrontState {
1456 fn pick(&self, keep: impl Fn(&FrontRoute) -> bool, untainted: bool) -> Option<u64> {
1463 let candidates: Vec<&FrontRoute> = self.routes.iter().filter(|r| keep(r)).collect();
1464 let candidates = match untainted {
1465 true => self.prefer_untainted(&candidates),
1466 false => candidates,
1467 };
1468 candidates
1469 .into_iter()
1470 .min_by_key(|r| route_order(&self.path.as_path(), r))
1471 .map(|r| r.id)
1472 }
1473
1474 fn best_route(&self) -> Option<u64> {
1479 self.pick(|_| true, true)
1480 }
1481
1482 fn dispatch(&self, exclude: Option<Origin>) -> Option<u64> {
1489 self.pick(|r| exclude.is_none_or(|origin| !r.route.hops.contains(&origin)), false)
1490 }
1491
1492 fn taints_a_reader(&self, route: &broadcast::Route) -> bool {
1499 route.hops.iter().any(|hop| self.excluded.contains_key(hop))
1500 }
1501
1502 fn prefer_untainted<'a>(&self, candidates: &[&'a FrontRoute]) -> Vec<&'a FrontRoute> {
1512 if self.excluded.is_empty() {
1513 return candidates.to_vec();
1514 }
1515 let clean: Vec<&FrontRoute> = candidates
1516 .iter()
1517 .copied()
1518 .filter(|r| !self.taints_a_reader(&r.route))
1519 .collect();
1520 match clean.is_empty() {
1521 true => candidates.to_vec(),
1522 false => clean,
1523 }
1524 }
1525
1526 fn serve_route(&self, skip: impl Fn(u64) -> bool) -> Option<u64> {
1535 if let Some(active) = self.active
1536 && !skip(active)
1537 && let Some(route) = self.routes.iter().find(|r| r.id == active)
1538 && !self.taints_a_reader(&route.route)
1539 {
1540 return Some(active);
1541 }
1542 self.pick(|r| !skip(r.id), true)
1543 }
1544
1545 fn reselect(&mut self, carrying: bool) {
1563 let best = self.best_route();
1564 if carrying
1565 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1566 && best_id != cur_id
1567 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1568 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1569 && incumbent.route.announce
1570 && candidate.route.cost < incumbent.route.cost
1571 && candidate.route.advertised == 0
1572 && candidate.route.hops.len() >= 2
1573 && !self.handover_allowed(&candidate.route)
1574 {
1575 return;
1577 }
1578 self.active = best;
1579 }
1580
1581 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1593 let name = self.path.as_path();
1594 match route.hops.iter().last() {
1595 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1596 None => true,
1597 }
1598 }
1599
1600 fn routes_snapshot(&self) -> Vec<broadcast::Route> {
1606 let mut routes: Vec<&FrontRoute> = self.routes.iter().collect();
1607 routes.sort_by_key(|r| route_order(&self.path.as_path(), r));
1608 routes.sort_by_key(|r| Some(r.id) != self.active);
1609 routes.into_iter().map(|r| r.route.clone()).collect()
1610 }
1611}
1612
1613fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1623 let mut leaf_guard = leaf.lock();
1628 let routes = state.read().routes_snapshot();
1629 if let Some(advert) = routes.first() {
1630 let announce = advert.announce;
1631 broadcast.clone().set_routes(routes);
1632 leaf_guard.set_announced(state, announce);
1633 }
1634}
1635
1636fn detach_source(
1649 state: &kio::Producer<FrontState>,
1650 broadcast: &broadcast::Producer,
1651 leaf: &Lock<OriginNode>,
1652 id: u64,
1653 graceful: bool,
1654) {
1655 let close = {
1656 let carrying = broadcast.demand().is_used();
1660 let Ok(mut s) = state.write() else { return };
1661 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1662 return;
1663 };
1664 s.routes.remove(pos);
1665 s.reselect(carrying);
1666 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1667 s.closed = true;
1670 true
1671 } else {
1672 false
1673 }
1674 };
1675 if close {
1676 broadcast.abort_spliced(Error::Dropped);
1677 }
1678 sync_front(state, broadcast, leaf);
1679}
1680
1681fn sync_announce(guard: &mut Option<stats::Announce>, announced: bool, ingress: &stats::Scope) {
1684 match (announced, guard.is_some()) {
1685 (true, false) => *guard = Some(ingress.announce()),
1686 (false, true) => *guard = None,
1687 _ => {}
1688 }
1689}
1690
1691async fn run_source(
1701 origin: Info,
1702 tree: Lock<OriginNode>,
1703 full: PathOwned,
1704 mut source: broadcast::Consumer,
1705 ingress: stats::Scope,
1706 _reservation: SourceReservation,
1708) {
1709 let ctx = AttachContext {
1710 origin: &origin,
1711 tree: &tree,
1712 full: &full,
1713 };
1714
1715 let Ok(mut route) = source.route_changed().await else {
1719 return;
1721 };
1722
1723 let mut announce = route.announce.then(|| ingress.announce());
1728 let mut may_take_over = true;
1734
1735 let leaf = if full.is_empty() {
1739 tree.clone()
1740 } else {
1741 tree.lock().leaf(&full)
1742 };
1743
1744 'attach: loop {
1745 let (state, broadcast, id) = match attach_source(&ctx, &leaf, &source, route.clone(), may_take_over) {
1746 Attach::Ready(state, broadcast, id) => (state, broadcast, id),
1747 Attach::Parked(incumbent) => {
1748 tracing::debug!(
1749 broadcast = %full,
1750 "path already live with a different publisher; parking this source until it ends",
1751 );
1752 let update = kio::wait(|waiter| {
1756 if let Poll::Ready(update) = source.poll_route_changed(waiter) {
1757 return Poll::Ready(Some(update));
1758 }
1759 match incumbent.poll(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending }) {
1762 Poll::Ready(_) => Poll::Ready(None),
1763 Poll::Pending => Poll::Pending,
1764 }
1765 })
1766 .await;
1767 match update {
1768 Some(Ok(update)) => {
1770 sync_announce(&mut announce, update.announce, &ingress);
1771 if update.hops.iter().next().copied() != route.hops.iter().next().copied() {
1776 may_take_over = true;
1777 }
1778 route = update;
1779 }
1780 Some(Err(_)) => return,
1782 None => {}
1784 }
1785 continue 'attach;
1786 }
1787 };
1788 let publisher = route.hops.iter().next().copied();
1789
1790 loop {
1791 let update = kio::wait(|waiter| {
1792 if state
1798 .poll_ref(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending })
1799 .is_ready()
1800 {
1801 return Poll::Ready(None);
1802 }
1803 source.poll_route_changed(waiter).map(Some)
1804 })
1805 .await;
1806 match update {
1807 None => {
1808 may_take_over = false;
1811 continue 'attach;
1812 }
1813 Some(Ok(update)) => {
1814 let announced = update.announce;
1815 if update.hops.iter().next().copied() != publisher {
1826 detach_source(&state, &broadcast, &leaf, id, true);
1827 sync_announce(&mut announce, announced, &ingress);
1828 may_take_over = true;
1832 route = update;
1833 continue 'attach;
1834 }
1835 {
1836 let carrying = broadcast.demand().is_used();
1837 let Ok(mut s) = state.write() else { return };
1838 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1839 return;
1840 };
1841 if entry.route == update {
1842 continue;
1843 }
1844 entry.route = update;
1845 s.reselect(carrying);
1846 }
1847 sync_announce(&mut announce, announced, &ingress);
1849 sync_front(&state, &broadcast, &leaf);
1850 }
1851 Some(Err(_)) => {
1852 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1855 return;
1856 }
1857 }
1858 }
1859 }
1860}
1861
1862enum Attach {
1864 Ready(kio::Producer<FrontState>, broadcast::Producer, u64),
1867 Parked(kio::Producer<FrontState>),
1874}
1875
1876struct AttachContext<'a> {
1878 origin: &'a Info,
1879 tree: &'a Lock<OriginNode>,
1882 full: &'a PathOwned,
1885}
1886
1887fn same_publisher(a: Option<Origin>, b: Option<Origin>) -> bool {
1896 if a == Some(Origin::UNKNOWN) || b == Some(Origin::UNKNOWN) {
1897 return false;
1898 }
1899 a == b
1900}
1901
1902fn attach_source(
1928 ctx: &AttachContext,
1929 leaf: &Lock<OriginNode>,
1930 source: &broadcast::Consumer,
1931 route: broadcast::Route,
1932 may_take_over: bool,
1933) -> Attach {
1934 let publisher = route.hops.iter().next().copied();
1935 let mut leaf_guard = leaf.lock();
1936
1937 if let Some(existing) = &leaf_guard.broadcast {
1940 let mut joined = None;
1941 let carrying = existing.broadcast.demand().is_used();
1942 if let Ok(mut s) = existing.state.write()
1943 && !s.closed
1944 {
1945 if same_publisher(s.publisher, publisher) {
1946 let id = s.next_route;
1947 s.next_route += 1;
1948 s.routes.push(FrontRoute {
1949 id,
1950 route: route.clone(),
1951 source: source.clone(),
1952 });
1953 s.reselect(carrying);
1954 joined = Some(id);
1955 } else if !may_take_over || !route.announce || s.taints_a_reader(&route) {
1956 return Attach::Parked(existing.state.clone());
1957 } else {
1958 s.closed = true;
1966 tracing::warn!(broadcast = %ctx.full, "replacing a live broadcast from a different publisher");
1967 }
1968 }
1969 if let Some(id) = joined {
1970 let state = existing.state.clone();
1971 let broadcast = existing.broadcast.clone();
1972 drop(leaf_guard);
1973 sync_front(&state, &broadcast, leaf);
1974 return Attach::Ready(state, broadcast, id);
1975 }
1976 }
1977
1978 let announce = route.announce;
1980 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
1981 origin: ctx.origin.clone(),
1982 });
1983 let _ = broadcast.clone().set_route(route.clone());
1984 let state = kio::Producer::new(FrontState {
1985 path: ctx.full.clone(),
1986 self_origin: ctx.origin.id,
1987 publisher,
1988 next_route: 1,
1989 excluded: HashMap::new(),
1990 routes: vec![FrontRoute {
1991 id: 0,
1992 route,
1993 source: source.clone(),
1994 }],
1995 active: Some(0),
1996 linger: ctx.origin.linger,
1997 closed: false,
1998 });
1999
2000 if let Some(stale) = leaf_guard.broadcast.take()
2004 && stale.announced
2005 {
2006 leaf_guard.notify.lock().unannounce(&stale.path);
2007 }
2008 let entry = OriginBroadcast {
2009 path: ctx.full.clone(),
2010 broadcast: broadcast.clone(),
2011 state: state.clone(),
2012 announced: announce,
2013 };
2014 if entry.announced {
2015 leaf_guard
2016 .notify
2017 .lock()
2018 .announce(ctx.full, &broadcast.consume(), &state);
2019 }
2020 leaf_guard.broadcast = Some(entry);
2021 drop(leaf_guard);
2022
2023 web_async::spawn(run_front(
2024 state.clone(),
2025 broadcast.clone(),
2026 ctx.tree.clone(),
2027 ctx.full.clone(),
2028 ));
2029
2030 Attach::Ready(state, broadcast, 0)
2031}
2032
2033async fn run_front(
2036 state: kio::Producer<FrontState>,
2037 mut broadcast: broadcast::Producer,
2038 tree: Lock<OriginNode>,
2039 full: PathOwned,
2040) {
2041 enum Step {
2042 Serve(Arc<str>, super::resume::Producer),
2043 Changed,
2045 Expired,
2047 Closed,
2048 }
2049
2050 let linger = state.read().linger;
2051 let mut deadline = kio::time::Deadline::new();
2056
2057 loop {
2058 let empty = {
2059 let s = state.read();
2060 !s.closed && s.routes.is_empty()
2061 };
2062 deadline.set(match (empty, deadline.deadline()) {
2063 (true, None) => web_async::time::Instant::now().checked_add(linger),
2066 (true, at) => at,
2067 (false, _) => None,
2068 });
2069
2070 let step = {
2071 kio::wait(|waiter| {
2072 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
2073 return Poll::Ready(Step::Serve(name, resume));
2074 }
2075 match state.poll(waiter, |s| {
2078 if s.closed || s.routes.is_empty() != empty {
2079 Poll::Ready(())
2080 } else {
2081 Poll::Pending
2082 }
2083 }) {
2084 Poll::Ready(Ok(guard)) => {
2085 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
2086 }
2087 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
2088 Poll::Pending => {}
2089 }
2090 deadline.poll(waiter).map(|_| Step::Expired)
2091 })
2092 .await
2093 };
2094
2095 match step {
2096 Step::Serve(name, resume) => {
2097 web_async::spawn(serve_track(state.clone(), name, resume));
2100 }
2101 Step::Changed => {}
2102 Step::Expired => {
2103 let close = {
2107 let Ok(mut s) = state.write() else { break };
2108 if !s.closed && s.routes.is_empty() {
2109 s.closed = true;
2110 true
2111 } else {
2112 false
2113 }
2114 };
2115 if close {
2116 break;
2117 }
2118 }
2119 Step::Closed => break,
2120 }
2121 }
2122
2123 broadcast.abort_spliced(Error::Dropped);
2125
2126 broadcast.finish();
2128
2129 tree.lock().remove(&state, &full);
2132}
2133
2134async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
2147 enum Step {
2148 Closed,
2149 Splice(u64, broadcast::Consumer),
2150 Complete,
2151 Failed(Error),
2152 NoRoute,
2157 Idle,
2159 Demand,
2161 }
2162
2163 let mut serving: Option<(u64, track::Consumer)> = None;
2165 let mut spliced_edge: Option<u64> = None;
2171 let mut refused: HashSet<u64> = HashSet::new();
2176 let mut refusal: Option<Error> = None;
2177 let mut dead: HashSet<u64> = HashSet::new();
2182 let mut idle_since: Option<web_async::time::Instant> = None;
2184 let mut deadline = kio::time::Deadline::new();
2185
2186 loop {
2187 let serving_id = serving.as_ref().map(|(id, _)| *id);
2188
2189 {
2196 let s = state.read();
2197 refused.retain(|id| s.routes.iter().any(|r| r.id == *id));
2198 dead.retain(|id| s.routes.iter().any(|r| r.id == *id));
2199 let exhausted = !s.routes.is_empty()
2200 && s.serve_route(|id| refused.contains(&id) || dead.contains(&id))
2201 .is_none();
2202 if exhausted && dead.is_empty() {
2203 drop(s);
2204 let err = refusal.take().unwrap_or(Error::NotFound);
2205 tracing::debug!(name = %name, %err, "every source refused track; aborting");
2206 let _ = resume.abort(err);
2207 return;
2208 }
2209 }
2210
2211 let used = resume.is_used();
2223 idle_since = match (resume.is_spliced(), used) {
2224 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
2225 _ => None,
2226 };
2227 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
2228
2229 let step = {
2230 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2231 kio::wait(|waiter| {
2232 match state.poll(waiter, |s| {
2238 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
2239 if s.closed
2240 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
2241 {
2242 Poll::Ready(())
2243 } else {
2244 Poll::Pending
2245 }
2246 }) {
2247 Poll::Ready(Ok(guard)) => {
2248 if guard.closed {
2249 return Poll::Ready(Step::Closed);
2250 }
2251 let Some(next) = guard.serve_route(skip) else {
2252 return Poll::Ready(Step::NoRoute);
2253 };
2254 let source = guard
2255 .routes
2256 .iter()
2257 .find(|r| r.id == next)
2258 .expect("servable source in table")
2259 .source
2260 .clone();
2261 return Poll::Ready(Step::Splice(next, source));
2262 }
2263 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
2264 Poll::Pending => {}
2265 }
2266
2267 let edge = match used {
2272 true => resume.poll_unused(waiter),
2273 false => resume.poll_used(waiter),
2274 };
2275 if edge.is_ready() {
2276 return Poll::Ready(Step::Demand);
2277 }
2278
2279 if let Some((_, track)) = &serving
2282 && let Poll::Ready(result) = track.poll_complete(waiter)
2283 {
2284 return Poll::Ready(match result {
2285 Ok(()) => Step::Complete,
2286 Err(err) => Step::Failed(err),
2287 });
2288 }
2289
2290 deadline.poll(waiter).map(|_| Step::Idle)
2291 })
2292 .await
2293 };
2294
2295 match step {
2296 Step::Closed => return,
2298 Step::Complete => {
2299 let _ = resume.finish();
2300 return;
2301 }
2302 Step::Failed(err) => {
2303 if resume.latest() == spliced_edge
2311 && let Some(id) = serving_id
2312 {
2313 let closing = state
2314 .read()
2315 .routes
2316 .iter()
2317 .find(|r| r.id == id)
2318 .is_some_and(|r| r.source.is_closing());
2319 if closing {
2320 dead.insert(id);
2321 } else {
2322 refused.insert(id);
2323 refusal = Some(err);
2324 }
2325 }
2326 serving = None;
2327 }
2328 Step::Demand => {}
2330 Step::NoRoute => serving = None,
2336 Step::Idle => {
2337 if resume.release().is_err() {
2342 return;
2344 }
2345 serving = None;
2346 }
2347 Step::Splice(id, source) => {
2348 let attempt = match source.track(&name) {
2352 Ok(track) => {
2353 let query = track.info().into_inner();
2356 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2357 let info = kio::wait(|waiter| {
2358 if let Poll::Ready(result) = query.poll(waiter) {
2359 return Poll::Ready(Some(result));
2360 }
2361 match state.poll(waiter, |s| {
2362 if s.closed || s.serve_route(skip) != Some(id) {
2363 Poll::Ready(())
2364 } else {
2365 Poll::Pending
2366 }
2367 }) {
2368 Poll::Ready(_) => Poll::Ready(None),
2369 Poll::Pending => Poll::Pending,
2370 }
2371 })
2372 .await;
2373 match info {
2374 None => continue,
2376 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2379 Poll::Ready(Err(err)) => Err(err),
2380 _ => Ok(track),
2381 },
2382 Some(Err(err)) => Err(err),
2383 }
2384 }
2385 Err(err) => Err(err),
2386 };
2387
2388 match attempt {
2389 Ok(track) => {
2390 if let Err(err) = resume.takeover(&track) {
2391 let _ = resume.abort(err);
2396 return;
2397 }
2398 spliced_edge = resume.latest();
2410 serving = Some((id, track));
2411 }
2412 Err(_) if source.is_closing() => {
2416 dead.insert(id);
2417 serving = None;
2418 }
2419 Err(err) => {
2425 tracing::debug!(name = %name, source = id, %err, "source refused track");
2426 refused.insert(id);
2427 refusal = Some(err);
2428 serving = None;
2429 }
2430 }
2431 }
2432 }
2433 }
2434}
2435
2436#[derive(Default)]
2442struct OriginDynamicState {
2443 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2446
2447 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2453}
2454
2455#[derive(Default)]
2462struct PendingBroadcast {
2463 resolved: Option<Result<broadcast::Consumer, Error>>,
2464}
2465
2466pub struct Dynamic {
2477 info: Origin,
2478 root: PathOwned,
2479 state: kio::Shared<OriginDynamicState>,
2480}
2481
2482impl Clone for Dynamic {
2483 fn clone(&self) -> Self {
2484 self.state.lock().requests.add_handler();
2488
2489 Self {
2490 info: self.info,
2491 root: self.root.clone(),
2492 state: self.state.clone(),
2493 }
2494 }
2495}
2496
2497impl Dynamic {
2498 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2499 state.lock().requests.add_handler();
2500
2501 Self { info, root, state }
2502 }
2503
2504 pub fn info(&self) -> &Origin {
2506 &self.info
2507 }
2508
2509 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2511 let mut state = ready!(self.state.poll(waiter, |state| {
2512 if state.requests.has_queued() {
2513 Poll::Ready(())
2514 } else {
2515 Poll::Pending
2516 }
2517 }));
2518
2519 let path = state.requests.pop().expect("predicate guaranteed a request");
2520 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2526 Poll::Ready(Ok(Request {
2527 path,
2528 producer,
2529 state: self.state.clone(),
2530 }))
2531 }
2532
2533 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2536 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2537 }
2538
2539 pub fn root(&self) -> &Path<'_> {
2541 &self.root
2542 }
2543}
2544
2545impl Drop for Dynamic {
2546 fn drop(&mut self) {
2547 let mut state = self.state.lock();
2550 if state.requests.remove_handler() {
2551 state.requests.drain_queued();
2555 }
2556 }
2557}
2558
2559pub struct Request {
2566 path: PathOwned,
2568
2569 producer: kio::Producer<PendingBroadcast>,
2572
2573 state: kio::Shared<OriginDynamicState>,
2575}
2576
2577impl Request {
2578 pub fn path(&self) -> &Path<'_> {
2580 &self.path
2581 }
2582
2583 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2589 let broadcast = broadcast.consume();
2590
2591 let resolved = {
2597 let mut state = self.state.lock();
2598 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2599 state
2600 .requests
2601 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2602 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2603 };
2604
2605 if let Ok(mut pending) = self.producer.write() {
2606 pending.resolved = Some(Ok(resolved));
2607 }
2608 }
2610
2611 pub fn reject(self, err: Error) {
2613 self.state
2614 .lock()
2615 .requests
2616 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2617 if let Ok(mut state) = self.producer.write() {
2618 state.resolved = Some(Err(err));
2619 }
2620 }
2621}
2622
2623impl Drop for Request {
2624 fn drop(&mut self) {
2625 self.state
2633 .lock()
2634 .requests
2635 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2636 }
2637}
2638
2639pub struct Requesting {
2646 inner: RequestState,
2647 stats: stats::Scope,
2650}
2651
2652enum RequestState {
2653 Ready(broadcast::Consumer),
2655 Failed(Error),
2658 Pending(kio::Consumer<PendingBroadcast>),
2660}
2661
2662impl Requesting {
2663 fn ready(broadcast: broadcast::Consumer) -> Self {
2664 Self {
2665 inner: RequestState::Ready(broadcast),
2666 stats: stats::Scope::default(),
2667 }
2668 }
2669
2670 fn failed(error: Error) -> Self {
2671 Self {
2672 inner: RequestState::Failed(error),
2673 stats: stats::Scope::default(),
2674 }
2675 }
2676
2677 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2678 Self {
2679 inner: RequestState::Pending(consumer),
2680 stats: stats::Scope::default(),
2681 }
2682 }
2683
2684 fn with_stats(mut self, scope: stats::Scope) -> Self {
2685 self.stats = scope;
2686 self
2687 }
2688
2689 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2691 match &self.inner {
2692 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2693 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2694 RequestState::Pending(consumer) => Poll::Ready(
2695 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2696 Some(result) => Poll::Ready(result.clone()),
2697 None => Poll::Pending,
2698 })) {
2699 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2700 Err(_closed) => Err(Error::Unroutable),
2702 },
2703 ),
2704 }
2705 }
2706}
2707
2708impl kio::Pollable for Requesting {
2709 type Output = Result<broadcast::Consumer, Error>;
2710
2711 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2712 self.poll_ok(waiter)
2713 }
2714}
2715
2716pub trait Consume<T> {
2724 fn consume(&self) -> T;
2726}
2727
2728impl<T, U: Consume<T>> Consume<T> for &U {
2729 fn consume(&self) -> T {
2730 (**self).consume()
2731 }
2732}
2733
2734impl Consume<Consumer> for Producer {
2735 fn consume(&self) -> Consumer {
2736 Consumer::new(
2740 self.info,
2741 self.root.clone(),
2742 self.nodes.clone(),
2743 self.dynamic.clone(),
2744 stats::Session::default(),
2745 )
2746 }
2747}
2748
2749impl Consume<Consumer> for Consumer {
2750 fn consume(&self) -> Consumer {
2751 self.clone()
2752 }
2753}
2754
2755impl Consume<broadcast::Consumer> for broadcast::Producer {
2756 fn consume(&self) -> broadcast::Consumer {
2757 self.consume()
2759 }
2760}
2761
2762impl Consume<broadcast::Consumer> for broadcast::Consumer {
2763 fn consume(&self) -> broadcast::Consumer {
2764 self.clone()
2765 }
2766}
2767
2768impl Consume<track::Consumer> for track::Producer {
2769 fn consume(&self) -> track::Consumer {
2770 self.consume()
2771 }
2772}
2773
2774impl Consume<track::Consumer> for track::Consumer {
2775 fn consume(&self) -> track::Consumer {
2776 self.clone()
2777 }
2778}
2779
2780#[derive(Clone)]
2786pub struct Consumer {
2787 info: Origin,
2789 nodes: OriginNodes,
2790
2791 root: PathOwned,
2793
2794 dynamic: kio::Shared<OriginDynamicState>,
2797
2798 stats: stats::Session,
2802
2803 exclude: Option<Origin>,
2807}
2808
2809impl std::ops::Deref for Consumer {
2810 type Target = Origin;
2811
2812 fn deref(&self) -> &Self::Target {
2813 &self.info
2814 }
2815}
2816
2817impl Consumer {
2818 fn new(
2819 info: Origin,
2820 root: PathOwned,
2821 nodes: OriginNodes,
2822 dynamic: kio::Shared<OriginDynamicState>,
2823 stats: stats::Session,
2824 ) -> Self {
2825 Self {
2826 info,
2827 nodes,
2828 root,
2829 dynamic,
2830 stats,
2831 exclude: None,
2832 }
2833 }
2834
2835 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2840 self.exclude = Some(peer);
2841 self
2842 }
2843
2844 pub fn with_stats(mut self, session: stats::Session) -> Self {
2848 self.stats = session;
2849 self
2850 }
2851
2852 fn untagged(&self) -> Self {
2856 Self {
2857 stats: stats::Session::default(),
2858 ..self.clone()
2859 }
2860 }
2861
2862 pub(crate) fn empty(&self) -> Self {
2867 Self {
2868 info: self.info,
2869 nodes: OriginNodes::empty(),
2870 root: self.root.clone(),
2871 dynamic: self.dynamic.clone(),
2872 stats: self.stats.clone(),
2873 exclude: self.exclude,
2874 }
2875 }
2876
2877 pub fn announced(&self) -> AnnounceConsumer {
2884 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone(), self.exclude)
2885 }
2886
2887 pub fn consume(&self) -> Self {
2889 self.clone()
2890 }
2891
2892 fn resolve(&self, path: impl AsPath) -> Resolved {
2901 let path = path.as_path();
2902 let Some(rest) = self.nodes.get(&path) else {
2903 return Resolved::Missing;
2904 };
2905 let state = self.nodes.tree.lock();
2906 state.resolve_broadcast(&rest, self.exclude)
2907 }
2908
2909 #[cfg(test)]
2911 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2912 match self.resolve(path) {
2913 Resolved::Found(broadcast) => Some(broadcast),
2914 Resolved::Excluded | Resolved::Missing => None,
2915 }
2916 }
2917
2918 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2929 let path = path.as_path();
2930
2931 let consumer = self.scope(std::slice::from_ref(&path))?;
2933
2934 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2938 return None;
2939 }
2940
2941 let mut announced = consumer.untagged().announced();
2945 let scope = self.stats.egress(self.root.join(&path).to_owned());
2946 loop {
2947 let OriginAnnounce {
2948 path: announced_path,
2949 broadcast,
2950 } = announced.next().await?;
2951 if announced_path.as_path() == path
2953 && let Some(broadcast) = broadcast
2954 {
2955 return Some(broadcast.with_stats(scope));
2956 }
2957 }
2958 }
2959
2960 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2966 let prefixes = PathPrefixes::new(prefixes);
2967 Some(Consumer {
2968 info: self.info,
2969 root: self.root.clone(),
2970 nodes: self.nodes.select(&prefixes)?,
2971 dynamic: self.dynamic.clone(),
2972 stats: self.stats.clone(),
2973 exclude: self.exclude,
2974 })
2975 }
2976
2977 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2996 let path = path.as_path();
2997
2998 let absolute = self.root.join(&path).to_owned();
3002 let scope = self.stats.egress(&absolute);
3003
3004 match self.resolve(&path) {
3010 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
3011 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
3012 Resolved::Missing => {}
3013 }
3014
3015 let mut state = self.dynamic.lock();
3016
3017 if let Some(weak) = state.served.get(&absolute) {
3021 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
3022 }
3023
3024 let consumer = if let Some(producer) = state.requests.join(&absolute) {
3027 producer.consume()
3028 } else {
3029 let producer = kio::Producer::<PendingBroadcast>::default();
3030 let consumer = producer.consume();
3031 if state.requests.insert(absolute, producer).is_err() {
3032 return kio::Pending::new(Requesting::failed(Error::Unroutable));
3033 }
3034 consumer
3035 };
3036
3037 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
3038 }
3039
3040 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
3045 let prefix = prefix.as_path();
3046
3047 Some(Self {
3048 info: self.info,
3049 root: self.root.join(&prefix).to_owned(),
3050 nodes: self.nodes.root(&prefix)?,
3051 dynamic: self.dynamic.clone(),
3052 stats: self.stats.clone(),
3053 exclude: self.exclude,
3054 })
3055 }
3056
3057 pub fn root(&self) -> &Path<'_> {
3059 &self.root
3060 }
3061
3062 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
3065 self.nodes.nodes.iter().map(|(root, _)| root)
3066 }
3067
3068 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
3070 self.root.join(path)
3071 }
3072}
3073
3074#[derive(Clone)]
3079pub struct AnnounceProducer {
3080 nodes: OriginNodes,
3081 root: PathOwned,
3082}
3083
3084impl AnnounceProducer {
3085 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
3086 Self { nodes, root }
3087 }
3088
3089 pub fn consume(&self) -> AnnounceConsumer {
3095 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default(), None)
3098 }
3099
3100 pub fn root(&self) -> &Path<'_> {
3102 &self.root
3103 }
3104}
3105
3106pub struct AnnounceConsumer {
3111 id: ConsumerId,
3112 nodes: OriginNodes,
3113 root: PathOwned,
3114
3115 state: kio::Producer<OriginConsumerState>,
3118
3119 stats: stats::Session,
3122
3123 guards: HashMap<PathOwned, stats::Announce>,
3127}
3128
3129impl AnnounceConsumer {
3130 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session, exclude: Option<Origin>) -> Self {
3131 let state = kio::Producer::<OriginConsumerState>::default();
3132 let id = ConsumerId::new();
3133
3134 for (_, absolute) in &nodes.nodes {
3135 let notify = AnnounceConsumerNotify {
3136 root: root.clone(),
3137 state: state.clone(),
3138 exclude,
3139 };
3140 nodes.tree.lock().consume_at(id, notify, absolute);
3141 }
3142
3143 Self {
3144 id,
3145 nodes,
3146 root,
3147 state,
3148 stats,
3149 guards: HashMap::new(),
3150 }
3151 }
3152
3153 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
3159 let OriginAnnounce { path, broadcast } = update;
3160 let absolute = self.root.join(&path).to_owned();
3161 match broadcast {
3162 Some(broadcast) => {
3163 let scope = self.stats.egress(&absolute);
3164 self.guards.entry(absolute).or_insert_with(|| scope.announce());
3165 OriginAnnounce {
3166 path,
3167 broadcast: Some(broadcast.with_stats(scope)),
3168 }
3169 }
3170 None => {
3171 self.guards.remove(&absolute);
3172 OriginAnnounce { path, broadcast: None }
3173 }
3174 }
3175 }
3176
3177 pub async fn next(&mut self) -> Option<OriginAnnounce> {
3184 kio::wait(|waiter| self.poll_next(waiter)).await
3185 }
3186
3187 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
3193 let update = {
3194 let mut state = match ready!(self.state.poll(waiter, |state| {
3195 if state.pending.is_empty() {
3196 Poll::Pending
3197 } else {
3198 Poll::Ready(())
3199 }
3200 })) {
3201 Ok(state) => state,
3202 Err(_) => return Poll::Ready(None),
3204 };
3205 state.take().expect("predicate guaranteed an update")
3206 };
3207 Poll::Ready(Some(self.attribute(update)))
3208 }
3209
3210 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
3215 let update = self.state.write().ok()?.take()?;
3216 Some(self.attribute(update))
3217 }
3218
3219 pub fn is_closed(&self) -> bool {
3221 self.state.write().is_err()
3222 }
3223
3224 pub fn root(&self) -> &Path<'_> {
3226 &self.root
3227 }
3228
3229 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
3231 self.root.join(path)
3232 }
3233}
3234
3235impl Drop for AnnounceConsumer {
3236 fn drop(&mut self) {
3237 for (_, absolute) in &self.nodes.nodes {
3238 self.nodes.tree.lock().detach(Claim::Consumer(self.id), absolute);
3239 }
3240 }
3241}
3242
3243#[cfg(test)]
3244use futures::FutureExt;
3245
3246#[cfg(test)]
3247#[allow(missing_docs)] impl AnnounceConsumer {
3249 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3250 let expected = expected.as_path();
3251 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3252 assert_eq!(announce.path, expected, "wrong path");
3253 let announced = announce.broadcast.expect("should be an active announce");
3254 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3255 }
3256
3257 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3261 let expected = expected.as_path();
3262 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3263 assert_eq!(announce.path, expected, "wrong path");
3264 announce.broadcast.expect("should be an active announce")
3265 }
3266
3267 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
3268 let expected = expected.as_path();
3269 let announce = self.try_next().expect("no next");
3270 assert_eq!(announce.path, expected, "wrong path");
3271 let announced = announce.broadcast.expect("should be an active announce");
3272 assert!(announced.is_clone(broadcast), "should be the same broadcast");
3273 }
3274
3275 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
3277 let expected = expected.as_path();
3278 let announce = self.try_next().expect("no next");
3279 assert_eq!(announce.path, expected, "wrong path");
3280 announce.broadcast.expect("should be an active announce")
3281 }
3282
3283 pub fn assert_next_none(&mut self, expected: impl AsPath) {
3284 let expected = expected.as_path();
3285 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3286 assert_eq!(announce.path, expected, "wrong path");
3287 assert!(announce.broadcast.is_none(), "should be unannounced");
3288 }
3289
3290 pub fn assert_next_wait(&mut self) {
3291 if let Some(res) = self.next().now_or_never() {
3292 panic!("next should block: got {:?}", res.map(|a| a.path));
3293 }
3294 }
3295
3296 }
3305
3306#[cfg(test)]
3307mod tests {
3308 use crate::coding::Decode;
3309 use crate::group;
3310
3311 use super::*;
3312
3313 fn announce() -> broadcast::Route {
3315 broadcast::Route::new().with_announce(true)
3316 }
3317
3318 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
3324 let name = Path::new(name);
3325 let peer_key = fnv_key(&name, [peer]);
3326 (100u64..)
3327 .map(|id| Origin::new(id).unwrap())
3328 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
3329 .unwrap()
3330 }
3331
3332 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
3335 let source = broadcast::Info::new().produce().consume();
3336 FrontState {
3337 path: Path::new("test").to_owned(),
3338 self_origin,
3339 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
3340 next_route: routes.len() as u64,
3341 excluded: HashMap::new(),
3342 routes: routes
3343 .into_iter()
3344 .enumerate()
3345 .map(|(id, route)| FrontRoute {
3346 id: id as u64,
3347 route,
3348 source: source.clone(),
3349 })
3350 .collect(),
3351 active: Some(0),
3352 linger: Duration::ZERO,
3353 closed: false,
3354 }
3355 }
3356
3357 fn sibling_route(peer: Origin) -> broadcast::Route {
3360 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3361 announce().with_hops(hops)
3362 }
3363
3364 fn upstream_route(cost: u64) -> broadcast::Route {
3366 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3367 announce().with_hops(hops).with_cost(cost)
3368 }
3369
3370 #[test]
3374 fn test_carrying_gate_keys() {
3375 let peer = Origin::new(3).unwrap();
3376
3377 let mut lost = front_state(
3379 origin_keyed("test", peer, false),
3380 vec![upstream_route(10), sibling_route(peer)],
3381 );
3382 lost.reselect(true);
3383 assert_eq!(
3384 lost.active,
3385 Some(0),
3386 "carrying front re-parented onto a higher-keyed peer"
3387 );
3388 lost.reselect(false);
3389 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3390
3391 let mut won = front_state(
3393 origin_keyed("test", peer, true),
3394 vec![upstream_route(10), sibling_route(peer)],
3395 );
3396 won.reselect(true);
3397 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3398 }
3399
3400 #[test]
3405 fn test_carrying_gate_symmetric_race() {
3406 let a = Origin::new(1).unwrap();
3407 let b = Origin::new(2).unwrap();
3408
3409 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3410 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3411 a_view.reselect(true);
3412 b_view.reselect(true);
3413
3414 let a_moved = a_view.active == Some(1);
3415 let b_moved = b_view.active == Some(1);
3416 assert!(
3417 a_moved != b_moved,
3418 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3419 );
3420 }
3421
3422 #[test]
3427 fn test_carrying_switches_to_benign_routes() {
3428 let peer = Origin::new(3).unwrap();
3429 let lost = origin_keyed("test", peer, false);
3430
3431 let mut forwarder = sibling_route(peer).with_cost(4);
3433 forwarder.advertised = 4;
3434 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3435 state.reselect(true);
3436 assert_eq!(
3437 state.active,
3438 Some(1),
3439 "a cheaper forwarder path must win while carrying"
3440 );
3441
3442 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3444 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3445 state.reselect(true);
3446 assert_eq!(
3447 state.active,
3448 Some(1),
3449 "a direct publisher route must win while carrying"
3450 );
3451
3452 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3457 state.reselect(true);
3458 assert_eq!(
3459 state.active,
3460 Some(1),
3461 "a reconnect on an identical chain must win while carrying"
3462 );
3463 }
3464
3465 #[test]
3468 fn test_carrying_gate_ignores_unannounced_incumbent() {
3469 let peer = Origin::new(3).unwrap();
3470 let unannounced = upstream_route(10).with_announce(false);
3471 let mut state = front_state(
3472 origin_keyed("test", peer, false),
3473 vec![unannounced, sibling_route(peer)],
3474 );
3475 state.reselect(true);
3476 assert_eq!(
3477 state.active,
3478 Some(1),
3479 "an unannounced incumbent must always be displaced"
3480 );
3481 }
3482
3483 #[test]
3488 fn test_reflection_through_an_exposed_peer_cannot_take_over() {
3489 let peer = Origin::new(42).unwrap();
3490 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3491 let reflected = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3492
3493 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3494
3495 assert!(!state.taints_a_reader(&reflected));
3497
3498 *state.excluded.entry(peer).or_default() += 1;
3501 assert!(state.taints_a_reader(&reflected));
3502 }
3503
3504 #[test]
3507 fn test_rival_publisher_through_another_peer_still_takes_over() {
3508 let peer = Origin::new(42).unwrap();
3509 let elsewhere = Origin::new(43).unwrap();
3510 let upstream = OriginList::try_from(vec![Origin::new(7).unwrap()]).unwrap();
3511 let rival = announce().with_hops(OriginList::try_from(vec![Origin::UNKNOWN, elsewhere]).unwrap());
3512
3513 let mut state = front_state(Origin::new(1).unwrap(), vec![announce().with_hops(upstream)]);
3514 *state.excluded.entry(peer).or_default() += 1;
3515
3516 assert!(!state.taints_a_reader(&rival), "only the peer we feed is a reflection");
3517 }
3518
3519 #[test]
3524 fn test_opaque_peer_understates_its_depth() {
3525 let us = Origin::new(1).unwrap();
3526
3527 let direct = || {
3529 announce()
3530 .with_hops(OriginList::try_from(vec![Origin::new(7).unwrap(), Origin::new(8).unwrap()]).unwrap())
3531 .with_cost(2)
3532 };
3533 let opaque = |cost| {
3536 announce()
3537 .with_hops(OriginList::try_from(vec![Origin::new(42).unwrap()]).unwrap())
3538 .with_cost(cost)
3539 };
3540
3541 let mut state = front_state(us, vec![direct(), opaque(1)]);
3543 state.reselect(false);
3544 assert_eq!(
3545 state.active,
3546 Some(1),
3547 "an unpriced opaque link out-ranks a shorter real path"
3548 );
3549
3550 let mut state = front_state(us, vec![direct(), opaque(16)]);
3552 state.reselect(false);
3553 assert_eq!(
3554 state.active,
3555 Some(0),
3556 "pricing the opaque link restores the intended order"
3557 );
3558 }
3559
3560 async fn settle() {
3563 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3564 }
3565
3566 fn origin_scoped(prefixes: &[Path]) -> Producer {
3568 Origin::random().produce().scope(prefixes).expect("in scope")
3569 }
3570
3571 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3574 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3575 .await
3576 .expect("timed out waiting for a track request")
3577 .expect("source closed");
3578 assert_eq!(request.name(), name, "unexpected track dispatched");
3579 request.accept(None)
3580 }
3581
3582 async fn accept_tracks(dynamic: &mut broadcast::Dynamic, count: usize) -> HashMap<String, track::Producer> {
3586 let mut accepted = HashMap::new();
3587 for _ in 0..count {
3588 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3589 .await
3590 .expect("timed out waiting for a track request")
3591 .expect("source closed");
3592 let name = request.name().to_string();
3593 accepted.insert(name, request.accept(None));
3594 }
3595 accepted
3596 }
3597
3598 #[tokio::test]
3602 async fn test_stats_tagged_end_to_end() {
3603 use crate::Timestamp;
3604 use crate::stats::{Config, Registry, Tier};
3605 use bytes::Bytes;
3606
3607 tokio::time::pause();
3608
3609 let registry = Registry::new(Config::new());
3610 let ctx = registry.tier(Tier::default()).session("acme");
3611
3612 let origin = Origin::random().produce();
3613 let ingress = origin.clone().with_stats(ctx.clone());
3614 let egress = origin.consume().with_stats(ctx.clone());
3615
3616 let mut announced = egress.announced();
3619
3620 let source = ingress.create_broadcast("demo", announce()).unwrap();
3622 let mut dynamic = source.dynamic();
3623 settle().await;
3624 settle().await;
3625
3626 let update = announced.next().await.unwrap();
3628 assert_eq!(update.path.as_str(), "demo");
3629 let broadcast = update.broadcast.unwrap();
3630
3631 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3633 let mut producer = accept_track(&mut dynamic, "video").await;
3634 settle().await;
3635 let mut sub = subscribing.await.unwrap();
3636
3637 let mut group = producer.append_group().unwrap();
3639 group
3640 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3641 .unwrap();
3642 group
3643 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3644 .unwrap();
3645 group.finish().unwrap();
3646
3647 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3649 let mut frames = 0;
3650 while let Some(frame) = group_c.read_frame().await.unwrap() {
3651 assert_eq!(frame.payload.len(), 5);
3652 frames += 1;
3653 }
3654 assert_eq!(frames, 2);
3655 settle().await;
3656
3657 let report = registry.report();
3658 let entry = report
3659 .traffic
3660 .iter()
3661 .find(|e| e.path.as_str() == "demo")
3662 .expect("demo tracked");
3663 let path_len = "demo".len() as u64;
3664
3665 let egress = &entry.publisher;
3667 assert_eq!(egress.announced, 1, "one egress announce");
3668 assert_eq!(egress.announced_bytes, path_len);
3669 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3670 assert_eq!(egress.broadcasts, 1, "one viewer");
3671 assert_eq!(egress.groups, 1);
3672 assert_eq!(egress.frames, 2);
3673 assert_eq!(egress.bytes, 10);
3674 assert_eq!(egress.fetches, 0);
3675
3676 let ingress = &entry.subscriber;
3678 assert_eq!(ingress.announced, 1, "one ingress announce");
3679 assert_eq!(ingress.announced_bytes, path_len);
3680 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3681 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3682 assert_eq!(ingress.groups, 1);
3683 assert_eq!(ingress.frames, 2);
3684 assert_eq!(ingress.bytes, 10);
3685
3686 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3688 let _ = fetched;
3689 settle().await;
3690 let report = registry.report();
3691 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3692 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3693 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3694 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3695 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3699 }
3700
3701 #[tokio::test]
3706 async fn test_stats_read_frame_counts_once() {
3707 use crate::Timestamp;
3708 use crate::stats::{Config, Registry, Tier};
3709 use bytes::Bytes;
3710
3711 tokio::time::pause();
3712
3713 let registry = Registry::new(Config::new());
3714 let ctx = registry.tier(Tier::default()).session("acme");
3715
3716 let origin = Origin::random().produce();
3717 let ingress = origin.clone().with_stats(ctx.clone());
3718 let egress = origin.consume().with_stats(ctx.clone());
3719
3720 let mut announced = egress.announced();
3721 let source = ingress.create_broadcast("demo", announce()).unwrap();
3722 let mut dynamic = source.dynamic();
3723 settle().await;
3724 settle().await;
3725
3726 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3727 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3728 let mut producer = accept_track(&mut dynamic, "video").await;
3729 settle().await;
3730 let mut sub = subscribing.await.unwrap();
3731
3732 producer
3734 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3735 .unwrap();
3736
3737 let frame = sub.read_frame().await.unwrap().expect("frame");
3738 assert_eq!(frame.payload.len(), 5);
3739 settle().await;
3740
3741 let report = registry.report();
3742 let entry = report
3743 .traffic
3744 .iter()
3745 .find(|e| e.path.as_str() == "demo")
3746 .expect("demo tracked");
3747 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3748 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3749 assert_eq!(
3750 entry.publisher.bytes, 5,
3751 "payload counted once, not zero and not doubled"
3752 );
3753 }
3754
3755 #[tokio::test]
3759 async fn test_stats_datagrams_counted_both_sides() {
3760 use crate::Timestamp;
3761 use crate::stats::{Config, Registry, Tier};
3762
3763 tokio::time::pause();
3764
3765 let registry = Registry::new(Config::new());
3766 let ctx = registry.tier(Tier::default()).session("acme");
3767
3768 let origin = Origin::random().produce();
3769 let ingress = origin.clone().with_stats(ctx.clone());
3770 let egress = origin.consume().with_stats(ctx.clone());
3771
3772 let mut announced = egress.announced();
3773 let source = ingress.create_broadcast("demo", announce()).unwrap();
3774 let mut dynamic = source.dynamic();
3775 settle().await;
3776 settle().await;
3777
3778 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3779 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3780 let mut producer = accept_track(&mut dynamic, "video").await;
3781 settle().await;
3782 let mut sub = subscribing.await.unwrap();
3783
3784 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3785 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3786 assert_eq!(&datagram.payload[..], b"hello");
3787 settle().await;
3788
3789 let report = registry.report();
3790 let entry = report
3791 .traffic
3792 .iter()
3793 .find(|e| e.path.as_str() == "demo")
3794 .expect("demo tracked");
3795
3796 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3797 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3798 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3799 assert_eq!(traffic.frames, 1, "{side}: one frame");
3800 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3801 }
3802 }
3803
3804 #[test]
3805 fn origin_rejects_reserved_ids() {
3806 assert!(Origin::new(0).is_err());
3807 assert!(Origin::new(1u64 << 62).is_err());
3808 assert_eq!(Origin::new(1).unwrap().id(), 1);
3809
3810 let mut zero = [0u8].as_slice();
3811 assert_eq!(
3812 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3813 Origin::UNKNOWN
3814 );
3815 }
3816
3817 #[test]
3818 fn origin_list_push_fails_at_limit() {
3819 let mut list = OriginList::new();
3820 for _ in 0..MAX_HOPS {
3821 list.push(Origin::random()).unwrap();
3822 }
3823 assert_eq!(list.len(), MAX_HOPS);
3824 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3825 }
3826
3827 #[test]
3828 fn origin_list_replace_first() {
3829 let mut list = OriginList::new();
3830 for _ in 0..3 {
3831 list.push(Origin::UNKNOWN).unwrap();
3832 }
3833
3834 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3836 assert_eq!(
3837 list.as_slice(),
3838 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3839 );
3840
3841 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3843 assert_eq!(list.len(), 3);
3844 }
3845
3846 #[test]
3847 fn origin_list_try_from_vec_enforces_limit() {
3848 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3849 assert!(OriginList::try_from(under).is_ok());
3850
3851 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3852 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3853 }
3854
3855 #[tokio::test]
3859 async fn test_idle_announce_prunes_its_nodes() {
3860 tokio::time::pause();
3861
3862 let origin = Origin::random().produce();
3863 let consumer = origin.consume();
3864 let start = origin.node_count();
3865
3866 for i in 0..32 {
3867 let path = format!("channel{i}/chat");
3868 let scoped = consumer.scope(&[Path::new(path.as_str())]).expect("in scope");
3869
3870 let mut announced = scoped.announced();
3871 announced.assert_next_wait();
3872 assert_eq!(origin.node_count(), start + 2, "the subtree exists while subscribed");
3873
3874 drop(announced);
3875 drop(scoped);
3876 assert_eq!(origin.node_count(), start, "cycle {i} left a node behind");
3877 }
3878 }
3879
3880 #[tokio::test]
3883 async fn test_prune_spares_live_nodes() {
3884 tokio::time::pause();
3885
3886 let origin = Origin::random().produce();
3887 let consumer = origin.consume();
3888 let bare = origin.node_count();
3889
3890 let mut broadcast = origin.create_broadcast("channel/chat", announce()).unwrap();
3891 settle().await;
3892 let live = origin.node_count();
3893 assert_eq!(live, bare + 2, "the broadcast should have created its subtree");
3894
3895 let mut announced = consumer.announced();
3897 announced.assert_next_some("channel/chat");
3898 drop(announced);
3899 assert_eq!(origin.node_count(), live, "an announced path was pruned");
3900 assert!(consumer.get_broadcast("channel/chat").is_some());
3901
3902 let idle = consumer.scope(&[Path::new("idle")]).expect("in scope");
3904 let first = idle.announced();
3905 let second = idle.announced();
3906 assert_eq!(origin.node_count(), live + 1);
3907 drop(first);
3908 assert_eq!(origin.node_count(), live + 1, "a node with a consumer was pruned");
3909 drop(second);
3910 assert_eq!(origin.node_count(), live, "the last cursor left the node behind");
3911
3912 broadcast.finish();
3914 settle().await;
3915 assert_eq!(origin.node_count(), bare);
3916 }
3917
3918 #[tokio::test]
3922 async fn test_multi_prefix_cursor_prunes_both_branches() {
3923 tokio::time::pause();
3924
3925 let origin = Origin::random().produce();
3926 let consumer = origin.consume();
3927 let bare = origin.node_count();
3928
3929 let scoped = consumer
3930 .scope(&[Path::new("room/a"), Path::new("room/b")])
3931 .expect("in scope");
3932 assert_eq!(origin.node_count(), bare, "scoping should not create nodes");
3933 assert!(consumer.with_root("room/c").is_some());
3934 assert_eq!(origin.node_count(), bare, "rooting should not create nodes");
3935
3936 let announced = scoped.announced();
3937 assert_eq!(origin.node_count(), bare + 3, "room, room/a and room/b");
3938
3939 drop(announced);
3940 assert_eq!(origin.node_count(), bare, "a two-branch cursor left nodes behind");
3941 }
3942
3943 #[tokio::test]
3946 async fn test_prune_stops_at_a_live_ancestor() {
3947 tokio::time::pause();
3948
3949 let origin = Origin::random().produce();
3950 let consumer = origin.consume();
3951 let bare = origin.node_count();
3952
3953 let outer = consumer.scope(&[Path::new("room/a")]).expect("in scope").announced();
3954 let inner = consumer
3955 .scope(&[Path::new("room/a/deep/leaf")])
3956 .expect("in scope")
3957 .announced();
3958 assert_eq!(origin.node_count(), bare + 4, "room, a, deep and leaf");
3959
3960 drop(inner);
3961 assert_eq!(origin.node_count(), bare + 2, "pruning ran past the cursor above it");
3962
3963 drop(outer);
3964 assert_eq!(origin.node_count(), bare);
3965 }
3966
3967 #[tokio::test]
3972 async fn test_pending_source_holds_its_node() {
3973 tokio::time::pause();
3974
3975 let origin = Origin::random().produce();
3976 let consumer = origin.consume();
3977 let bare = origin.node_count();
3978
3979 let waiting = consumer
3981 .scope(&[Path::new("channel/chat")])
3982 .expect("in scope")
3983 .announced();
3984 assert_eq!(origin.node_count(), bare + 2);
3985
3986 let _broadcast = origin.create_broadcast("channel/chat", announce()).unwrap();
3987
3988 drop(waiting);
3990 assert_eq!(origin.node_count(), bare + 2, "a pending source's node was pruned");
3991
3992 settle().await;
3993 assert!(
3994 consumer.get_broadcast("channel/chat").is_some(),
3995 "the source attached into an orphan"
3996 );
3997 }
3998
3999 #[tokio::test]
4003 async fn test_scoped_handles_survive_a_prune() {
4004 tokio::time::pause();
4005
4006 let scope = [Path::new("channel")];
4007 let producer = origin_scoped(&scope);
4008 let consumer = producer.consume().scope(&scope).expect("in scope");
4009
4010 drop(consumer.announced());
4012
4013 let _broadcast = producer.create_broadcast("channel/chat", announce()).unwrap();
4014 settle().await;
4015
4016 let mut announced = consumer.announced();
4017 announced.assert_next_some("channel/chat");
4018 assert!(consumer.get_broadcast("channel/chat").is_some());
4019 }
4020
4021 #[tokio::test]
4022 async fn test_announce() {
4023 tokio::time::pause();
4024
4025 let origin = Origin::random().produce();
4026
4027 let mut consumer1 = origin.consume().announced();
4028 consumer1.assert_next_wait();
4029
4030 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
4032 settle().await;
4033
4034 consumer1.assert_next_some("test1");
4035 consumer1.assert_next_wait();
4036
4037 let mut consumer2 = origin.consume().announced();
4040
4041 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
4043 settle().await;
4044
4045 consumer1.assert_next_some("test2");
4046 consumer1.assert_next_wait();
4047
4048 consumer2.assert_next_some("test1");
4049 consumer2.assert_next_some("test2");
4050 consumer2.assert_next_wait();
4051
4052 broadcast1.finish();
4054 settle().await;
4055
4056 consumer1.assert_next_none("test1");
4058 consumer2.assert_next_none("test1");
4059 consumer1.assert_next_wait();
4060 consumer2.assert_next_wait();
4061
4062 let mut consumer3 = origin.consume().announced();
4064 consumer3.assert_next_some("test2");
4065 consumer3.assert_next_wait();
4066
4067 broadcast2.finish();
4068 settle().await;
4069
4070 consumer1.assert_next_none("test2");
4071 consumer2.assert_next_none("test2");
4072 consumer3.assert_next_none("test2");
4073 }
4074
4075 #[tokio::test]
4079 async fn test_duplicate() {
4080 tokio::time::pause();
4081
4082 let origin = Origin::random().produce();
4083 let consumer = origin.consume();
4084 let mut announced = consumer.announced();
4085
4086 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
4087 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
4088 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
4089 settle().await;
4090 assert!(consumer.get_broadcast("test").is_some());
4091
4092 announced.assert_next_some("test");
4093 announced.assert_next_wait();
4094
4095 broadcast2.finish();
4097 settle().await;
4098 assert!(consumer.get_broadcast("test").is_some());
4099 announced.assert_next_wait();
4100
4101 broadcast1.finish();
4103 settle().await;
4104 assert!(consumer.get_broadcast("test").is_some());
4105 announced.assert_next_wait();
4106
4107 broadcast3.finish();
4109 settle().await;
4110 assert!(consumer.get_broadcast("test").is_none());
4111
4112 announced.assert_next_none("test");
4113 announced.assert_next_wait();
4114 }
4115
4116 #[tokio::test]
4119 async fn test_route_failover() {
4120 tokio::time::pause();
4121
4122 let origin = Origin::random().produce();
4123 let consumer = origin.consume();
4124 let mut announced = consumer.announced();
4125
4126 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4129 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4130
4131 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4133 let mut dynamic_a = source_a.dynamic();
4134 settle().await;
4135 settle().await;
4136 let broadcast = consumer.request_broadcast("test").await.unwrap();
4137 announced.assert_next_some("test");
4138
4139 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4141 let mut dynamic_b = source_b.dynamic();
4142 settle().await;
4143 settle().await;
4144 announced.assert_next_wait();
4145
4146 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4148 let mut producer = accept_track(&mut dynamic_a, "video").await;
4149 settle().await;
4150 dynamic_b.assert_no_request();
4151
4152 let mut sub = subscribing.await.unwrap();
4153 sub.assert_no_group();
4156 assert_eq!(producer.subscription().unwrap().group_start, None);
4157
4158 producer.append_group().unwrap();
4159 producer.append_group().unwrap();
4160 assert_eq!(sub.assert_group().sequence, 0);
4161 assert_eq!(sub.assert_group().sequence, 1);
4162
4163 producer.abort(Error::Dropped).unwrap();
4167 source_a.abort(Error::Dropped).unwrap();
4168 drop(dynamic_a);
4169 settle().await;
4170 announced.assert_next_wait();
4171
4172 let mut producer = accept_track(&mut dynamic_b, "video").await;
4176 settle().await;
4177 sub.assert_no_group();
4178 assert_eq!(producer.subscription().unwrap().group_start, None);
4179 producer.create_group(group::Info { sequence: 1 }).unwrap();
4180 producer.create_group(group::Info { sequence: 2 }).unwrap();
4181 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
4182 sub.assert_not_closed();
4183 }
4184
4185 #[tokio::test]
4192 async fn test_route_failover_restores_every_track() {
4193 tokio::time::pause();
4194
4195 let origin = Origin::random().produce();
4196 let consumer = origin.consume();
4197
4198 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4200 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4201
4202 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4203 let mut dynamic_a = source_a.dynamic();
4204 settle().await;
4205 settle().await;
4206 let broadcast = consumer.request_broadcast("test").await.unwrap();
4207
4208 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4210 let mut dynamic_b = source_b.dynamic();
4211 settle().await;
4212 settle().await;
4213
4214 const TRACKS: [(&str, u64); 3] = [("video", 3), ("audio", 1), ("data", 2)];
4218
4219 let subscribing: Vec<_> = TRACKS
4220 .iter()
4221 .map(|(name, _)| broadcast.track(name).unwrap().subscribe(None))
4222 .collect();
4223 let mut producers_a = accept_tracks(&mut dynamic_a, TRACKS.len()).await;
4224 settle().await;
4225
4226 let mut subs = Vec::new();
4227 for ((name, groups), subscribing) in TRACKS.iter().zip(subscribing) {
4228 let mut sub = subscribing.await.unwrap();
4229 let producer = producers_a
4230 .get_mut(*name)
4231 .unwrap_or_else(|| panic!("{name} was never dispatched"));
4232 for expected in 0..*groups {
4233 producer.append_group().unwrap();
4234 assert_eq!(sub.assert_group().sequence, expected, "{name} did not start");
4235 }
4236 subs.push((*name, *groups, sub));
4237 }
4238
4239 for (_, producer) in producers_a.drain() {
4241 producer.abort(Error::Dropped).unwrap();
4242 }
4243 source_a.abort(Error::Dropped).unwrap();
4244 drop(dynamic_a);
4245 settle().await;
4246
4247 let mut producers_b = accept_tracks(&mut dynamic_b, TRACKS.len()).await;
4250 settle().await;
4251
4252 for (_, _, sub) in subs.iter_mut() {
4255 sub.assert_no_group();
4256 }
4257 settle().await;
4258
4259 for (name, groups, sub) in subs.iter_mut() {
4260 let producer = producers_b
4261 .get_mut(*name)
4262 .unwrap_or_else(|| panic!("{name} was never re-dispatched to the standby"));
4263 assert_eq!(
4266 producer
4267 .subscription()
4268 .unwrap_or_else(|| panic!("{name} resumed without a subscription"))
4269 .group_start,
4270 None,
4271 "{name} must keep the subscriber's live-edge demand"
4272 );
4273 let boundary = *groups;
4274
4275 producer.create_group(group::Info { sequence: boundary - 1 }).unwrap();
4277 producer.create_group(group::Info { sequence: boundary }).unwrap();
4278 assert_eq!(sub.assert_group().sequence, boundary, "{name} did not resume");
4279 sub.assert_not_closed();
4280 }
4281 }
4282
4283 #[tokio::test]
4286 async fn test_broadcast_route_watch() {
4287 let mut producer = broadcast::Info::new().produce();
4288 let mut consumer = producer.consume();
4289
4290 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
4292
4293 producer.set_route(broadcast::Route::default()).unwrap();
4295 assert!(consumer.route_changed().now_or_never().is_none());
4296
4297 let mut hops = OriginList::new();
4298 hops.push(Origin::new(7).unwrap()).unwrap();
4299 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
4300 producer.set_route(route.clone()).unwrap();
4301 assert_eq!(consumer.route_changed().await.unwrap(), route);
4302
4303 let mut fresh = producer.consume();
4305 assert_eq!(fresh.route_changed().await.unwrap(), route);
4306
4307 drop(producer);
4308 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
4309 }
4310
4311 #[tokio::test]
4315 async fn test_route_cost_update() {
4316 tokio::time::pause();
4317
4318 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
4322 let consumer = origin.consume();
4323 let mut announced = consumer.announced();
4324
4325 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4328 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4329
4330 let mut source_a = origin
4332 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4333 .unwrap();
4334 let mut dynamic_a = source_a.dynamic();
4335 settle().await;
4336 let broadcast = consumer.request_broadcast("test").await.unwrap();
4337 announced.assert_next_some("test");
4338
4339 let mut watch = broadcast.clone();
4340 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
4341
4342 let mut source_b = origin
4343 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4344 .unwrap();
4345 let mut dynamic_b = source_b.dynamic();
4346 settle().await;
4347 assert!(
4348 watch.route_changed().now_or_never().is_none(),
4349 "a losing standby must not change the advertised route"
4350 );
4351
4352 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4354 let mut producer = accept_track(&mut dynamic_a, "video").await;
4355 settle().await;
4356 let mut sub = subscribing.await.unwrap();
4357 producer.append_group().unwrap();
4358 assert_eq!(sub.assert_group().sequence, 0);
4359
4360 source_a
4363 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
4364 .unwrap();
4365 settle().await;
4366 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
4367 announced.assert_next_wait();
4368
4369 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
4370 settle().await;
4371 sub.assert_no_group();
4374 assert_eq!(producer_b.subscription().unwrap().group_start, None);
4375 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
4376 assert_eq!(sub.assert_group().sequence, 1);
4377 sub.assert_not_closed();
4378
4379 source_b
4381 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
4382 .unwrap();
4383 settle().await;
4384 let advertised = watch.route_changed().await.unwrap();
4385 assert_eq!(advertised.hops, hops_b);
4386 assert_eq!(advertised.cost, 5);
4387 announced.assert_next_wait();
4388 }
4389
4390 #[tokio::test]
4393 async fn test_completed_track_survives_route_churn() {
4394 tokio::time::pause();
4395
4396 let origin = Origin::random().produce();
4397 let consumer = origin.consume();
4398
4399 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4401 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4402
4403 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4404 let mut dynamic_a = source_a.dynamic();
4405 settle().await;
4406 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4407 let mut dynamic_b = source_b.dynamic();
4408 settle().await;
4409 settle().await;
4410 let broadcast = consumer.request_broadcast("test").await.unwrap();
4411
4412 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4414 let mut producer = accept_track(&mut dynamic_a, "video").await;
4415 settle().await;
4416 let mut sub = subscribing.await.unwrap();
4417 producer.append_group().unwrap();
4418 assert_eq!(sub.assert_group().sequence, 0);
4419 producer.finish().unwrap();
4420 drop(producer);
4421 settle().await;
4422 sub.assert_closed();
4423
4424 source_a.abort(Error::Dropped).unwrap();
4426 drop(dynamic_a);
4427 settle().await;
4428 dynamic_b.assert_no_request();
4429
4430 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
4432 late.assert_closed();
4433 }
4434
4435 #[tokio::test]
4439 async fn test_refused_track_aborts_instantly() {
4440 tokio::time::pause();
4441
4442 let origin = Origin::random().produce();
4443 let consumer = origin.consume();
4444
4445 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4446 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4447 let mut dynamic = source.dynamic();
4448 settle().await;
4449 settle().await;
4450 let broadcast = consumer.request_broadcast("test").await.unwrap();
4451
4452 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4453 let request = dynamic.requested_track().await.unwrap();
4454 request.reject(Error::NotFound);
4455 settle().await;
4456
4457 assert!(matches!(subscribing.await, Err(Error::NotFound)));
4459 dynamic.assert_no_request();
4460 }
4461
4462 #[tokio::test]
4467 async fn test_stale_rejection_does_not_abort_a_handover() {
4468 tokio::time::pause();
4469
4470 let origin = Origin::random().produce();
4471 let consumer = origin.consume();
4472
4473 let publisher = Origin::new(1).unwrap();
4474 let peer = Origin::new(5).unwrap();
4475 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4476 let local = OriginList::try_from(vec![publisher]).unwrap();
4477
4478 let source_remote = origin
4480 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4481 .unwrap();
4482 let mut dynamic_remote = source_remote.dynamic();
4483 settle().await;
4484 settle().await;
4485 let broadcast = consumer.request_broadcast("test").await.unwrap();
4486 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4487 let request_remote = dynamic_remote.requested_track().await.unwrap();
4488
4489 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4492 let mut dynamic_local = source_local.dynamic();
4493 request_remote.reject(Error::NotFound);
4494 settle().await;
4495
4496 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4498 settle().await;
4499 let mut sub = subscribing
4500 .await
4501 .expect("the handover must win over the stale rejection");
4502 producer_local.append_group().unwrap();
4503 assert_eq!(sub.assert_group().sequence, 0);
4504 sub.assert_not_closed();
4505 }
4506
4507 #[tokio::test]
4511 async fn test_route_handover() {
4512 tokio::time::pause();
4513
4514 let origin = Origin::random().produce();
4515 let consumer = origin.consume();
4516 let mut announced = consumer.announced();
4517
4518 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4520 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4521
4522 let source_a = origin
4523 .create_broadcast("test", announce().with_hops(hops_long))
4524 .unwrap();
4525 let mut dynamic_a = source_a.dynamic();
4526 settle().await;
4527 settle().await;
4528 let broadcast = consumer.request_broadcast("test").await.unwrap();
4529 announced.assert_next_some("test");
4530
4531 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4532 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
4533 settle().await;
4534 let mut sub = subscribing.await.unwrap();
4535 producer_a.append_group().unwrap();
4536 producer_a.append_group().unwrap();
4537 assert_eq!(sub.assert_group().sequence, 0);
4538 assert_eq!(sub.assert_group().sequence, 1);
4539
4540 let source_b = origin
4543 .create_broadcast("test", announce().with_hops(hops_short))
4544 .unwrap();
4545 let mut dynamic_b = source_b.dynamic();
4546 settle().await;
4547 settle().await;
4548 announced.assert_next_wait();
4549
4550 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
4551 settle().await;
4552
4553 sub.assert_no_group();
4557 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
4558 assert_eq!(producer_b.subscription().unwrap().group_start, None);
4559
4560 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
4562 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
4563 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
4564 assert_eq!(sub.assert_group().sequence, 2);
4565 assert_eq!(sub.assert_group().sequence, 3);
4566 sub.assert_no_group();
4567 sub.assert_not_closed();
4568 }
4569
4570 #[tokio::test(start_paused = true)]
4573 async fn test_route_unannounce_immediate() {
4574 let origin = Origin::random().produce();
4575 let consumer = origin.consume();
4576 let mut announced = consumer.announced();
4577
4578 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4579 let mut source = origin
4580 .create_broadcast("test", announce().with_hops(hops.clone()))
4581 .unwrap();
4582 settle().await;
4583 let broadcast = consumer.request_broadcast("test").await.unwrap();
4584 announced.assert_next_some("test");
4585
4586 source.finish();
4589 settle().await;
4590 announced.assert_next_none("test");
4591
4592 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4594 settle().await;
4595 let fresh = consumer.request_broadcast("test").await.unwrap();
4596 announced.assert_next_some("test");
4597 assert!(
4598 !fresh.is_clone(&broadcast),
4599 "re-create must not splice the old broadcast"
4600 );
4601 }
4602
4603 #[tokio::test(start_paused = true)]
4608 async fn test_route_detach_immediate() {
4609 let origin = Origin::random().produce();
4610 let consumer = origin.consume();
4611 let mut announced = consumer.announced();
4612
4613 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4614 let source = origin
4615 .create_broadcast("test", announce().with_hops(hops.clone()))
4616 .unwrap();
4617 let mut dynamic = source.dynamic();
4618 settle().await;
4619 settle().await;
4620 let broadcast = consumer.request_broadcast("test").await.unwrap();
4621 announced.assert_next_some("test");
4622
4623 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4624 let producer = accept_track(&mut dynamic, "video").await;
4625 settle().await;
4626 let mut sub = subscribing.await.unwrap();
4627
4628 drop(producer);
4630 source.abort(Error::Dropped).unwrap();
4631 drop(dynamic);
4632
4633 settle().await;
4634 announced.assert_next_none("test");
4635 sub.assert_error();
4636
4637 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4640 settle().await;
4641 settle().await;
4642 let fresh = consumer.request_broadcast("test").await.unwrap();
4643 announced.assert_next_some("test");
4644 assert!(
4645 !fresh.is_clone(&broadcast),
4646 "re-create must not splice the old broadcast"
4647 );
4648 }
4649
4650 #[tokio::test(start_paused = true)]
4655 async fn test_idle_track_releases_without_respinning() {
4656 let origin = Info::new(Origin::random()).produce();
4657 let consumer = origin.consume();
4658
4659 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4660 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4661 let mut dynamic = source.dynamic();
4662 settle().await;
4663 let broadcast = consumer.request_broadcast("test").await.unwrap();
4664
4665 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4666 let producer = accept_track(&mut dynamic, "video").await;
4667 settle().await;
4668 let sub = subscribing.await.unwrap();
4669
4670 drop(sub);
4673 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
4674 settle().await;
4675 assert!(
4676 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
4677 "the copy must stay spliced inside the linger",
4678 );
4679
4680 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4683 settle().await;
4684 assert!(
4685 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4686 "an idle copy must be released after the linger",
4687 );
4688
4689 for _ in 0..3 {
4693 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4694 settle().await;
4695 assert!(
4696 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4697 "an unread copy must stay released, not be re-spliced",
4698 );
4699 }
4700 assert!(
4701 dynamic.requested_track().now_or_never().is_none(),
4702 "an unread track must not be re-requested",
4703 );
4704 drop(producer);
4705
4706 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4708 let mut producer = accept_track(&mut dynamic, "video").await;
4709 settle().await;
4710 let mut sub = subscribing.await.unwrap();
4711 producer.append_group().unwrap();
4712 assert_eq!(sub.assert_group().sequence, 0);
4713 }
4714
4715 #[tokio::test(start_paused = true)]
4719 async fn test_back_to_back_fetches_reuse_the_track() {
4720 let origin = Info::new(Origin::random()).produce();
4721 let consumer = origin.consume();
4722
4723 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4724 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4725 let mut dynamic = source.dynamic();
4726 settle().await;
4727 let broadcast = consumer.request_broadcast("test").await.unwrap();
4728
4729 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4731 let mut producer = accept_track(&mut dynamic, "video").await;
4732 producer.append_group().unwrap().finish().unwrap();
4733 settle().await;
4734 let first = fetching.await.expect("first fetch");
4735 drop(first);
4736
4737 settle().await;
4739 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4740 settle().await;
4741 assert!(
4742 dynamic.requested_track().now_or_never().is_none(),
4743 "a fetch inside the linger must reuse the track, not re-request it",
4744 );
4745 drop(fetching.await.expect("second fetch"));
4746
4747 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4749 settle().await;
4750 assert!(
4751 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4752 "the copy must be released once the fetches stop",
4753 );
4754 drop(producer);
4755
4756 settle().await;
4758 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4759 let mut producer = accept_track(&mut dynamic, "video").await;
4760 producer.append_group().unwrap().finish().unwrap();
4761 settle().await;
4762 fetching.await.expect("fetch after the linger");
4763 }
4764
4765 #[tokio::test(start_paused = true)]
4769 async fn test_linger_reconnect_splices() {
4770 let origin = Info::new(Origin::random())
4771 .with_linger(Duration::from_secs(5))
4772 .produce();
4773 let consumer = origin.consume();
4774 let mut announced = consumer.announced();
4775
4776 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4777 let source = origin
4778 .create_broadcast("test", announce().with_hops(hops.clone()))
4779 .unwrap();
4780 let mut dynamic = source.dynamic();
4781 settle().await;
4782 settle().await;
4783 let broadcast = consumer.request_broadcast("test").await.unwrap();
4784 announced.assert_next_some("test");
4785
4786 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4787 let mut producer = accept_track(&mut dynamic, "video").await;
4788 settle().await;
4789 let mut sub = subscribing.await.unwrap();
4790
4791 producer.append_group().unwrap();
4792 producer.append_group().unwrap();
4793 assert_eq!(sub.assert_group().sequence, 0);
4794 assert_eq!(sub.assert_group().sequence, 1);
4795
4796 drop(producer);
4799 source.abort(Error::Dropped).unwrap();
4800 drop(dynamic);
4801 settle().await;
4802
4803 announced.assert_next_wait();
4805 sub.assert_no_group();
4806 sub.assert_not_closed();
4807
4808 let during = consumer.request_broadcast("test").await.unwrap();
4810 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4811
4812 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4815 let mut dynamic = source.dynamic();
4816 settle().await;
4817 settle().await;
4818 announced.assert_next_wait();
4819 let again = consumer.request_broadcast("test").await.unwrap();
4820 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4821
4822 let mut producer = accept_track(&mut dynamic, "video").await;
4826 settle().await;
4827 sub.assert_no_group();
4828 assert_eq!(producer.subscription().unwrap().group_start, None);
4829 producer.create_group(group::Info { sequence: 2 }).unwrap();
4830 assert_eq!(sub.assert_group().sequence, 2);
4831 sub.assert_not_closed();
4832 }
4833
4834 #[tokio::test(start_paused = true)]
4837 async fn test_linger_expiry_closes() {
4838 let origin = Info::new(Origin::random())
4839 .with_linger(Duration::from_secs(5))
4840 .produce();
4841 let consumer = origin.consume();
4842 let mut announced = consumer.announced();
4843
4844 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4845 let source = origin
4846 .create_broadcast("test", announce().with_hops(hops.clone()))
4847 .unwrap();
4848 let mut dynamic = source.dynamic();
4849 settle().await;
4850 settle().await;
4851 let broadcast = consumer.request_broadcast("test").await.unwrap();
4852 announced.assert_next_some("test");
4853
4854 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4855 let producer = accept_track(&mut dynamic, "video").await;
4856 settle().await;
4857 let mut sub = subscribing.await.unwrap();
4858
4859 drop(producer);
4860 source.abort(Error::Dropped).unwrap();
4861 drop(dynamic);
4862 settle().await;
4863 announced.assert_next_wait();
4864
4865 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4867 settle().await;
4868 announced.assert_next_none("test");
4869 sub.assert_error();
4870
4871 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4873 settle().await;
4874 settle().await;
4875 let fresh = consumer.request_broadcast("test").await.unwrap();
4876 announced.assert_next_some("test");
4877 assert!(
4878 !fresh.is_clone(&broadcast),
4879 "a late re-create must not splice the expired broadcast"
4880 );
4881 }
4882
4883 #[tokio::test(start_paused = true)]
4887 async fn test_linger_forever() {
4888 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4889 let consumer = origin.consume();
4890 let mut announced = consumer.announced();
4891
4892 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4893 let source = origin
4894 .create_broadcast("test", announce().with_hops(hops.clone()))
4895 .unwrap();
4896 settle().await;
4897 let broadcast = consumer.request_broadcast("test").await.unwrap();
4898 announced.assert_next_some("test");
4899
4900 source.abort(Error::Dropped).unwrap();
4901 settle().await;
4902
4903 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4905 announced.assert_next_wait();
4906 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4907 settle().await;
4908 settle().await;
4909 let again = consumer.request_broadcast("test").await.unwrap();
4910 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4911 drop(source);
4912 }
4913
4914 #[tokio::test(start_paused = true)]
4928 async fn test_linger_parks_a_live_subscription() {
4929 let origin = Info::new(Origin::random())
4930 .with_linger(Duration::from_secs(5))
4931 .produce();
4932 let consumer = origin.consume();
4933
4934 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4935 let source = origin
4936 .create_broadcast("test", announce().with_hops(hops.clone()))
4937 .unwrap();
4938 let mut dynamic = source.dynamic();
4939 settle().await;
4940 settle().await;
4941 let broadcast = consumer.request_broadcast("test").await.unwrap();
4942
4943 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4946 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4947 settle().await;
4948 let mut sub = subscribing.await.unwrap();
4949 producer.append_group().unwrap();
4950 assert_eq!(sub.assert_group().sequence, 0);
4951
4952 source.abort(Error::Dropped).unwrap();
4957 settle().await;
4958 settle().await;
4959 sub.assert_not_closed();
4960
4961 tokio::time::sleep(Duration::from_secs(4)).await;
4964 settle().await;
4965 sub.assert_not_closed();
4966
4967 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4969 let mut dynamic = source.dynamic();
4970 settle().await;
4971 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4972 settle().await;
4973 producer.create_group(group::Info { sequence: 1 }).unwrap();
4974 assert_eq!(sub.assert_group().sequence, 1);
4975 sub.assert_not_closed();
4976 }
4977
4978 #[tokio::test(start_paused = true)]
4989 async fn test_idle_release_survives_the_route_leaving() {
4990 let origin = Info::new(Origin::random())
4991 .with_linger(Duration::from_secs(600))
4992 .produce();
4993 let consumer = origin.consume();
4994
4995 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4996 let source = origin
4997 .create_broadcast("test", announce().with_hops(hops.clone()))
4998 .unwrap();
4999 let mut dynamic = source.dynamic();
5000 settle().await;
5001 settle().await;
5002 let broadcast = consumer.request_broadcast("test").await.unwrap();
5003
5004 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
5005 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
5006 settle().await;
5007 let mut sub = subscribing.await.unwrap();
5008 producer.append_group().unwrap();
5009 assert_eq!(sub.assert_group().sequence, 0);
5010
5011 source.abort(Error::Dropped).unwrap();
5014 settle().await;
5015 settle().await;
5016 drop(sub);
5017 drop(producer);
5018 drop(dynamic);
5019 settle().await;
5020 tokio::time::sleep(TRACK_IDLE_LINGER + Duration::from_secs(1)).await;
5021 settle().await;
5022
5023 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5026 let mut dynamic = source.dynamic();
5027 settle().await;
5028 settle().await;
5029 let broadcast = consumer.request_broadcast("test").await.unwrap();
5030 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
5031 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
5032 settle().await;
5033 let mut sub = subscribing.await.unwrap();
5034 producer.append_group().unwrap();
5035 assert_eq!(
5036 sub.assert_group().sequence,
5037 0,
5038 "the reconnect's first group must not be filtered by a stale boundary"
5039 );
5040 }
5041
5042 #[tokio::test(start_paused = true)]
5045 async fn test_linger_skipped_on_finish() {
5046 let origin = Info::new(Origin::random())
5047 .with_linger(Duration::from_secs(5))
5048 .produce();
5049 let consumer = origin.consume();
5050 let mut announced = consumer.announced();
5051
5052 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5053 let mut source = origin
5054 .create_broadcast("test", announce().with_hops(hops.clone()))
5055 .unwrap();
5056 settle().await;
5057 let broadcast = consumer.request_broadcast("test").await.unwrap();
5058 announced.assert_next_some("test");
5059
5060 source.finish();
5063 settle().await;
5064 announced.assert_next_none("test");
5065
5066 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5068 settle().await;
5069 let fresh = consumer.request_broadcast("test").await.unwrap();
5070 announced.assert_next_some("test");
5071 assert!(
5072 !fresh.is_clone(&broadcast),
5073 "a finish must not leave a lingering broadcast to splice into"
5074 );
5075 }
5076
5077 #[tokio::test]
5080 async fn test_announce_toggle() {
5081 tokio::time::pause();
5082
5083 let origin = Origin::random().produce();
5084 let consumer = origin.consume();
5085 let mut announced = consumer.announced();
5086
5087 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
5088 settle().await;
5089
5090 announced.assert_next_wait();
5092 let broadcast = consumer
5093 .get_broadcast("test")
5094 .expect("offline broadcast is still routable");
5095 assert!(!broadcast.route().announce);
5096
5097 let requested = consumer.request_broadcast("test").await.unwrap();
5099 assert!(requested.is_clone(&broadcast));
5100
5101 source.set_route(announce()).unwrap();
5103 settle().await;
5104 let face = announced.assert_next_some("test");
5105 assert!(face.is_clone(&broadcast));
5106
5107 let mut fresh = origin.consume().announced();
5109 fresh.assert_next_some("test");
5110 fresh.assert_next_wait();
5111
5112 source.set_route(broadcast::Route::new()).unwrap();
5114 settle().await;
5115 announced.assert_next_none("test");
5116 assert!(consumer.get_broadcast("test").is_some());
5117 let mut fresh = origin.consume().announced();
5118 fresh.assert_next_wait();
5119
5120 source.finish();
5121 settle().await;
5122 assert!(consumer.get_broadcast("test").is_none());
5123 }
5124
5125 #[tokio::test]
5128 async fn test_announce_beats_offline() {
5129 tokio::time::pause();
5130
5131 let origin = Origin::random().produce();
5132 let consumer = origin.consume();
5133 let mut announced = consumer.announced();
5134
5135 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
5137 settle().await;
5138 announced.assert_next_wait();
5139
5140 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
5143 settle().await;
5144 announced.assert_next_some("test");
5145 let face = consumer.get_broadcast("test").unwrap();
5146 assert!(face.route().announce);
5147 assert_eq!(face.route().cost, 10);
5148
5149 announced_source.finish();
5152 settle().await;
5153 announced.assert_next_none("test");
5154 assert!(consumer.get_broadcast("test").is_some());
5155 }
5156
5157 #[tokio::test]
5160 async fn test_better_source_no_churn() {
5161 tokio::time::pause();
5162
5163 let origin = Origin::random().produce();
5164 let mut announced = origin.consume().announced();
5165
5166 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
5169 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5170 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
5171 settle().await;
5172 let face = announced.assert_next_some("test");
5173
5174 let _b = origin
5175 .create_broadcast("test", announce().with_hops(hops_b.clone()))
5176 .unwrap();
5177 settle().await;
5178 announced.assert_next_wait();
5179 let current = origin.consume().get_broadcast("test").unwrap();
5180 assert!(current.is_clone(&face), "the broadcast identity must not change");
5181 assert_eq!(current.route().hops, hops_b);
5183 }
5184
5185 #[tokio::test]
5191 async fn test_publisher_mismatch_replaces() {
5192 tokio::time::pause();
5193
5194 let origin = Origin::random().produce();
5195 let consumer = origin.consume();
5196 let mut announced = consumer.announced();
5197
5198 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5199 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5200
5201 let mut source_a = origin
5202 .create_broadcast("test", announce().with_hops(hops_a.clone()))
5203 .unwrap();
5204 settle().await;
5205 let face_a = announced.assert_next_some("test");
5206
5207 let _source_b = origin
5210 .create_broadcast("test", announce().with_hops(hops_b.clone()))
5211 .unwrap();
5212 settle().await;
5213 settle().await;
5214 announced.assert_next_none("test");
5215 let face_b = announced.assert_next_some("test");
5216 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
5217 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
5218 assert!(face_a.is_closed(), "the displaced front must close");
5222
5223 source_a.finish();
5225 settle().await;
5226 settle().await;
5227 announced.assert_next_wait();
5228 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
5229 }
5230
5231 #[tokio::test]
5234 async fn test_displaced_publisher_reclaims_path_when_replacement_leaves() {
5235 tokio::time::pause();
5236
5237 let origin = Origin::random().produce();
5238 let consumer = origin.consume();
5239 let mut announced = consumer.announced();
5240
5241 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5242 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5243
5244 let mut source_a = origin
5246 .create_broadcast("test", announce().with_hops(hops_a.clone()))
5247 .unwrap();
5248 settle().await;
5249 announced.assert_next_some("test");
5250
5251 let mut source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
5254 settle().await;
5255 settle().await;
5256 announced.assert_next_none("test");
5257 announced.assert_next_some("test");
5258 announced.assert_next_wait();
5261
5262 source_b.finish();
5264 settle().await;
5265 settle().await;
5266
5267 assert!(
5268 consumer.get_broadcast("test").is_some(),
5269 "the still-live publisher A should reclaim the path once its replacement leaves"
5270 );
5271 let recovered = consumer.request_broadcast("test").await.unwrap();
5272 assert_eq!(recovered.route().hops, hops_a);
5273
5274 announced.assert_next_none("test");
5275 announced.assert_next_some("test");
5276 announced.assert_next_wait();
5277
5278 source_a.finish();
5279 }
5280
5281 #[tokio::test]
5285 async fn test_reclaimed_publisher_can_still_replace_itself() {
5286 tokio::time::pause();
5287
5288 let origin = Origin::random().produce();
5289 let consumer = origin.consume();
5290
5291 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5292 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5293 let hops_c = OriginList::try_from(vec![Origin::new(3).unwrap()]).unwrap();
5294
5295 let mut source_a1 = origin
5297 .create_broadcast("test", announce().with_hops(hops_a.clone()))
5298 .unwrap();
5299 let mut source_a2 = origin
5300 .create_broadcast("test", announce().with_hops(hops_a.clone()))
5301 .unwrap();
5302 settle().await;
5303 settle().await;
5304 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_a);
5305
5306 let mut source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
5308 settle().await;
5309 settle().await;
5310
5311 source_b.finish();
5313 settle().await;
5314 settle().await;
5315 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_a);
5316
5317 source_a1.set_route(announce().with_hops(hops_c.clone())).unwrap();
5321 settle().await;
5322 settle().await;
5323 assert_eq!(
5324 consumer.get_broadcast("test").unwrap().route().hops,
5325 hops_c,
5326 "the new publisher must take the path over, not stand by behind the old front"
5327 );
5328
5329 source_a1.finish();
5330 source_a2.finish();
5331 }
5332
5333 #[tokio::test]
5336 async fn test_repricing_does_not_earn_a_takeover() {
5337 tokio::time::pause();
5338
5339 let origin = Origin::random().produce();
5340 let consumer = origin.consume();
5341 let mut announced = consumer.announced();
5342
5343 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5344 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5345
5346 let mut source_a = origin
5347 .create_broadcast("test", announce().with_hops(hops_a.clone()).with_cost(5))
5348 .unwrap();
5349 settle().await;
5350 announced.assert_next_some("test");
5351
5352 let mut source_b = origin
5354 .create_broadcast("test", announce().with_hops(hops_b.clone()))
5355 .unwrap();
5356 settle().await;
5357 settle().await;
5358 announced.assert_next_none("test");
5359 announced.assert_next_some("test");
5360 announced.assert_next_wait();
5361
5362 source_a
5365 .set_route(announce().with_hops(hops_a.clone()).with_cost(9))
5366 .unwrap();
5367 settle().await;
5368 settle().await;
5369 assert_eq!(
5370 consumer.get_broadcast("test").unwrap().route().hops,
5371 hops_b,
5372 "a repricing must not take the path back from the live front"
5373 );
5374 announced.assert_next_wait();
5375
5376 source_a.finish();
5377 source_b.finish();
5378 }
5379
5380 #[tokio::test]
5385 async fn test_reconnect_wins_over_stale_route() {
5386 tokio::time::pause();
5387
5388 let origin = Origin::random().produce();
5389 let consumer = origin.consume();
5390
5391 let publisher = Origin::new(1).unwrap();
5392 let hops = OriginList::try_from(vec![publisher]).unwrap();
5393
5394 let stale = origin
5397 .create_broadcast("test", announce().with_hops(hops.clone()))
5398 .unwrap();
5399 let mut stale_dynamic = stale.dynamic();
5400 settle().await;
5401
5402 let fresh = origin
5404 .create_broadcast("test", announce().with_hops(hops.clone()))
5405 .unwrap();
5406 let mut fresh_dynamic = fresh.dynamic();
5407 settle().await;
5408 settle().await;
5409
5410 let broadcast = consumer.request_broadcast("test").await.unwrap();
5412 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5413 settle().await;
5414 let _producer = accept_track(&mut fresh_dynamic, "video").await;
5415 settle().await;
5416 subscribing.await.unwrap();
5417 stale_dynamic.assert_no_request();
5418 }
5419
5420 #[tokio::test]
5426 async fn test_carrying_reconnect_switches_immediately() {
5427 tokio::time::pause();
5428
5429 let origin = Origin::random().produce();
5430 let consumer = origin.consume();
5431
5432 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5433
5434 let stale = origin
5435 .create_broadcast("test", announce().with_hops(hops.clone()))
5436 .unwrap();
5437 let mut stale_dynamic = stale.dynamic();
5438 settle().await;
5439
5440 let broadcast = consumer.request_broadcast("test").await.unwrap();
5442 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5443 settle().await;
5444 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
5445 settle().await;
5446 let _subscription = subscribing.await.unwrap();
5448
5449 let fresh = origin
5451 .create_broadcast("test", announce().with_hops(hops.clone()))
5452 .unwrap();
5453 let mut fresh_dynamic = fresh.dynamic();
5454 settle().await;
5455 settle().await;
5456
5457 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
5460 }
5461
5462 #[tokio::test]
5472 async fn test_offline_mismatch_never_evicts_a_live_front() {
5473 tokio::time::pause();
5474
5475 let origin = Origin::random().produce();
5476 let consumer = origin.consume();
5477 let mut announced = consumer.announced();
5478
5479 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5480 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
5481
5482 let mut live = origin
5483 .create_broadcast("test", announce().with_hops(hops_live.clone()))
5484 .unwrap();
5485 let mut live_dynamic = live.dynamic();
5486 settle().await;
5487 let face = announced.assert_next_some("test");
5488
5489 let broadcast = consumer.request_broadcast("test").await.unwrap();
5491 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5492 settle().await;
5493 let _producer = accept_track(&mut live_dynamic, "video").await;
5494 settle().await;
5495 subscribing.await.unwrap();
5496
5497 let cache = origin
5500 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
5501 .unwrap();
5502 settle().await;
5503 settle().await;
5504 announced.assert_next_wait();
5505 assert!(!face.is_closed(), "the live front must survive");
5506 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
5507
5508 live.finish();
5511 settle().await;
5512 settle().await;
5513 announced.assert_next_none("test");
5514 let taken = consumer
5515 .get_broadcast("test")
5516 .expect("the parked source must take over");
5517 assert_eq!(taken.route().hops, hops_cache);
5518 announced.assert_next_wait();
5520 drop(cache);
5521 }
5522
5523 #[tokio::test]
5529 async fn test_dispatch_excludes_requester() {
5530 tokio::time::pause();
5531
5532 let origin = Origin::random().produce();
5533 let consumer = origin.consume();
5534
5535 let peer = Origin::new(5).unwrap();
5536 let publisher = Origin::new(1).unwrap();
5537 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
5539 let clean = OriginList::try_from(vec![publisher]).unwrap();
5540
5541 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5542 let mut dynamic_a = source_a.dynamic();
5543 settle().await;
5544 let source_b = origin
5545 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
5546 .unwrap();
5547 let mut dynamic_b = source_b.dynamic();
5548 settle().await;
5549 settle().await;
5550
5551 let shared = consumer.request_broadcast("test").await.unwrap();
5554 let subscribing = shared.track("video").unwrap().subscribe(None);
5555 let _producer_a = accept_track(&mut dynamic_a, "video").await;
5556 settle().await;
5557 subscribing.await.unwrap();
5558
5559 let scoped = consumer.clone().excluding(peer);
5564 let pinned = scoped.request_broadcast("test").await.unwrap();
5565 let subscribing = pinned.track("video").unwrap().subscribe(None);
5566 let _producer_b = accept_track(&mut dynamic_b, "video").await;
5567 settle().await;
5568 subscribing.await.unwrap();
5569 dynamic_a.assert_no_request();
5570 }
5571
5572 #[tokio::test]
5577 async fn test_unknown_publishers_do_not_splice() {
5578 tokio::time::pause();
5579
5580 let origin = Origin::random().produce();
5581 let consumer = origin.consume();
5582 let mut announced = consumer.announced();
5583
5584 let unknown_a = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
5585 let unknown_b = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
5586
5587 let mut source_a = origin
5588 .create_broadcast("test", announce().with_hops(unknown_a.clone()))
5589 .unwrap();
5590 settle().await;
5591 settle().await;
5592 announced.assert_next_some("test");
5593
5594 let source_b = origin
5598 .create_broadcast("test", announce().with_hops(unknown_b))
5599 .unwrap();
5600 settle().await;
5601 settle().await;
5602 announced.assert_next_none("test");
5603 let live = announced.assert_next_some("test");
5604
5605 source_a
5608 .set_route(announce().with_hops(unknown_a).with_cost(9))
5609 .unwrap();
5610 settle().await;
5611 settle().await;
5612 assert!(
5613 consumer.get_broadcast("test").unwrap().is_clone(&live),
5614 "UNKNOWN-to-UNKNOWN repricing must not replace the live front"
5615 );
5616 announced.assert_next_wait();
5617
5618 drop(source_a);
5619 drop(source_b);
5620 }
5621
5622 #[tokio::test]
5625 async fn test_known_publishers_still_splice() {
5626 tokio::time::pause();
5627
5628 let origin = Origin::random().produce();
5629 let consumer = origin.consume();
5630 let mut announced = consumer.announced();
5631
5632 let publisher = Origin::new(1).unwrap();
5633 let hops_a = OriginList::try_from(vec![publisher]).unwrap();
5634 let hops_b = OriginList::try_from(vec![publisher, Origin::new(3).unwrap()]).unwrap();
5635
5636 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
5637 settle().await;
5638 settle().await;
5639 announced.assert_next_some("test");
5640
5641 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
5643 settle().await;
5644 settle().await;
5645 announced.assert_next_wait();
5646
5647 drop(source_a);
5648 drop(source_b);
5649 }
5650
5651 #[tokio::test]
5657 async fn test_standby_join_splices_live_subscriber() {
5658 tokio::time::pause();
5659
5660 let origin = Origin::random().produce();
5661 let consumer = origin.consume();
5662
5663 let publisher = Origin::new(1).unwrap();
5664 let peer = Origin::new(5).unwrap();
5665 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5666 let local = OriginList::try_from(vec![publisher]).unwrap();
5667
5668 let source_remote = origin
5670 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5671 .unwrap();
5672 let mut dynamic_remote = source_remote.dynamic();
5673 settle().await;
5674 settle().await;
5675 let broadcast = consumer.request_broadcast("test").await.unwrap();
5676 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5677 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
5678 settle().await;
5679 let mut sub = subscribing.await.unwrap();
5680 producer_remote.append_group().unwrap();
5681 assert_eq!(sub.assert_group().sequence, 0);
5682
5683 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5687 let mut dynamic_local = source_local.dynamic();
5688 settle().await;
5689 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
5690 settle().await;
5691 sub.assert_no_group();
5692 assert_eq!(producer_local.subscription().unwrap().group_start, None);
5693 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
5694 assert_eq!(sub.assert_group().sequence, 1);
5695 sub.assert_not_closed();
5696 }
5697
5698 #[tokio::test]
5705 async fn test_standby_with_a_partial_track_list_splits_per_track() {
5706 tokio::time::pause();
5707
5708 let origin = Origin::random().produce();
5709 let consumer = origin.consume();
5710
5711 let publisher = Origin::new(1).unwrap();
5712 let peer = Origin::new(5).unwrap();
5713 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5714 let local = OriginList::try_from(vec![publisher]).unwrap();
5715
5716 let source_remote = origin
5718 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5719 .unwrap();
5720 let mut dynamic_remote = source_remote.dynamic();
5721 settle().await;
5722 settle().await;
5723 let broadcast = consumer.request_broadcast("test").await.unwrap();
5724
5725 let subscribing_video = broadcast.track("video").unwrap().subscribe(None);
5726 let subscribing_audio = broadcast.track("audio").unwrap().subscribe(None);
5727 let mut producers_remote = accept_tracks(&mut dynamic_remote, 2).await;
5728 settle().await;
5729
5730 let mut sub_video = subscribing_video.await.unwrap();
5731 let mut sub_audio = subscribing_audio.await.unwrap();
5732 for name in ["video", "audio"] {
5733 producers_remote.get_mut(name).unwrap().append_group().unwrap();
5734 }
5735 assert_eq!(sub_video.assert_group().sequence, 0);
5736 assert_eq!(sub_audio.assert_group().sequence, 0);
5737
5738 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5740 let mut dynamic_local = source_local.dynamic();
5741 settle().await;
5742
5743 let mut producer_local = None;
5744 for _ in 0..2 {
5745 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic_local.requested_track())
5746 .await
5747 .expect("timed out waiting for a track request")
5748 .expect("source closed");
5749 match request.name() {
5750 "video" => producer_local = Some(request.accept(None)),
5751 "audio" => request.reject(Error::NotFound),
5752 other => panic!("unexpected track dispatched: {other}"),
5753 }
5754 }
5755 settle().await;
5756 let mut producer_local = producer_local.expect("the standby was never asked for video");
5757
5758 sub_video.assert_no_group();
5761 assert_eq!(producer_local.subscription().unwrap().group_start, None);
5762 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
5763 assert_eq!(
5764 sub_video.assert_group().sequence,
5765 1,
5766 "video did not move to the standby"
5767 );
5768
5769 producers_remote.get_mut("audio").unwrap().append_group().unwrap();
5771 assert_eq!(
5772 sub_audio.assert_group().sequence,
5773 1,
5774 "audio did not stay on the incumbent"
5775 );
5776 sub_video.assert_not_closed();
5777 sub_audio.assert_not_closed();
5778 }
5779
5780 #[tokio::test]
5787 async fn test_standby_missing_track_keeps_incumbent() {
5788 tokio::time::pause();
5789
5790 let origin = Origin::random().produce();
5791 let consumer = origin.consume();
5792
5793 let publisher = Origin::new(1).unwrap();
5794 let peer = Origin::new(5).unwrap();
5795 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5796 let local = OriginList::try_from(vec![publisher]).unwrap();
5797
5798 let source_remote = origin
5800 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5801 .unwrap();
5802 let mut dynamic_remote = source_remote.dynamic();
5803 settle().await;
5804 settle().await;
5805 let broadcast = consumer.request_broadcast("test").await.unwrap();
5806 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5807 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
5808 settle().await;
5809 let mut sub = subscribing.await.unwrap();
5810 producer_remote.append_group().unwrap();
5811 assert_eq!(sub.assert_group().sequence, 0);
5812
5813 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5816 let mut dynamic_local = source_local.dynamic();
5817 settle().await;
5818 let request = dynamic_local.requested_track().await.unwrap();
5819 assert_eq!(request.name(), "audio");
5820 request.reject(Error::NotFound);
5821 settle().await;
5822
5823 producer_remote.append_group().unwrap();
5825 assert_eq!(sub.assert_group().sequence, 1);
5826 sub.assert_not_closed();
5827
5828 source_remote.abort(Error::Dropped).unwrap();
5831 settle().await;
5832 settle().await;
5833 sub.assert_closed();
5834 dynamic_local.assert_no_request();
5835
5836 let retry = broadcast.track("audio").unwrap().subscribe(None);
5838 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
5839 settle().await;
5840 let mut sub = retry.await.expect("a fresh request must reach the standby");
5841 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
5842 assert_eq!(sub.assert_group().sequence, 2);
5843 }
5844
5845 #[tokio::test]
5850 async fn test_unservable_track_retried_by_a_later_request() {
5851 tokio::time::pause();
5852
5853 let origin = Origin::random().produce();
5854 let consumer = origin.consume();
5855
5856 let source = origin.create_broadcast("test", announce()).unwrap();
5857 let mut dynamic = source.dynamic();
5858 settle().await;
5859 settle().await;
5860 let broadcast = consumer.request_broadcast("test").await.unwrap();
5861
5862 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
5864 let request = dynamic.requested_track().await.unwrap();
5865 request.reject(Error::NotFound);
5866 settle().await;
5867 assert!(matches!(subscribing.await, Err(Error::NotFound)));
5868
5869 let retry = broadcast.track("audio").unwrap().subscribe(None);
5871 let mut producer = accept_track(&mut dynamic, "audio").await;
5872 settle().await;
5873 let mut sub = retry.await.expect("a fresh request must reach the source");
5874 producer.append_group().unwrap();
5875 assert_eq!(sub.assert_group().sequence, 0);
5876 }
5877
5878 #[tokio::test]
5884 async fn test_track_dying_without_progress_aborts() {
5885 tokio::time::pause();
5886
5887 let origin = Origin::random().produce();
5888 let consumer = origin.consume();
5889
5890 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5891 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5892 let mut dynamic = source.dynamic();
5893 settle().await;
5894 settle().await;
5895 let broadcast = consumer.request_broadcast("test").await.unwrap();
5896
5897 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5898 let producer = accept_track(&mut dynamic, "video").await;
5899 settle().await;
5900 let mut sub = subscribing.await.unwrap();
5901
5902 drop(producer);
5905 settle().await;
5906 sub.assert_closed();
5907 dynamic.assert_no_request();
5908
5909 let retry = broadcast.track("video").unwrap().subscribe(None);
5911 let mut producer = accept_track(&mut dynamic, "video").await;
5912 settle().await;
5913 let mut sub = retry.await.expect("a fresh request must reach the source");
5914 producer.append_group().unwrap();
5915 assert_eq!(sub.assert_group().sequence, 0);
5916 }
5917
5918 #[tokio::test]
5923 async fn test_delivered_copy_death_survives_unrelated_wakes() {
5924 tokio::time::pause();
5925
5926 let origin = Origin::random().produce();
5927 let consumer = origin.consume();
5928
5929 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5930 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5931 let mut dynamic = source.dynamic();
5932 settle().await;
5933 settle().await;
5934 let broadcast = consumer.request_broadcast("test").await.unwrap();
5935
5936 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5937 let mut producer = accept_track(&mut dynamic, "video").await;
5938 settle().await;
5939 let mut sub = subscribing.await.unwrap();
5940 producer.append_group().unwrap();
5941 assert_eq!(sub.assert_group().sequence, 0);
5942
5943 drop(sub);
5946 settle().await;
5947 let resubscribing = broadcast.track("video").unwrap().subscribe(None);
5948 settle().await;
5949 let mut sub = resubscribing.await.unwrap();
5950 assert_eq!(sub.assert_group().sequence, 0, "cached group re-served");
5951
5952 drop(producer);
5955 let mut producer = accept_track(&mut dynamic, "video").await;
5956 settle().await;
5957 producer.create_group(group::Info { sequence: 1 }).unwrap();
5958 assert_eq!(sub.assert_group().sequence, 1);
5959 sub.assert_not_closed();
5960 }
5961
5962 #[tokio::test]
5968 async fn test_per_track_fallback_respects_exclusion() {
5969 tokio::time::pause();
5970
5971 let origin = Origin::random().produce();
5972 let consumer = origin.consume();
5973
5974 let publisher = Origin::new(1).unwrap();
5975 let peer = Origin::new(5).unwrap();
5976 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5977 let local = OriginList::try_from(vec![publisher]).unwrap();
5978
5979 let source_tainted = origin
5981 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5982 .unwrap();
5983 let mut dynamic_tainted = source_tainted.dynamic();
5984 settle().await;
5985 settle().await;
5986
5987 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5990 let mut dynamic_clean = source_clean.dynamic();
5991 settle().await;
5992
5993 let scoped = consumer.clone().excluding(peer);
5994 let broadcast = scoped.request_broadcast("test").await.unwrap();
5995 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5996 settle().await;
5997
5998 let request = dynamic_clean.requested_track().await.unwrap();
6002 request.reject(Error::NotFound);
6003 settle().await;
6004 dynamic_tainted.assert_no_request();
6005 }
6006
6007 #[tokio::test]
6011 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
6012 tokio::time::pause();
6013
6014 let origin = Origin::random().produce();
6015 let consumer = origin.consume();
6016
6017 let publisher = Origin::new(1).unwrap();
6018 let peer = Origin::new(5).unwrap();
6019 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
6020 let local = OriginList::try_from(vec![publisher]).unwrap();
6021
6022 let source_tainted = origin
6023 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
6024 .unwrap();
6025 let mut dynamic_tainted = source_tainted.dynamic();
6026 settle().await;
6027 settle().await;
6028 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
6029 let mut dynamic_clean = source_clean.dynamic();
6030 settle().await;
6031
6032 let scoped = consumer.clone().excluding(peer);
6033 let broadcast = scoped.request_broadcast("test").await.unwrap();
6034 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
6035 let _clean = accept_track(&mut dynamic_clean, "video").await;
6036 settle().await;
6037
6038 source_clean.abort(Error::Dropped).unwrap();
6040 settle().await;
6041 settle().await;
6042 dynamic_tainted.assert_no_request();
6043
6044 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
6047 }
6048
6049 #[tokio::test]
6054 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
6055 tokio::time::pause();
6056
6057 let origin = Origin::random().produce();
6058 let consumer = origin.consume();
6059
6060 let publisher = Origin::new(1).unwrap();
6061 let peer = Origin::new(5).unwrap();
6062 let local = OriginList::try_from(vec![publisher]).unwrap();
6063 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
6064
6065 let source_clean = origin
6069 .create_broadcast("test", announce().with_hops(local).with_cost(5))
6070 .unwrap();
6071 let mut dynamic_clean = source_clean.dynamic();
6072 settle().await;
6073 settle().await;
6074 let scoped = consumer.clone().excluding(peer);
6075 let broadcast = scoped.request_broadcast("test").await.unwrap();
6076 let subscribing = broadcast.track("video").unwrap().subscribe(None);
6077 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
6078 settle().await;
6079 let mut sub = subscribing.await.unwrap();
6080 producer_clean.append_group().unwrap();
6081 assert_eq!(sub.assert_group().sequence, 0);
6082
6083 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
6090 tainted.advertised = 1;
6091 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
6092 let mut dynamic_tainted = source_tainted.dynamic();
6093 settle().await;
6094 settle().await;
6095 dynamic_tainted.assert_no_request();
6096 producer_clean.append_group().unwrap();
6097 assert_eq!(sub.assert_group().sequence, 1);
6098 sub.assert_not_closed();
6099
6100 drop(sub);
6103 drop(broadcast);
6104 drop(scoped);
6105 settle().await;
6106 let mut bumped = announce().with_hops(via_peer).with_cost(1);
6107 bumped.advertised = 1;
6108 source_tainted.set_route(bumped).unwrap();
6109 settle().await;
6110 let plain = consumer.request_broadcast("test").await.unwrap();
6111 let _plain_track = plain.track("video").unwrap().subscribe(None);
6112 settle().await;
6113 settle().await;
6114 assert!(
6115 dynamic_tainted.requested_track().now_or_never().is_some(),
6116 "the front must be free to use the route again once the peer is gone"
6117 );
6118 }
6119
6120 #[tokio::test]
6124 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
6125 tokio::time::pause();
6126
6127 let origin = Origin::random().produce();
6128 let consumer = origin.consume();
6129 let mut dynamic = origin.dynamic();
6130
6131 let peer = Origin::new(5).unwrap();
6132 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
6133 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
6134 settle().await;
6135 settle().await;
6136
6137 let scoped = consumer.clone().excluding(peer);
6138 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
6139 assert!(
6140 dynamic.requested_broadcast().now_or_never().is_none(),
6141 "the dynamic handler was asked to route around the exclusion"
6142 );
6143
6144 let _pending = scoped.request_broadcast("other");
6146 settle().await;
6147 assert!(
6148 dynamic.requested_broadcast().now_or_never().is_some(),
6149 "a genuinely missing path must still fall back"
6150 );
6151 }
6152
6153 #[tokio::test]
6156 async fn test_dispatch_all_tainted_unroutable() {
6157 tokio::time::pause();
6158
6159 let origin = Origin::random().produce();
6160 let consumer = origin.consume();
6161
6162 let peer = Origin::new(5).unwrap();
6163 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
6164 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
6165 settle().await;
6166 settle().await;
6167
6168 let scoped = consumer.clone().excluding(peer);
6169 match scoped.request_broadcast("test").await {
6170 Err(Error::Unroutable) => {}
6171 Err(err) => panic!("expected Unroutable, got {err:?}"),
6172 Ok(_) => panic!("expected Unroutable, got a broadcast"),
6173 }
6174
6175 consumer.request_broadcast("test").await.unwrap();
6177 }
6178
6179 #[tokio::test]
6180 async fn test_duplicate_reverse() {
6181 tokio::time::pause();
6182
6183 let origin = Origin::random().produce();
6184
6185 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6186 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6187 settle().await;
6188 assert!(origin.consume().get_broadcast("test").is_some());
6189
6190 broadcast2.finish();
6192 settle().await;
6193 assert!(origin.consume().get_broadcast("test").is_some());
6194
6195 broadcast1.finish();
6196 settle().await;
6197 assert!(origin.consume().get_broadcast("test").is_none());
6198 }
6199
6200 #[tokio::test]
6201 async fn test_deterministic_tiebreak() {
6202 tokio::time::pause();
6203
6204 fn hops(ids: &[u64]) -> OriginList {
6205 OriginList::try_from(
6206 ids.iter()
6207 .copied()
6208 .map(|id| Origin::new(id).unwrap())
6209 .collect::<Vec<_>>(),
6210 )
6211 .unwrap()
6212 }
6213
6214 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
6217 let origin = Origin::random().produce();
6218 let _a = origin
6219 .create_broadcast("test", announce().with_hops(hops(first)))
6220 .unwrap();
6221 let _b = origin
6222 .create_broadcast("test", announce().with_hops(hops(second)))
6223 .unwrap();
6224 settle().await;
6225 origin.consume().get_broadcast("test").unwrap().route().hops
6226 }
6227
6228 let forward = winner(&[5, 20], &[5, 40]).await;
6232 let reverse = winner(&[5, 40], &[5, 20]).await;
6233 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
6234
6235 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
6237 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
6238 }
6239
6240 #[tokio::test]
6245 async fn test_many_announces() {
6246 let origin = Origin::random().produce();
6247
6248 let mut consumer = origin.consume().announced();
6249 let mut broadcasts = Vec::new();
6251 for i in 0..256 {
6252 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
6253 settle().await;
6254 }
6255
6256 for i in 0..256 {
6257 consumer.assert_next_some(format!("test{i:03}"));
6258 }
6259 consumer.assert_next_wait();
6260 }
6261
6262 #[tokio::test]
6263 async fn test_many_announces_try() {
6264 let origin = Origin::random().produce();
6265
6266 let mut consumer = origin.consume().announced();
6267 let mut broadcasts = Vec::new();
6269 for i in 0..256 {
6270 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
6271 settle().await;
6272 }
6273
6274 for i in 0..256 {
6275 consumer.assert_try_next_some(format!("test{i:03}"));
6276 }
6277 }
6278
6279 #[tokio::test]
6280 async fn test_with_root_basic() {
6281 let origin = Origin::random().produce();
6282
6283 let foo_producer = origin.with_root("foo").expect("should create root");
6285 assert_eq!(foo_producer.root().as_str(), "foo");
6286
6287 let mut consumer = origin.consume().announced();
6288
6289 let _broadcast = foo_producer
6291 .create_broadcast("bar/baz", announce())
6292 .expect("publish allowed");
6293 settle().await;
6294 consumer.assert_next_some("foo/bar/baz");
6296
6297 let mut foo_consumer = foo_producer.consume().announced();
6299 foo_consumer.assert_next_some("bar/baz");
6300 }
6301
6302 #[tokio::test]
6303 async fn test_with_root_nested() {
6304 let origin = Origin::random().produce();
6305
6306 let foo_producer = origin.with_root("foo").expect("should create foo root");
6308 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
6309 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
6310
6311 let mut consumer = origin.consume().announced();
6312
6313 let _broadcast = foo_bar_producer
6315 .create_broadcast("baz", announce())
6316 .expect("publish allowed");
6317 settle().await;
6318 consumer.assert_next_some("foo/bar/baz");
6320
6321 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
6323 foo_bar_consumer.assert_next_some("baz");
6324 }
6325
6326 #[tokio::test]
6327 async fn test_publish_scope_allows() {
6328 let origin = Origin::random().produce();
6329
6330 let limited_producer = origin
6332 .scope(&["allowed/path1".into(), "allowed/path2".into()])
6333 .expect("should create limited producer");
6334
6335 let _broadcast = limited_producer
6337 .create_broadcast("allowed/path1", announce())
6338 .expect("publish allowed");
6339 let _keep2 = limited_producer
6340 .create_broadcast("allowed/path1/nested", announce())
6341 .expect("publish allowed");
6342 let _keep3 = limited_producer
6343 .create_broadcast("allowed/path2", announce())
6344 .expect("publish allowed");
6345 settle().await;
6346
6347 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
6349 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
6351 }
6352
6353 #[tokio::test]
6354 async fn test_publish_max_parts() {
6355 let origin = Origin::random().produce();
6356
6357 let at_limit = (0..Path::MAX_PARTS)
6358 .map(|i| i.to_string())
6359 .collect::<Vec<_>>()
6360 .join("/");
6361 let _broadcast = origin
6362 .create_broadcast(at_limit.as_str(), announce())
6363 .expect("publish allowed");
6364 settle().await;
6365
6366 let too_deep = format!("{at_limit}/extra");
6367 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
6368
6369 let rooted = origin.with_root("root").expect("wildcard allows any root");
6371 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
6372 }
6373
6374 #[tokio::test]
6375 async fn test_publish_scope_empty() {
6376 let origin = Origin::random().produce();
6377
6378 assert!(origin.scope(&[]).is_none());
6380 }
6381
6382 #[tokio::test]
6383 async fn test_consume_scope_filters() {
6384 let origin = Origin::random().produce();
6385
6386 let mut consumer = origin.consume().announced();
6387
6388 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
6390 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
6391 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
6392 settle().await;
6393
6394 let mut limited_consumer = origin
6396 .consume()
6397 .scope(&["allowed".into()])
6398 .expect("should create limited consumer")
6399 .announced();
6400
6401 limited_consumer.assert_next_some("allowed");
6403 limited_consumer.assert_next_some("allowed/nested");
6404 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
6408 consumer.assert_next_some("allowed/nested");
6409 consumer.assert_next_some("notallowed");
6410 }
6411
6412 #[tokio::test]
6413 async fn test_consume_scope_multiple_prefixes() {
6414 let origin = Origin::random().produce();
6415
6416 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
6417 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
6418 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
6419 settle().await;
6420
6421 let mut limited_consumer = origin
6423 .consume()
6424 .scope(&["foo".into(), "bar".into()])
6425 .expect("should create limited consumer")
6426 .announced();
6427
6428 limited_consumer.assert_next_some("bar/test");
6430 limited_consumer.assert_next_some("foo/test");
6431 limited_consumer.assert_next_wait(); }
6433
6434 #[tokio::test]
6435 async fn test_with_root_and_publish_scope() {
6436 let origin = Origin::random().produce();
6437
6438 let foo_producer = origin.with_root("foo").expect("should create foo root");
6440
6441 let limited_producer = foo_producer
6443 .scope(&["bar".into(), "goop/pee".into()])
6444 .expect("should create limited producer");
6445
6446 let mut consumer = origin.consume().announced();
6447
6448 let _broadcast = limited_producer
6450 .create_broadcast("bar", announce())
6451 .expect("publish allowed");
6452 let _keep2 = limited_producer
6453 .create_broadcast("bar/nested", announce())
6454 .expect("publish allowed");
6455 let _keep3 = limited_producer
6456 .create_broadcast("goop/pee", announce())
6457 .expect("publish allowed");
6458 let _keep4 = limited_producer
6459 .create_broadcast("goop/pee/nested", announce())
6460 .expect("publish allowed");
6461 settle().await;
6462
6463 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
6465 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
6467
6468 consumer.assert_next_some("foo/bar");
6470 consumer.assert_next_some("foo/bar/nested");
6471 consumer.assert_next_some("foo/goop/pee");
6472 consumer.assert_next_some("foo/goop/pee/nested");
6473 }
6474
6475 #[tokio::test]
6476 async fn test_with_root_and_consume_scope() {
6477 let origin = Origin::random().produce();
6478
6479 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
6481 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
6482 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
6483 settle().await;
6484
6485 let foo_producer = origin.with_root("foo").expect("should create foo root");
6487
6488 let mut limited_consumer = foo_producer
6490 .consume()
6491 .scope(&["bar".into(), "goop/pee".into()])
6492 .expect("should create limited consumer")
6493 .announced();
6494
6495 limited_consumer.assert_next_some("bar/test");
6497 limited_consumer.assert_next_some("goop/pee/test");
6498 limited_consumer.assert_next_wait(); }
6500
6501 #[tokio::test]
6502 async fn test_with_root_unauthorized() {
6503 let origin = Origin::random().produce();
6504
6505 let limited_producer = origin
6507 .scope(&["allowed".into()])
6508 .expect("should create limited producer");
6509
6510 assert!(limited_producer.with_root("notallowed").is_none());
6512
6513 let allowed_root = limited_producer
6515 .with_root("allowed")
6516 .expect("should create allowed root");
6517 assert_eq!(allowed_root.root().as_str(), "allowed");
6518 }
6519
6520 #[tokio::test]
6521 async fn test_wildcard_permission() {
6522 let origin = Origin::random().produce();
6523
6524 let root_producer = origin.clone();
6526
6527 let _broadcast = root_producer
6529 .create_broadcast("any/path", announce())
6530 .expect("publish allowed");
6531 let _keep2 = root_producer
6532 .create_broadcast("other/path", announce())
6533 .expect("publish allowed");
6534 settle().await;
6535
6536 let foo_producer = root_producer.with_root("foo").expect("should create any root");
6538 assert_eq!(foo_producer.root().as_str(), "foo");
6539 }
6540
6541 #[tokio::test]
6542 async fn test_consume_broadcast_with_permissions() {
6543 let origin = Origin::random().produce();
6544
6545 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
6546 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
6547 settle().await;
6548
6549 let limited_consumer = origin
6551 .consume()
6552 .scope(&["allowed".into()])
6553 .expect("should create limited consumer");
6554
6555 let result = limited_consumer.get_broadcast("allowed/test");
6557 assert!(result.is_some());
6558 assert!(
6559 result
6560 .unwrap()
6561 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
6562 );
6563
6564 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
6566
6567 let consumer = origin.consume();
6569 assert!(consumer.get_broadcast("allowed/test").is_some());
6570 assert!(consumer.get_broadcast("notallowed/test").is_some());
6571 }
6572
6573 #[tokio::test]
6574 async fn test_nested_paths_with_permissions() {
6575 let origin = Origin::random().produce();
6576
6577 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
6579
6580 let _broadcast = limited_producer
6582 .create_broadcast("a/b/c", announce())
6583 .expect("publish allowed");
6584 let _keep2 = limited_producer
6585 .create_broadcast("a/b/c/d", announce())
6586 .expect("publish allowed");
6587 let _keep3 = limited_producer
6588 .create_broadcast("a/b/c/d/e", announce())
6589 .expect("publish allowed");
6590 settle().await;
6591
6592 assert!(limited_producer.create_broadcast("a", announce()).is_err());
6594 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
6595 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
6596 }
6597
6598 #[tokio::test]
6599 async fn test_multiple_consumers_with_different_permissions() {
6600 let origin = Origin::random().produce();
6601
6602 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
6604 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
6605 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
6606 settle().await;
6607
6608 let mut foo_consumer = origin
6610 .consume()
6611 .scope(&["foo".into()])
6612 .expect("should create foo consumer")
6613 .announced();
6614
6615 let mut bar_consumer = origin
6616 .consume()
6617 .scope(&["bar".into()])
6618 .expect("should create bar consumer")
6619 .announced();
6620
6621 let mut foobar_consumer = origin
6622 .consume()
6623 .scope(&["foo".into(), "bar".into()])
6624 .expect("should create foobar consumer")
6625 .announced();
6626
6627 foo_consumer.assert_next_some("foo/test");
6629 foo_consumer.assert_next_wait();
6630
6631 bar_consumer.assert_next_some("bar/test");
6632 bar_consumer.assert_next_wait();
6633
6634 foobar_consumer.assert_next_some("bar/test");
6635 foobar_consumer.assert_next_some("foo/test");
6636 foobar_consumer.assert_next_wait();
6637 }
6638
6639 #[tokio::test]
6640 async fn test_select_with_empty_prefix() {
6641 let origin = Origin::random().produce();
6642
6643 let demo_producer = origin.with_root("demo").expect("should create demo root");
6645 let limited_producer = demo_producer
6646 .scope(&["worm-node".into(), "foobar".into()])
6647 .expect("should create limited producer");
6648
6649 let _broadcast1 = limited_producer
6651 .create_broadcast("worm-node/test", announce())
6652 .expect("publish allowed");
6653 let _broadcast2 = limited_producer
6654 .create_broadcast("foobar/test", announce())
6655 .expect("publish allowed");
6656 settle().await;
6657
6658 let mut consumer = limited_producer
6660 .consume()
6661 .scope(&["".into()])
6662 .expect("should create consumer with empty prefix")
6663 .announced();
6664
6665 let a1 = consumer.try_next().expect("expected first announcement");
6667 let a2 = consumer.try_next().expect("expected second announcement");
6668 consumer.assert_next_wait();
6669
6670 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6671 paths.sort();
6672 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
6673 }
6674
6675 #[tokio::test]
6676 async fn test_select_narrowing_scope() {
6677 let origin = Origin::random().produce();
6678
6679 let demo_producer = origin.with_root("demo").expect("should create demo root");
6681 let limited_producer = demo_producer
6682 .scope(&["worm-node".into(), "foobar".into()])
6683 .expect("should create limited producer");
6684
6685 let _broadcast1 = limited_producer
6687 .create_broadcast("worm-node", announce())
6688 .expect("publish allowed");
6689 let _broadcast2 = limited_producer
6690 .create_broadcast("worm-node/foo", announce())
6691 .expect("publish allowed");
6692 let _broadcast3 = limited_producer
6693 .create_broadcast("foobar/bar", announce())
6694 .expect("publish allowed");
6695 settle().await;
6696
6697 let mut worm_consumer = limited_producer
6699 .consume()
6700 .scope(&["worm-node".into()])
6701 .expect("should create worm-node consumer")
6702 .announced();
6703
6704 worm_consumer.assert_next_some("worm-node");
6706 worm_consumer.assert_next_some("worm-node/foo");
6707 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
6711 .consume()
6712 .scope(&["worm-node/foo".into()])
6713 .expect("should create worm-node/foo consumer")
6714 .announced();
6715
6716 foo_consumer.assert_next_some("worm-node/foo");
6717 foo_consumer.assert_next_wait(); }
6719
6720 #[tokio::test]
6721 async fn test_select_multiple_roots_with_empty_prefix() {
6722 let origin = Origin::random().produce();
6723
6724 let limited_producer = origin
6726 .scope(&["app1".into(), "app2".into(), "shared".into()])
6727 .expect("should create limited producer");
6728
6729 let _broadcast1 = limited_producer
6731 .create_broadcast("app1/data", announce())
6732 .expect("publish allowed");
6733 let _broadcast2 = limited_producer
6734 .create_broadcast("app2/config", announce())
6735 .expect("publish allowed");
6736 let _broadcast3 = limited_producer
6737 .create_broadcast("shared/resource", announce())
6738 .expect("publish allowed");
6739 settle().await;
6740
6741 let mut consumer = limited_producer
6743 .consume()
6744 .scope(&["".into()])
6745 .expect("should create consumer with empty prefix")
6746 .announced();
6747
6748 consumer.assert_next_some("app1/data");
6750 consumer.assert_next_some("app2/config");
6751 consumer.assert_next_some("shared/resource");
6752 consumer.assert_next_wait();
6753 }
6754
6755 #[tokio::test]
6756 async fn test_publish_scope_with_empty_prefix() {
6757 let origin = Origin::random().produce();
6758
6759 let limited_producer = origin
6761 .scope(&["services/api".into(), "services/web".into()])
6762 .expect("should create limited producer");
6763
6764 let same_producer = limited_producer
6766 .scope(&["".into()])
6767 .expect("should create producer with empty prefix");
6768
6769 let _broadcast = same_producer
6771 .create_broadcast("services/api", announce())
6772 .expect("publish allowed");
6773 let _keep2 = same_producer
6774 .create_broadcast("services/web", announce())
6775 .expect("publish allowed");
6776 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
6777 assert!(same_producer.create_broadcast("other", announce()).is_err());
6778 }
6779
6780 #[tokio::test]
6781 async fn test_select_narrowing_to_deeper_path() {
6782 let origin = Origin::random().produce();
6783
6784 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
6786
6787 let _broadcast1 = limited_producer
6789 .create_broadcast("org/team1/project1", announce())
6790 .expect("publish allowed");
6791 let _broadcast2 = limited_producer
6792 .create_broadcast("org/team1/project2", announce())
6793 .expect("publish allowed");
6794 let _broadcast3 = limited_producer
6795 .create_broadcast("org/team2/project1", announce())
6796 .expect("publish allowed");
6797 settle().await;
6798
6799 let mut team2_consumer = limited_producer
6801 .consume()
6802 .scope(&["org/team2".into()])
6803 .expect("should create team2 consumer")
6804 .announced();
6805
6806 team2_consumer.assert_next_some("org/team2/project1");
6807 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
6811 .consume()
6812 .scope(&["org/team1/project1".into()])
6813 .expect("should create project1 consumer")
6814 .announced();
6815
6816 project1_consumer.assert_next_some("org/team1/project1");
6818 project1_consumer.assert_next_wait();
6819 }
6820
6821 #[tokio::test]
6822 async fn test_select_with_non_matching_prefix() {
6823 let origin = Origin::random().produce();
6824
6825 let limited_producer = origin
6827 .scope(&["allowed/path".into()])
6828 .expect("should create limited producer");
6829
6830 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
6832
6833 assert!(limited_producer.scope(&["other/path".into()]).is_none());
6835 }
6836
6837 #[tokio::test]
6840 async fn test_with_root_trailing_slash_consumer() {
6841 let origin = Origin::random().produce();
6842
6843 let prefix = "some_prefix/".to_string();
6845 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6846
6847 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6848 settle().await;
6849 consumer.assert_next_some("test");
6850 }
6851
6852 #[tokio::test]
6854 async fn test_with_root_trailing_slash_producer() {
6855 let origin = Origin::random().produce();
6856
6857 let prefix = "some_prefix/".to_string();
6859 let rooted = origin.with_root(prefix).unwrap();
6860
6861 let _b = rooted.create_broadcast("test", announce()).unwrap();
6862 settle().await;
6863
6864 let mut consumer = rooted.consume().announced();
6865 consumer.assert_next_some("test");
6866 }
6867
6868 #[tokio::test]
6870 async fn test_with_root_trailing_slash_unannounce() {
6871 tokio::time::pause();
6872
6873 let origin = Origin::random().produce();
6874
6875 let prefix = "some_prefix/".to_string();
6876 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
6877
6878 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
6879 settle().await;
6880 consumer.assert_next_some("test");
6881
6882 b.finish();
6884 settle().await;
6885
6886 consumer.assert_next_none("test");
6888 }
6889
6890 #[tokio::test]
6891 async fn test_select_maintains_access_with_wider_prefix() {
6892 let origin = Origin::random().produce();
6893
6894 let demo_producer = origin.with_root("demo").expect("should create demo root");
6896 let user_producer = demo_producer
6897 .scope(&["worm-node".into(), "foobar".into()])
6898 .expect("should create user producer");
6899
6900 let _broadcast1 = user_producer
6902 .create_broadcast("worm-node/data", announce())
6903 .expect("publish allowed");
6904 let _broadcast2 = user_producer
6905 .create_broadcast("foobar", announce())
6906 .expect("publish allowed");
6907 settle().await;
6908
6909 let mut consumer = user_producer
6911 .consume()
6912 .scope(&["".into()])
6913 .expect("scope with empty prefix should not fail when user has specific permissions")
6914 .announced();
6915
6916 let a1 = consumer.try_next().expect("expected first announcement");
6918 let a2 = consumer.try_next().expect("expected second announcement");
6919 consumer.assert_next_wait();
6920
6921 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6922 paths.sort();
6923 assert_eq!(paths, ["foobar", "worm-node/data"]);
6924
6925 let mut narrow_consumer = user_producer
6927 .consume()
6928 .scope(&["worm-node".into()])
6929 .expect("should be able to narrow scope to worm-node")
6930 .announced();
6931
6932 narrow_consumer.assert_next_some("worm-node/data");
6933 narrow_consumer.assert_next_wait(); }
6935
6936 #[tokio::test]
6937 async fn test_duplicate_prefixes_deduped() {
6938 let origin = Origin::random().produce();
6939
6940 let producer = origin
6942 .scope(&["demo".into(), "demo".into()])
6943 .expect("should create producer");
6944
6945 let _broadcast = producer
6946 .create_broadcast("demo/stream", announce())
6947 .expect("publish allowed");
6948 settle().await;
6949
6950 let mut consumer = producer.consume().announced();
6951 consumer.assert_next_some("demo/stream");
6952 consumer.assert_next_wait();
6953 }
6954
6955 #[tokio::test]
6956 async fn test_overlapping_prefixes_deduped() {
6957 let origin = Origin::random().produce();
6958
6959 let producer = origin
6961 .scope(&["demo".into(), "demo/foo".into()])
6962 .expect("should create producer");
6963
6964 let _broadcast = producer
6966 .create_broadcast("demo/bar/stream", announce())
6967 .expect("publish allowed");
6968 settle().await;
6969
6970 let mut consumer = producer.consume().announced();
6971 consumer.assert_next_some("demo/bar/stream");
6972 consumer.assert_next_wait();
6973 }
6974
6975 #[tokio::test]
6976 async fn test_overlapping_prefixes_no_duplicate_announcements() {
6977 let origin = Origin::random().produce();
6978
6979 let producer = origin
6981 .scope(&["demo".into(), "demo/foo".into()])
6982 .expect("should create producer");
6983
6984 let _broadcast = producer
6985 .create_broadcast("demo/foo/stream", announce())
6986 .expect("publish allowed");
6987 settle().await;
6988
6989 let mut consumer = producer.consume().announced();
6990 consumer.assert_next_some("demo/foo/stream");
6992 consumer.assert_next_wait();
6993 }
6994
6995 #[tokio::test]
6996 async fn test_allowed_returns_deduped_prefixes() {
6997 let origin = Origin::random().produce();
6998
6999 let producer = origin
7000 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
7001 .expect("should create producer");
7002
7003 let allowed: Vec<_> = producer.allowed().collect();
7004 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
7005 }
7006
7007 #[tokio::test]
7008 async fn test_announced_broadcast_already_announced() {
7009 let origin = Origin::random().produce();
7010
7011 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
7012 settle().await;
7013
7014 let consumer = origin.consume();
7015 let result = consumer.announced_broadcast("test").await.expect("should find it");
7016 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
7017 }
7018
7019 #[tokio::test]
7020 async fn test_announced_broadcast_delayed() {
7021 tokio::time::pause();
7022
7023 let origin = Origin::random().produce();
7024
7025 let consumer = origin.consume();
7026
7027 let wait = tokio::spawn({
7029 let consumer = consumer.clone();
7030 async move { consumer.announced_broadcast("test").await }
7031 });
7032
7033 tokio::task::yield_now().await;
7035
7036 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
7037 settle().await;
7038
7039 let result = wait.await.unwrap().expect("should find it");
7040 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
7041 }
7042
7043 #[tokio::test]
7044 async fn test_announced_broadcast_ignores_unrelated_paths() {
7045 tokio::time::pause();
7046
7047 let origin = Origin::random().produce();
7048
7049 let consumer = origin.consume();
7050
7051 let wait = tokio::spawn({
7052 let consumer = consumer.clone();
7053 async move { consumer.announced_broadcast("target").await }
7054 });
7055
7056 tokio::task::yield_now().await;
7057
7058 let _other = origin.create_broadcast("other", announce()).unwrap();
7060 settle().await;
7061 tokio::task::yield_now().await;
7062 assert!(!wait.is_finished(), "must not resolve on unrelated path");
7063
7064 let _target = origin.create_broadcast("target", announce()).unwrap();
7065 settle().await;
7066 let result = wait.await.unwrap().expect("should find target");
7067 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
7068 }
7069
7070 #[tokio::test]
7071 async fn test_announced_broadcast_skips_nested_paths() {
7072 tokio::time::pause();
7073
7074 let origin = Origin::random().produce();
7075
7076 let consumer = origin.consume();
7077
7078 let wait = tokio::spawn({
7079 let consumer = consumer.clone();
7080 async move { consumer.announced_broadcast("foo").await }
7081 });
7082
7083 tokio::task::yield_now().await;
7084
7085 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
7087 settle().await;
7088 tokio::task::yield_now().await;
7089 assert!(!wait.is_finished(), "must not resolve on a nested path");
7090
7091 let _exact = origin.create_broadcast("foo", announce()).unwrap();
7092 settle().await;
7093 let result = wait.await.unwrap().expect("should find foo exactly");
7094 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
7095 }
7096
7097 #[tokio::test]
7098 async fn test_announced_broadcast_disallowed() {
7099 let origin = Origin::random().produce();
7100 let limited = origin
7101 .consume()
7102 .scope(&["allowed".into()])
7103 .expect("should create limited");
7104
7105 assert!(limited.announced_broadcast("notallowed").await.is_none());
7107 }
7108
7109 #[tokio::test]
7110 async fn test_announced_broadcast_scope_too_narrow() {
7111 let origin = Origin::random().produce();
7114 let limited = origin
7115 .consume()
7116 .scope(&["foo/specific".into()])
7117 .expect("should create limited");
7118
7119 let result = limited
7121 .announced_broadcast("foo")
7122 .now_or_never()
7123 .expect("must not block");
7124 assert!(result.is_none());
7125 }
7126
7127 #[tokio::test]
7131 async fn test_coalesce_announce_then_unannounce() {
7132 tokio::time::pause();
7134
7135 let origin = Origin::random().produce();
7136 let mut announced = origin.consume().announced();
7137
7138 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
7139 settle().await;
7140 broadcast.finish();
7141
7142 settle().await;
7143
7144 announced.assert_next_wait();
7145 }
7146
7147 #[tokio::test]
7148 async fn test_coalesce_announce_unannounce_announce() {
7149 tokio::time::pause();
7152
7153 let origin = Origin::random().produce();
7154 let mut announced = origin.consume().announced();
7155
7156 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
7157 settle().await;
7158 broadcast1.finish();
7159 settle().await;
7160 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
7161 settle().await;
7162
7163 announced.assert_next_some("test");
7164 announced.assert_next_wait();
7165 }
7166
7167 #[tokio::test]
7168 async fn test_coalesce_unannounce_announce_preserved() {
7169 tokio::time::pause();
7172
7173 let origin = Origin::random().produce();
7174 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
7175 settle().await;
7176
7177 let mut announced = origin.consume().announced();
7178 announced.assert_next_some("test");
7179
7180 broadcast1.finish();
7182 settle().await;
7183
7184 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
7185 settle().await;
7186
7187 announced.assert_next_none("test");
7189 announced.assert_next_some("test");
7190 announced.assert_next_wait();
7191 }
7192
7193 #[tokio::test]
7194 async fn test_coalesce_unannounce_announce_unannounce() {
7195 tokio::time::pause();
7198
7199 let origin = Origin::random().produce();
7200 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
7201 settle().await;
7202
7203 let mut announced = origin.consume().announced();
7204 announced.assert_next_some("test");
7205
7206 broadcast1.finish();
7207 settle().await;
7208
7209 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
7210 settle().await;
7211 broadcast2.finish();
7212 settle().await;
7213
7214 announced.assert_next_none("test");
7215 announced.assert_next_wait();
7216 }
7217
7218 #[tokio::test]
7219 async fn test_coalesce_churn_bounded() {
7220 tokio::time::pause();
7225
7226 let origin = Origin::random().produce();
7227 let mut announced = origin.consume().announced();
7228
7229 for _ in 0..1000 {
7230 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
7231 settle().await;
7232 broadcast.finish();
7233 }
7234 settle().await;
7235
7236 let mut collected = Vec::new();
7237 while let Some(update) = announced.try_next() {
7238 collected.push(update);
7239 }
7240 assert!(
7241 collected.len() <= 1,
7242 "expected at most one pending update, got {}",
7243 collected.len()
7244 );
7245 assert!(
7246 collected.iter().all(|a| a.path == Path::new("test")),
7247 "unexpected path in pending updates",
7248 );
7249 }
7250
7251 #[tokio::test]
7255 async fn test_consumer_clone_is_side_effect_free() {
7256 let origin = Origin::random().produce();
7257
7258 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
7259 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
7260 settle().await;
7261
7262 let consumer = origin.consume();
7263 let mut announced = consumer.announced();
7264
7265 for _ in 0..16 {
7268 let cloned = consumer.clone();
7269 assert!(cloned.get_broadcast("test1").is_some());
7270 assert!(cloned.get_broadcast("test2").is_some());
7271 }
7272
7273 let a1 = announced.try_next().expect("first announcement");
7276 let a2 = announced.try_next().expect("second announcement");
7277 announced.assert_next_wait();
7278
7279 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
7280 paths.sort();
7281 assert_eq!(paths, ["test1", "test2"]);
7282
7283 let mut fresh = consumer.announced();
7285 let b1 = fresh.try_next().expect("backlog: first");
7286 let b2 = fresh.try_next().expect("backlog: second");
7287 fresh.assert_next_wait();
7288
7289 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
7290 paths.sort();
7291 assert_eq!(paths, ["test1", "test2"]);
7292 }
7293
7294 #[tokio::test]
7296 async fn dynamic_request_unroutable_without_handler() {
7297 let origin = Origin::random().produce();
7298 let consumer = origin.consume();
7299 assert!(matches!(
7300 consumer.request_broadcast("missing").await,
7301 Err(Error::Unroutable)
7302 ));
7303 }
7304
7305 #[tokio::test(start_paused = true)]
7308 async fn dynamic_request_served_not_announced() {
7309 let origin = Origin::random().produce();
7310 let mut dynamic = origin.dynamic();
7311 let consumer = origin.consume();
7312
7313 let mut announced = origin.consume().announced();
7315 announced.assert_next_wait();
7316
7317 let served = broadcast::Info::new().produce();
7318 let request_fut = consumer.request_broadcast("fallback");
7321
7322 let mut served_dynamic = served.dynamic();
7324
7325 let request = dynamic.requested_broadcast().await.unwrap();
7326 assert_eq!(request.path(), &Path::new("fallback"));
7327 request.accept(&served);
7328
7329 let broadcast = request_fut.await.unwrap();
7330 assert!(broadcast.is_clone(&served.consume()));
7331
7332 let track_fut = broadcast.track("video").unwrap().subscribe(None);
7334 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
7335 let mut track = track_fut.await.unwrap();
7336 producer.append_group().unwrap();
7337 track.assert_group();
7338
7339 announced.assert_next_wait();
7341 }
7342
7343 #[tokio::test(start_paused = true)]
7345 async fn dynamic_request_coalesces() {
7346 let origin = Origin::random().produce();
7347 let mut dynamic = origin.dynamic();
7348 let consumer = origin.consume();
7349
7350 let f1 = consumer.request_broadcast("dup");
7352 let f2 = consumer.request_broadcast("dup");
7353
7354 let request = dynamic.requested_broadcast().await.unwrap();
7356 assert_eq!(request.path(), &Path::new("dup"));
7357 assert!(
7358 dynamic.requested_broadcast().now_or_never().is_none(),
7359 "a coalesced request must not be served twice"
7360 );
7361
7362 let served = broadcast::Info::new().produce();
7364 request.accept(&served);
7365 assert!(f1.await.unwrap().is_clone(&served.consume()));
7366 assert!(f2.await.unwrap().is_clone(&served.consume()));
7367 }
7368
7369 #[tokio::test(start_paused = true)]
7372 async fn dynamic_request_dedups_served() {
7373 let origin = Origin::random().produce();
7374 let mut dynamic = origin.dynamic();
7375 let consumer = origin.consume();
7376
7377 let request_fut = consumer.request_broadcast("fallback");
7378 let request = dynamic.requested_broadcast().await.unwrap();
7379 let served = broadcast::Info::new().produce();
7380 request.accept(&served);
7381 let first = request_fut.await.unwrap();
7382 assert!(first.is_clone(&served.consume()));
7383
7384 let second = consumer.request_broadcast("fallback").await.unwrap();
7386 assert!(second.is_clone(&served.consume()));
7387
7388 assert!(
7390 dynamic.requested_broadcast().now_or_never().is_none(),
7391 "a still-live served broadcast must not be re-requested from the handler"
7392 );
7393 }
7394
7395 #[tokio::test(start_paused = true)]
7397 async fn dynamic_request_reserves_after_close() {
7398 let origin = Origin::random().produce();
7399 let mut dynamic = origin.dynamic();
7400 let consumer = origin.consume();
7401
7402 let request_fut = consumer.request_broadcast("fallback");
7403 let request = dynamic.requested_broadcast().await.unwrap();
7404 let served = broadcast::Info::new().produce();
7405 request.accept(&served);
7406 request_fut.await.unwrap();
7407
7408 drop(served);
7410
7411 let request_fut = consumer.request_broadcast("fallback");
7413 let request = dynamic.requested_broadcast().await.unwrap();
7414 assert_eq!(request.path(), &Path::new("fallback"));
7415 let served = broadcast::Info::new().produce();
7416 request.accept(&served);
7417 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
7418 }
7419
7420 #[tokio::test(start_paused = true)]
7423 async fn dynamic_request_served_cache_bounded() {
7424 let origin = Origin::random().produce();
7425 let mut dynamic = origin.dynamic();
7426 let consumer = origin.consume();
7427
7428 for i in 0..100 {
7429 let path = format!("one-shot/{i}");
7430 let request_fut = consumer.request_broadcast(&path);
7431 let request = dynamic.requested_broadcast().await.unwrap();
7432 let served = broadcast::Info::new().produce();
7433 request.accept(&served);
7434 request_fut.await.unwrap();
7435 drop(served);
7437 }
7438
7439 assert!(
7442 origin.dynamic.read().served.len() <= 4,
7443 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
7444 origin.dynamic.read().served.len()
7445 );
7446 }
7447
7448 #[tokio::test(start_paused = true)]
7451 async fn dynamic_request_coalesces_after_handoff() {
7452 let origin = Origin::random().produce();
7453 let mut dynamic = origin.dynamic();
7454 let consumer = origin.consume();
7455
7456 let f1 = consumer.request_broadcast("fallback");
7457 let request = dynamic.requested_broadcast().await.unwrap();
7459
7460 let f2 = consumer.request_broadcast("fallback");
7462 assert!(
7463 dynamic.requested_broadcast().now_or_never().is_none(),
7464 "a repeat request during hand-off must coalesce, not re-queue"
7465 );
7466
7467 let served = broadcast::Info::new().produce();
7469 request.accept(&served);
7470 assert!(f1.await.unwrap().is_clone(&served.consume()));
7471 assert!(f2.await.unwrap().is_clone(&served.consume()));
7472 }
7473
7474 #[tokio::test(start_paused = true)]
7476 async fn dynamic_request_dropped_after_handoff() {
7477 let origin = Origin::random().produce();
7478 let mut dynamic = origin.dynamic();
7479 let consumer = origin.consume();
7480
7481 let f1 = consumer.request_broadcast("fallback");
7482 let request = dynamic.requested_broadcast().await.unwrap();
7483 let f2 = consumer.request_broadcast("fallback");
7484
7485 drop(request);
7487 assert!(matches!(f1.await, Err(Error::Unroutable)));
7488 assert!(matches!(f2.await, Err(Error::Unroutable)));
7489 }
7490
7491 #[tokio::test(start_paused = true)]
7493 async fn dynamic_request_rejected() {
7494 let origin = Origin::random().produce();
7495 let mut dynamic = origin.dynamic();
7496 let consumer = origin.consume();
7497
7498 let request_fut = consumer.request_broadcast("fallback");
7499
7500 let request = dynamic.requested_broadcast().await.unwrap();
7501 request.reject(Error::Cancel);
7502
7503 assert!(matches!(request_fut.await, Err(Error::Cancel)));
7504 }
7505
7506 #[tokio::test(start_paused = true)]
7510 async fn dynamic_request_rerequest_after_reject() {
7511 let origin = Origin::random().produce();
7512 let mut dynamic = origin.dynamic();
7513 let consumer = origin.consume();
7514
7515 let f1 = consumer.request_broadcast("fallback");
7516 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
7517 assert!(matches!(f1.await, Err(Error::Unroutable)));
7518
7519 let served = broadcast::Info::new().produce();
7520 let f2 = consumer.request_broadcast("fallback");
7522 let request = dynamic.requested_broadcast().await.unwrap();
7523 assert_eq!(request.path(), &Path::new("fallback"));
7524 request.accept(&served);
7525 assert!(f2.await.unwrap().is_clone(&served.consume()));
7526 }
7527
7528 #[tokio::test(start_paused = true)]
7531 async fn dynamic_request_handler_dropped() {
7532 let origin = Origin::random().produce();
7533 let dynamic = origin.dynamic();
7534 let consumer = origin.consume();
7535
7536 let request_fut = consumer.request_broadcast("fallback");
7537 drop(dynamic);
7538 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
7539
7540 assert!(matches!(
7542 consumer.request_broadcast("again").await,
7543 Err(Error::Unroutable)
7544 ));
7545 }
7546
7547 #[tokio::test(start_paused = true)]
7551 async fn dynamic_request_accept_after_handler_dropped() {
7552 let origin = Origin::random().produce();
7553 let mut dynamic = origin.dynamic();
7554 let consumer = origin.consume();
7555
7556 let request_fut = consumer.request_broadcast("fallback");
7557
7558 let request = dynamic.requested_broadcast().await.unwrap();
7560 drop(dynamic);
7561
7562 let served = broadcast::Info::new().produce();
7563 request.accept(&served);
7565 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
7566 }
7567
7568 #[tokio::test(start_paused = true)]
7570 async fn dynamic_request_prefers_announced() {
7571 let origin = Origin::random().produce();
7572 let mut dynamic = origin.dynamic();
7573 let consumer = origin.consume();
7574
7575 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
7576 settle().await;
7577
7578 let got = consumer.request_broadcast("live").await.unwrap();
7579 assert!(
7580 got.is_clone(&consumer.get_broadcast("live").unwrap()),
7581 "should return the published broadcast"
7582 );
7583 assert!(
7584 dynamic.requested_broadcast().now_or_never().is_none(),
7585 "a published path must not queue a fallback request"
7586 );
7587 }
7588
7589 #[tokio::test(start_paused = true)]
7591 async fn dynamic_clone_keeps_alive() {
7592 let origin = Origin::random().produce();
7593 let dynamic = origin.dynamic();
7594 let consumer = origin.consume();
7595
7596 drop(dynamic.clone());
7597
7598 let request_fut = consumer.request_broadcast("fallback");
7601 assert!(
7602 request_fut.now_or_never().is_none(),
7603 "request should stay pending until served"
7604 );
7605 }
7606
7607 fn wedge_watchdog<F>(name: &str, secs: u64, scenario: F)
7613 where
7614 F: std::future::Future<Output = ()> + Send + 'static,
7615 {
7616 let (done_tx, done_rx) = std::sync::mpsc::channel::<()>();
7617 let handle = std::thread::spawn(move || {
7618 let rt = ::tokio::runtime::Builder::new_current_thread()
7619 .enable_time()
7620 .build()
7621 .unwrap();
7622 rt.block_on(scenario);
7623 let _ = done_tx.send(());
7624 });
7625 match done_rx.recv_timeout(std::time::Duration::from_secs(secs)) {
7626 Ok(()) => {
7627 let _ = handle.join();
7628 }
7629 Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
7630 let err = handle.join().unwrap_err();
7632 std::panic::resume_unwind(err);
7633 }
7634 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
7635 panic!("{name}: scenario wedged; a task is spinning inside a single poll")
7636 }
7637 }
7638 }
7639
7640 #[test]
7653 fn test_active_corpse_does_not_livelock_takeover() {
7654 wedge_watchdog("active-corpse", 20, async {
7655 let origin = Origin::random().produce();
7656 let consumer = origin.consume();
7657 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
7658
7659 let source_a = origin
7662 .create_broadcast("test", announce().with_hops(hops.clone()))
7663 .unwrap();
7664 let mut dynamic_a = source_a.dynamic();
7665 settle().await;
7666 settle().await;
7667
7668 let broadcast = consumer.request_broadcast("test").await.unwrap();
7669 let subscribing = broadcast.track("video").unwrap().subscribe(None);
7670 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
7671 settle().await;
7672 let mut sub = subscribing.await.unwrap();
7673 producer_a.append_group().unwrap();
7674 sub.assert_group();
7675
7676 let source_b = origin
7680 .create_broadcast("test", announce().with_hops(hops.clone()))
7681 .unwrap();
7682 let dynamic_b = source_b.dynamic();
7683 settle().await;
7684
7685 drop(dynamic_b);
7693 source_b.abort(Error::Dropped).unwrap();
7694 settle().await;
7695
7696 producer_a.append_group().unwrap();
7699 sub.assert_group();
7700 sub.assert_not_closed();
7701 });
7702 }
7703
7704 #[test]
7712 fn test_route_churn_never_wedges() {
7713 for seed in 1..=8u64 {
7714 wedge_watchdog(&format!("churn seed {seed}"), 30, churn_scenario(seed));
7715 }
7716 }
7717
7718 async fn churn_scenario(seed: u64) {
7719 let mut rng = seed.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
7721 let mut next = move || {
7722 rng = rng.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
7723 rng >> 33
7724 };
7725
7726 let origin = Origin::random().produce();
7727 let consumer = origin.consume();
7728 let names: Vec<Arc<str>> = (0..8).map(|i| Arc::from(format!("t{i}"))).collect();
7729
7730 let mut subs: Vec<track::Subscriber> = Vec::new();
7731 let mut pending_subs: Vec<kio::Pending<track::Subscribing>> = Vec::new();
7732
7733 struct Source {
7734 producer: Option<broadcast::Producer>,
7735 server: ::tokio::task::JoinHandle<()>,
7736 }
7737 let mut sources: Vec<Source> = Vec::new();
7738
7739 for step in 0..400u64 {
7740 match next() % 10 {
7741 0 | 1 => {
7744 if sources.len() >= 3 {
7745 continue;
7746 }
7747 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
7748 let route = announce().with_hops(hops).with_cost(next() % 4);
7749 let Ok(source) = origin.create_broadcast("test", route) else {
7750 continue;
7751 };
7752 let mut dynamic = source.dynamic();
7753 let behavior = next();
7754 let server = ::tokio::spawn(async move {
7755 let mut round = 0u64;
7756 let mut kept: Vec<track::Producer> = Vec::new();
7757 while let Ok(request) = dynamic.requested_track().await {
7758 round += 1;
7759 match (behavior >> (round % 16)) % 4 {
7760 0 => drop(request), 1 => {
7762 let mut producer = request.accept(None);
7763 let _ = producer.create_group(group::Info { sequence: round });
7764 let _ = producer.abort(Error::Dropped);
7765 }
7766 2 => {
7767 let mut producer = request.accept(None);
7768 let _ = producer.create_group(group::Info { sequence: round });
7769 kept.push(producer);
7770 }
7771 _ => {
7772 let mut producer = request.accept(None);
7773 let _ = producer.finish();
7774 }
7775 }
7776 }
7777 });
7778 sources.push(Source {
7779 producer: Some(source),
7780 server,
7781 });
7782 }
7783 2 | 3 => {
7785 if sources.is_empty() {
7786 continue;
7787 }
7788 let i = (next() as usize) % sources.len();
7789 let mut source = sources.swap_remove(i);
7790 if next() % 2 == 0
7791 && let Some(producer) = source.producer.take()
7792 {
7793 let _ = producer.abort(Error::Dropped);
7794 }
7795 source.server.abort();
7796 }
7797 4..=6 => {
7799 if subs.len() + pending_subs.len() >= 24 {
7800 continue;
7801 }
7802 let Some(broadcast) = consumer.get_broadcast("test") else {
7803 continue;
7804 };
7805 let name = &names[(next() as usize) % names.len()];
7806 if let Ok(track) = broadcast.track(name.as_ref()) {
7807 pending_subs.push(track.subscribe(None));
7808 }
7809 }
7810 7 => {
7812 if subs.is_empty() {
7813 continue;
7814 }
7815 let i = (next() as usize) % subs.len();
7816 subs.swap_remove(i);
7817 }
7818 _ => {
7820 for sub in pending_subs.drain(..) {
7821 match ::tokio::time::timeout(std::time::Duration::from_millis(5), sub).await {
7822 Ok(Ok(sub)) => subs.push(sub),
7823 Ok(Err(_)) => {}
7824 Err(_) => {}
7826 }
7827 }
7828 for sub in subs.iter_mut() {
7829 while let Some(Ok(Some(_))) = sub.recv_group().now_or_never() {}
7830 }
7831 }
7832 }
7833 if step % 16 == 0 {
7834 settle().await;
7835 }
7836 for _ in 0..(next() % 3) {
7838 ::tokio::task::yield_now().await;
7839 }
7840 }
7841
7842 for source in sources {
7843 source.server.abort();
7844 }
7845 settle().await;
7846 }
7847}