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 linger: Duration,
129}
130
131impl Default for Info {
132 fn default() -> Self {
135 Self {
136 id: Origin::UNKNOWN,
137 pool: cache::Pool::default(),
138 cache_duration: Duration::MAX,
139 linger: Duration::ZERO,
140 }
141 }
142}
143
144impl Info {
145 pub fn new(id: Origin) -> Self {
147 Self { id, ..Self::default() }
148 }
149
150 pub fn with_pool(mut self, pool: cache::Pool) -> Self {
152 self.pool = pool;
153 self
154 }
155
156 pub fn with_cache_duration(mut self, cache_duration: Duration) -> Self {
159 self.cache_duration = cache_duration;
160 self
161 }
162
163 pub fn with_linger(mut self, linger: Duration) -> Self {
166 self.linger = linger;
167 self
168 }
169
170 pub fn produce(self) -> Producer {
172 Producer::new(self)
173 }
174}
175
176impl TryFrom<u64> for Origin {
177 type Error = InvalidOrigin;
178
179 fn try_from(id: u64) -> Result<Self, Self::Error> {
180 Self::new(id)
181 }
182}
183
184impl fmt::Display for Origin {
185 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
186 self.id.fmt(f)
187 }
188}
189
190impl<V: Copy> Encode<V> for Origin
191where
192 u64: Encode<V>,
193{
194 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
195 self.id.encode(w, version)
196 }
197}
198
199impl<V: Copy> Decode<V> for Origin
200where
201 u64: Decode<V>,
202{
203 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
204 let id = u64::decode(r, version)?;
205 if id >= 1u64 << 62 {
206 return Err(DecodeError::InvalidValue);
207 }
208 Ok(Self { id })
209 }
210}
211
212pub(crate) const MAX_HOPS: usize = 32;
218
219#[derive(Debug, Clone, Default, PartialEq, Eq)]
224pub struct OriginList(Vec<Origin>);
225
226#[derive(Debug, Clone, Copy, PartialEq, Eq)]
228#[non_exhaustive]
229pub struct TooManyOrigins;
230
231impl fmt::Display for TooManyOrigins {
232 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
233 write!(f, "too many origins (max {MAX_HOPS})")
234 }
235}
236
237impl std::error::Error for TooManyOrigins {}
238
239impl From<TooManyOrigins> for DecodeError {
240 fn from(_: TooManyOrigins) -> Self {
241 DecodeError::BoundsExceeded
242 }
243}
244
245impl OriginList {
246 pub fn new() -> Self {
248 Self(Vec::new())
249 }
250
251 pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
253 if self.0.len() >= MAX_HOPS {
254 return Err(TooManyOrigins);
255 }
256 self.0.push(origin);
257 Ok(())
258 }
259
260 pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
263 for entry in &mut self.0 {
264 if *entry == target {
265 *entry = replacement;
266 return true;
267 }
268 }
269 false
270 }
271
272 pub fn contains(&self, origin: &Origin) -> bool {
274 self.0.contains(origin)
275 }
276
277 pub fn len(&self) -> usize {
279 self.0.len()
280 }
281
282 pub fn is_empty(&self) -> bool {
284 self.0.is_empty()
285 }
286
287 pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
289 self.0.iter()
290 }
291
292 pub fn as_slice(&self) -> &[Origin] {
294 &self.0
295 }
296}
297
298impl TryFrom<Vec<Origin>> for OriginList {
299 type Error = TooManyOrigins;
300
301 fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
302 if v.len() > MAX_HOPS {
303 return Err(TooManyOrigins);
304 }
305 Ok(Self(v))
306 }
307}
308
309impl<'a> IntoIterator for &'a OriginList {
310 type Item = &'a Origin;
311 type IntoIter = std::slice::Iter<'a, Origin>;
312
313 fn into_iter(self) -> Self::IntoIter {
314 self.iter()
315 }
316}
317
318impl<V: Copy> Encode<V> for OriginList
319where
320 u64: Encode<V>,
321 Origin: Encode<V>,
322{
323 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
324 (self.0.len() as u64).encode(w, version)?;
325 for origin in &self.0 {
326 origin.encode(w, version)?;
327 }
328 Ok(())
329 }
330}
331
332impl<V: Copy> Decode<V> for OriginList
333where
334 u64: Decode<V>,
335 Origin: Decode<V>,
336{
337 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
338 let count = u64::decode(r, version)? as usize;
339 if count > MAX_HOPS {
340 return Err(DecodeError::BoundsExceeded);
341 }
342 let mut list = Vec::with_capacity(count);
343 for _ in 0..count {
344 list.push(Origin::decode(r, version)?);
345 }
346 Ok(Self(list))
347 }
348}
349
350static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
351
352#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
353struct ConsumerId(u64);
354
355impl ConsumerId {
356 fn new() -> Self {
357 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
358 }
359}
360
361struct OriginBroadcast {
367 path: PathOwned,
368 broadcast: broadcast::Producer,
370 state: kio::Producer<FrontState>,
373 announced: bool,
374}
375
376fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
385 (hops.len(), fnv_key(name, hops.iter().copied()))
386}
387
388fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
401 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
403
404 let mut hash = SEED;
405 for &byte in name.as_str().as_bytes() {
406 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
407 }
408 for origin in origins {
409 for &byte in &origin.id().to_le_bytes() {
410 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
411 }
412 }
413
414 hash
415}
416
417fn route_order(name: &Path, route: &FrontRoute) -> (bool, u64, usize, u64, Reverse<u64>) {
436 let (len, hash) = route_key(name, &route.route.hops);
437 (!route.route.announce, route.route.cost, len, hash, Reverse(route.id))
438}
439
440enum PendingUpdate {
449 Announce(broadcast::Consumer),
450 Unannounce,
451 UnannounceAnnounce(broadcast::Consumer),
452}
453
454#[derive(Default)]
459struct OriginConsumerState {
460 pending: BTreeMap<PathOwned, PendingUpdate>,
461}
462
463impl OriginConsumerState {
464 fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
465 let new = match self.pending.remove(&path) {
466 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
468 Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
470 PendingUpdate::UnannounceAnnounce(broadcast)
471 }
472 };
473 self.pending.insert(path, new);
474 }
475
476 fn apply_unannounce(&mut self, path: PathOwned) {
477 match self.pending.remove(&path) {
478 Some(PendingUpdate::Announce(_)) => {}
480 None | Some(PendingUpdate::Unannounce) => {
481 self.pending.insert(path, PendingUpdate::Unannounce);
482 }
483 Some(PendingUpdate::UnannounceAnnounce(_)) => {
486 self.pending.insert(path, PendingUpdate::Unannounce);
487 }
488 }
489 }
490
491 fn take(&mut self) -> Option<OriginAnnounce> {
493 let path = self.pending.keys().next()?.clone();
494 let broadcast = match self.pending.remove(&path).unwrap() {
495 PendingUpdate::Announce(broadcast) => Some(broadcast),
496 PendingUpdate::Unannounce => None,
497 PendingUpdate::UnannounceAnnounce(broadcast) => {
498 self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
501 None
502 }
503 };
504 Some(OriginAnnounce { path, broadcast })
505 }
506}
507
508#[derive(Clone)]
509struct AnnounceConsumerNotify {
510 root: PathOwned,
511 state: kio::Producer<OriginConsumerState>,
512}
513
514impl AnnounceConsumerNotify {
515 fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer) {
516 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
517 self.state
518 .write()
519 .ok()
520 .expect("consumer closed")
521 .apply_announce(path, broadcast);
522 }
523
524 fn unannounce(&self, path: impl AsPath) {
525 let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
526 self.state.write().ok().expect("consumer closed").apply_unannounce(path);
527 }
528}
529
530struct NotifyNode {
531 parent: Option<Lock<NotifyNode>>,
532
533 consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
536}
537
538impl NotifyNode {
539 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
540 Self {
541 parent,
542 consumers: HashMap::new(),
543 }
544 }
545
546 fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer) {
547 for consumer in self.consumers.values() {
548 consumer.announce(path.as_path(), broadcast.clone());
549 }
550
551 if let Some(parent) = &self.parent {
552 parent.lock().announce(path, broadcast);
553 }
554 }
555
556 fn unannounce(&mut self, path: impl AsPath) {
557 for consumer in self.consumers.values() {
558 consumer.unannounce(path.as_path());
559 }
560
561 if let Some(parent) = &self.parent {
562 parent.lock().unannounce(path);
563 }
564 }
565}
566
567pub(crate) struct ExclusionGuard {
574 state: kio::Producer<FrontState>,
575 peer: Origin,
576}
577
578impl ExclusionGuard {
579 fn new(state: &kio::Producer<FrontState>, peer: Origin) -> Option<Arc<Self>> {
582 let mut s = state.write().ok()?;
583 if s.closed {
584 return None;
585 }
586 *s.excluded.entry(peer).or_default() += 1;
587 drop(s);
588 Some(Arc::new(Self {
589 state: state.clone(),
590 peer,
591 }))
592 }
593}
594
595impl Drop for ExclusionGuard {
596 fn drop(&mut self) {
597 let Ok(mut state) = self.state.write() else { return };
598 if let std::collections::hash_map::Entry::Occupied(mut entry) = state.excluded.entry(self.peer) {
599 match entry.get() {
600 1 => drop(entry.remove()),
601 n => *entry.get_mut() = n - 1,
602 }
603 }
604 }
605}
606
607enum Resolved {
614 Found(broadcast::Consumer),
617 Excluded,
619 Missing,
621}
622
623struct OriginNode {
624 broadcast: Option<OriginBroadcast>,
627
628 nested: HashMap<String, Lock<OriginNode>>,
630
631 notify: Lock<NotifyNode>,
633}
634
635impl OriginNode {
636 fn new(parent: Option<Lock<NotifyNode>>) -> Self {
637 Self {
638 broadcast: None,
639 nested: HashMap::new(),
640 notify: Lock::new(NotifyNode::new(parent)),
641 }
642 }
643
644 fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
645 let (dir, rest) = path.next_part().expect("leaf called with empty path");
646
647 let next = self.entry(dir);
648 if rest.is_empty() { next } else { next.lock().leaf(&rest) }
649 }
650
651 fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
652 match self.nested.get(dir) {
653 Some(next) => next.clone(),
654 None => {
655 let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
656 self.nested.insert(dir.to_string(), next.clone());
657 next
658 }
659 }
660 }
661
662 fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
666 let Some(existing) = &mut self.broadcast else { return };
667 if !existing.state.same_channel(expect) || existing.announced == announce {
668 return;
669 }
670 existing.announced = announce;
671 let path = existing.path.clone();
672 let consumer = existing.broadcast.consume();
673 let mut notify = self.notify.lock();
674 if announce {
675 notify.announce(&path, &consumer);
676 } else {
677 notify.unannounce(&path);
678 }
679 }
680
681 fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
682 self.consume_initial(&mut notify);
683 self.notify.lock().consumers.insert(id, notify);
684 }
685
686 fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
687 if let Some(broadcast) = &self.broadcast
690 && broadcast.announced
691 {
692 notify.announce(&broadcast.path, broadcast.broadcast.consume());
693 }
694
695 for nested in self.nested.values() {
697 nested.lock().consume_initial(notify);
698 }
699 }
700
701 fn resolve_broadcast(&self, rest: impl AsPath, exclude: Option<Origin>) -> Resolved {
702 let rest = rest.as_path();
703
704 if let Some((dir, rest)) = rest.next_part() {
705 let Some(node) = self.nested.get(dir) else {
706 return Resolved::Missing;
707 };
708 let node = node.lock();
709 return node.resolve_broadcast(&rest, exclude);
710 }
711
712 let Some(broadcast) = self.broadcast.as_ref() else {
713 return Resolved::Missing;
714 };
715 let Some(origin) = exclude else {
716 return Resolved::Found(broadcast.broadcast.consume());
717 };
718
719 let state = broadcast.state.read();
735 if !state.routes.iter().any(|r| r.route.hops.contains(&origin)) {
736 drop(state);
737 let shared = broadcast.broadcast.consume();
738 return match ExclusionGuard::new(&broadcast.state, origin) {
739 Some(guard) => Resolved::Found(shared.with_exclusion(guard)),
740 None => Resolved::Found(shared),
743 };
744 }
745 match state
746 .dispatch(Some(origin))
747 .and_then(|clean| state.routes.iter().find(|r| r.id == clean))
748 {
749 Some(route) => Resolved::Found(route.source.clone()),
750 None => Resolved::Excluded,
751 }
752 }
753
754 fn unconsume(&mut self, id: ConsumerId) {
755 self.notify.lock().consumers.remove(&id).expect("consumer not found");
756 if self.is_empty() {
757 }
760 }
761
762 fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
766 let relative = relative.as_path();
767
768 if let Some((dir, relative)) = relative.next_part() {
769 let Some(nested) = self.nested.get(dir) else { return };
770 let nested = nested.clone();
771 let mut locked = nested.lock();
772 locked.remove(expect, &relative);
773
774 if locked.is_empty() {
775 drop(locked);
776 self.nested.remove(dir);
777 }
778 } else if let Some(existing) = &self.broadcast
779 && existing.state.same_channel(expect)
780 {
781 let existing = self.broadcast.take().expect("checked above");
782 if existing.announced {
783 self.notify.lock().unannounce(&existing.path);
784 }
785 }
786 }
787
788 fn is_empty(&self) -> bool {
789 self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
790 }
791}
792
793#[derive(Clone)]
794struct OriginNodes {
795 nodes: Vec<(PathOwned, Lock<OriginNode>)>,
796}
797
798impl OriginNodes {
799 pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
802 let mut roots = Vec::new();
803
804 for (root, state) in &self.nodes {
805 for prefix in prefixes {
806 if root.has_prefix(prefix) {
807 roots.push((root.to_owned(), state.clone()));
809 continue;
810 }
811
812 if let Some(suffix) = prefix.strip_prefix(root) {
813 let nested = state.lock().leaf(&suffix);
815 roots.push((prefix.to_owned(), nested));
816 }
817 }
818 }
819
820 if roots.is_empty() {
821 None
822 } else {
823 Some(Self { nodes: roots })
824 }
825 }
826
827 pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
828 let new_root = new_root.as_path();
829 let mut roots = Vec::new();
830
831 if new_root.is_empty() {
832 return Some(self.clone());
833 }
834
835 for (root, state) in &self.nodes {
836 if let Some(suffix) = root.strip_prefix(&new_root) {
837 roots.push((suffix.to_owned(), state.clone()));
839 } else if let Some(suffix) = new_root.strip_prefix(root) {
840 let nested = state.lock().leaf(&suffix);
843 roots.push(("".into(), nested));
844 }
845 }
846
847 if roots.is_empty() {
848 None
849 } else {
850 Some(Self { nodes: roots })
851 }
852 }
853
854 pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
856 let path = path.as_path();
857
858 for (root, state) in &self.nodes {
859 if let Some(suffix) = path.strip_prefix(root) {
860 return Some((state.clone(), suffix.to_owned()));
861 }
862 }
863
864 None
865 }
866}
867
868impl Default for OriginNodes {
869 fn default() -> Self {
870 Self {
871 nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
872 }
873 }
874}
875
876#[derive(Clone)]
878pub struct OriginAnnounce {
879 pub path: PathOwned,
881 pub broadcast: Option<broadcast::Consumer>,
887}
888
889#[derive(Clone)]
891pub struct Producer {
892 info: Origin,
896
897 nodes: OriginNodes,
900
901 root: PathOwned,
903
904 dynamic: kio::Shared<OriginDynamicState>,
908
909 pool: cache::Pool,
912
913 cache_duration: Duration,
916
917 linger: Duration,
920
921 stats: stats::Session,
925}
926
927impl std::ops::Deref for Producer {
928 type Target = Origin;
929
930 fn deref(&self) -> &Self::Target {
931 &self.info
932 }
933}
934
935impl Producer {
936 pub fn new(info: Info) -> Self {
940 Self {
941 info: info.id,
942 nodes: OriginNodes::default(),
943 root: PathOwned::default(),
944 dynamic: kio::Shared::default(),
945 pool: info.pool,
946 cache_duration: info.cache_duration,
947 linger: info.linger,
948 stats: stats::Session::default(),
949 }
950 }
951
952 pub fn with_stats(mut self, session: stats::Session) -> Self {
956 self.stats = session;
957 self
958 }
959
960 pub fn with_linger(mut self, linger: Duration) -> Self {
969 self.linger = linger;
970 self
971 }
972
973 pub fn info(&self) -> Info {
976 Info {
977 id: self.info,
978 pool: self.pool.clone(),
979 cache_duration: self.cache_duration,
980 linger: self.linger,
981 }
982 }
983
984 pub(crate) fn empty(info: Origin) -> Self {
989 Self {
990 info,
991 nodes: OriginNodes { nodes: Vec::new() },
992 root: PathOwned::default(),
993 dynamic: kio::Shared::default(),
994 pool: cache::Pool::default(),
995 cache_duration: Duration::MAX,
996 linger: Duration::ZERO,
997 stats: stats::Session::default(),
998 }
999 }
1000
1001 pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
1052 let path = path.as_path();
1053
1054 debug_assert!(
1055 !route.hops.contains(&self.info),
1056 "create_broadcast called with a looping hop chain",
1057 );
1058
1059 let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
1060 let full = self.root.join(&path).to_owned();
1061
1062 if full.parts().count() > Path::MAX_PARTS {
1066 return Err(BoundsExceeded.into());
1067 }
1068
1069 let ingress = self.stats.ingress(&full);
1073
1074 let mut source = broadcast::Info { origin: self.info() }
1075 .produce()
1076 .with_stats(ingress.clone());
1077 source.set_route(route).expect("fresh producer");
1078
1079 web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
1080
1081 Ok(source)
1082 }
1083
1084 pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
1090 let prefixes = PathPrefixes::new(prefixes);
1091 Some(Producer {
1092 info: self.info,
1093 nodes: self.nodes.select(&prefixes)?,
1094 root: self.root.clone(),
1095 dynamic: self.dynamic.clone(),
1096 pool: self.pool.clone(),
1097 cache_duration: self.cache_duration,
1098 linger: self.linger,
1099 stats: self.stats.clone(),
1100 })
1101 }
1102
1103 pub fn dynamic(&self) -> Dynamic {
1112 Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
1113 }
1114
1115 pub fn consume(&self) -> Consumer {
1120 Consumer::new(
1123 self.info,
1124 self.root.clone(),
1125 self.nodes.clone(),
1126 self.dynamic.clone(),
1127 stats::Session::default(),
1128 )
1129 }
1130
1131 pub fn announces(&self) -> AnnounceProducer {
1137 AnnounceProducer::new(self.root.clone(), self.nodes.clone())
1138 }
1139
1140 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
1145 let prefix = prefix.as_path();
1146
1147 Some(Self {
1148 info: self.info,
1149 root: self.root.join(&prefix).to_owned(),
1150 nodes: self.nodes.root(&prefix)?,
1151 dynamic: self.dynamic.clone(),
1152 pool: self.pool.clone(),
1153 cache_duration: self.cache_duration,
1154 linger: self.linger,
1155 stats: self.stats.clone(),
1156 })
1157 }
1158
1159 pub fn root(&self) -> &Path<'_> {
1161 &self.root
1162 }
1163
1164 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
1167 self.nodes.nodes.iter().map(|(root, _)| root)
1168 }
1169
1170 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1172 self.root.join(path)
1173 }
1174}
1175
1176const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1189
1190struct FrontRoute {
1192 id: u64,
1193 route: broadcast::Route,
1196 source: broadcast::Consumer,
1198}
1199
1200struct FrontState {
1202 path: PathOwned,
1204 self_origin: Origin,
1206 publisher: Option<Origin>,
1214 next_route: u64,
1217 routes: Vec<FrontRoute>,
1218 excluded: HashMap<Origin, usize>,
1225 active: Option<u64>,
1227 linger: Duration,
1230 closed: bool,
1235}
1236
1237impl FrontState {
1238 fn pick(&self, keep: impl Fn(&FrontRoute) -> bool, untainted: bool) -> Option<u64> {
1245 let candidates: Vec<&FrontRoute> = self.routes.iter().filter(|r| keep(r)).collect();
1246 let candidates = match untainted {
1247 true => self.prefer_untainted(&candidates),
1248 false => candidates,
1249 };
1250 candidates
1251 .into_iter()
1252 .min_by_key(|r| route_order(&self.path.as_path(), r))
1253 .map(|r| r.id)
1254 }
1255
1256 fn best_route(&self) -> Option<u64> {
1261 self.pick(|_| true, true)
1262 }
1263
1264 fn dispatch(&self, exclude: Option<Origin>) -> Option<u64> {
1271 self.pick(|r| exclude.is_none_or(|origin| !r.route.hops.contains(&origin)), false)
1272 }
1273
1274 fn taints_a_reader(&self, route: &broadcast::Route) -> bool {
1276 route.hops.iter().any(|hop| self.excluded.contains_key(hop))
1277 }
1278
1279 fn prefer_untainted<'a>(&self, candidates: &[&'a FrontRoute]) -> Vec<&'a FrontRoute> {
1289 if self.excluded.is_empty() {
1290 return candidates.to_vec();
1291 }
1292 let clean: Vec<&FrontRoute> = candidates
1293 .iter()
1294 .copied()
1295 .filter(|r| !self.taints_a_reader(&r.route))
1296 .collect();
1297 match clean.is_empty() {
1298 true => candidates.to_vec(),
1299 false => clean,
1300 }
1301 }
1302
1303 fn serve_route(&self, skip: impl Fn(u64) -> bool) -> Option<u64> {
1312 if let Some(active) = self.active
1313 && !skip(active)
1314 && let Some(route) = self.routes.iter().find(|r| r.id == active)
1315 && !self.taints_a_reader(&route.route)
1316 {
1317 return Some(active);
1318 }
1319 self.pick(|r| !skip(r.id), true)
1320 }
1321
1322 fn reselect(&mut self, carrying: bool) {
1340 let best = self.best_route();
1341 if carrying
1342 && let (Some(best_id), Some(cur_id)) = (best, self.active)
1343 && best_id != cur_id
1344 && let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
1345 && let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
1346 && incumbent.route.announce
1347 && candidate.route.cost < incumbent.route.cost
1348 && candidate.route.advertised == 0
1349 && candidate.route.hops.len() >= 2
1350 && !self.handover_allowed(&candidate.route)
1351 {
1352 return;
1354 }
1355 self.active = best;
1356 }
1357
1358 fn handover_allowed(&self, route: &broadcast::Route) -> bool {
1370 let name = self.path.as_path();
1371 match route.hops.iter().last() {
1372 Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
1373 None => true,
1374 }
1375 }
1376
1377 fn routes_snapshot(&self) -> Vec<broadcast::Route> {
1383 let mut routes: Vec<&FrontRoute> = self.routes.iter().collect();
1384 routes.sort_by_key(|r| route_order(&self.path.as_path(), r));
1385 routes.sort_by_key(|r| Some(r.id) != self.active);
1386 routes.into_iter().map(|r| r.route.clone()).collect()
1387 }
1388}
1389
1390fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
1400 let mut leaf_guard = leaf.lock();
1405 let routes = state.read().routes_snapshot();
1406 if let Some(advert) = routes.first() {
1407 let announce = advert.announce;
1408 broadcast.clone().set_routes(routes);
1409 leaf_guard.set_announced(state, announce);
1410 }
1411}
1412
1413fn detach_source(
1426 state: &kio::Producer<FrontState>,
1427 broadcast: &broadcast::Producer,
1428 leaf: &Lock<OriginNode>,
1429 id: u64,
1430 graceful: bool,
1431) {
1432 let close = {
1433 let carrying = broadcast.demand().is_used();
1437 let Ok(mut s) = state.write() else { return };
1438 let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
1439 return;
1440 };
1441 s.routes.remove(pos);
1442 s.reselect(carrying);
1443 if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
1444 s.closed = true;
1447 true
1448 } else {
1449 false
1450 }
1451 };
1452 if close {
1453 broadcast.abort_spliced(Error::Dropped);
1454 }
1455 sync_front(state, broadcast, leaf);
1456}
1457
1458fn sync_announce(guard: &mut Option<stats::Announce>, announced: bool, ingress: &stats::Scope) {
1461 match (announced, guard.is_some()) {
1462 (true, false) => *guard = Some(ingress.announce()),
1463 (false, true) => *guard = None,
1464 _ => {}
1465 }
1466}
1467
1468async fn run_source(
1478 origin: Info,
1479 node: Lock<OriginNode>,
1480 full: PathOwned,
1481 rest: PathOwned,
1482 mut source: broadcast::Consumer,
1483 ingress: stats::Scope,
1484) {
1485 let ctx = AttachContext {
1486 origin: &origin,
1487 node: &node,
1488 full: &full,
1489 rest: &rest,
1490 };
1491
1492 let Ok(mut route) = source.route_changed().await else {
1496 return;
1498 };
1499
1500 let mut announce = route.announce.then(|| ingress.announce());
1505
1506 'attach: loop {
1507 let leaf = if rest.is_empty() {
1512 node.clone()
1513 } else {
1514 node.lock().leaf(&rest)
1515 };
1516
1517 let (state, broadcast, id) = match attach_source(&ctx, &leaf, &source, route.clone()) {
1518 Attach::Ready(state, broadcast, id) => (state, broadcast, id),
1519 Attach::Parked(incumbent) => {
1520 tracing::warn!(
1521 broadcast = %full,
1522 "path already live with a different publisher; parking an offline source until it ends",
1523 );
1524 let update = kio::wait(|waiter| {
1528 if let Poll::Ready(update) = source.poll_route_changed(waiter) {
1529 return Poll::Ready(Some(update));
1530 }
1531 match incumbent.poll(waiter, |s| if s.closed { Poll::Ready(()) } else { Poll::Pending }) {
1534 Poll::Ready(_) => Poll::Ready(None),
1535 Poll::Pending => Poll::Pending,
1536 }
1537 })
1538 .await;
1539 match update {
1540 Some(Ok(update)) => {
1542 sync_announce(&mut announce, update.announce, &ingress);
1543 route = update;
1544 }
1545 Some(Err(_)) => return,
1547 None => {}
1549 }
1550 continue 'attach;
1551 }
1552 };
1553 let publisher = route.hops.iter().next().copied();
1554
1555 loop {
1556 match source.route_changed().await {
1557 Ok(update) => {
1558 let announced = update.announce;
1559 if update.hops.iter().next().copied() != publisher {
1570 detach_source(&state, &broadcast, &leaf, id, true);
1571 sync_announce(&mut announce, announced, &ingress);
1572 route = update;
1573 continue 'attach;
1574 }
1575 {
1576 let carrying = broadcast.demand().is_used();
1577 let Ok(mut s) = state.write() else { return };
1578 let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
1579 return;
1580 };
1581 if entry.route == update {
1582 continue;
1583 }
1584 entry.route = update;
1585 s.reselect(carrying);
1586 }
1587 sync_announce(&mut announce, announced, &ingress);
1589 sync_front(&state, &broadcast, &leaf);
1590 }
1591 Err(_) => {
1592 detach_source(&state, &broadcast, &leaf, id, source.is_finished());
1595 return;
1596 }
1597 }
1598 }
1599 }
1600}
1601
1602enum Attach {
1604 Ready(kio::Producer<FrontState>, broadcast::Producer, u64),
1607 Parked(kio::Producer<FrontState>),
1612}
1613
1614struct AttachContext<'a> {
1616 origin: &'a Info,
1617 node: &'a Lock<OriginNode>,
1618 full: &'a PathOwned,
1620 rest: &'a PathOwned,
1622}
1623
1624fn same_publisher(a: Option<Origin>, b: Option<Origin>) -> bool {
1649 if a == Some(Origin::UNKNOWN) || b == Some(Origin::UNKNOWN) {
1650 return false;
1651 }
1652 a == b
1653}
1654
1655fn attach_source(
1656 ctx: &AttachContext,
1657 leaf: &Lock<OriginNode>,
1658 source: &broadcast::Consumer,
1659 route: broadcast::Route,
1660) -> Attach {
1661 let publisher = route.hops.iter().next().copied();
1662 let mut leaf_guard = leaf.lock();
1663
1664 if let Some(existing) = &leaf_guard.broadcast {
1667 let mut joined = None;
1668 let carrying = existing.broadcast.demand().is_used();
1669 if let Ok(mut s) = existing.state.write()
1670 && !s.closed
1671 {
1672 if same_publisher(s.publisher, publisher) {
1673 let id = s.next_route;
1674 s.next_route += 1;
1675 s.routes.push(FrontRoute {
1676 id,
1677 route: route.clone(),
1678 source: source.clone(),
1679 });
1680 s.reselect(carrying);
1681 joined = Some(id);
1682 } else if !route.announce {
1683 return Attach::Parked(existing.state.clone());
1684 } else {
1685 s.closed = true;
1693 tracing::warn!(broadcast = %ctx.full, "replacing a live broadcast from a different publisher");
1694 }
1695 }
1696 if let Some(id) = joined {
1697 let state = existing.state.clone();
1698 let broadcast = existing.broadcast.clone();
1699 drop(leaf_guard);
1700 sync_front(&state, &broadcast, leaf);
1701 return Attach::Ready(state, broadcast, id);
1702 }
1703 }
1704
1705 let announce = route.announce;
1707 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
1708 origin: ctx.origin.clone(),
1709 });
1710 let _ = broadcast.clone().set_route(route.clone());
1711 let state = kio::Producer::new(FrontState {
1712 path: ctx.full.clone(),
1713 self_origin: ctx.origin.id,
1714 publisher,
1715 next_route: 1,
1716 excluded: HashMap::new(),
1717 routes: vec![FrontRoute {
1718 id: 0,
1719 route,
1720 source: source.clone(),
1721 }],
1722 active: Some(0),
1723 linger: ctx.origin.linger,
1724 closed: false,
1725 });
1726
1727 if let Some(stale) = leaf_guard.broadcast.take()
1731 && stale.announced
1732 {
1733 leaf_guard.notify.lock().unannounce(&stale.path);
1734 }
1735 let entry = OriginBroadcast {
1736 path: ctx.full.clone(),
1737 broadcast: broadcast.clone(),
1738 state: state.clone(),
1739 announced: announce,
1740 };
1741 if entry.announced {
1742 leaf_guard.notify.lock().announce(ctx.full, &broadcast.consume());
1743 }
1744 leaf_guard.broadcast = Some(entry);
1745 drop(leaf_guard);
1746
1747 web_async::spawn(run_front(
1748 state.clone(),
1749 broadcast.clone(),
1750 ctx.node.clone(),
1751 ctx.rest.clone(),
1752 ));
1753
1754 Attach::Ready(state, broadcast, 0)
1755}
1756
1757async fn run_front(
1760 state: kio::Producer<FrontState>,
1761 mut broadcast: broadcast::Producer,
1762 node: Lock<OriginNode>,
1763 rest: PathOwned,
1764) {
1765 enum Step {
1766 Serve(Arc<str>, super::resume::Producer),
1767 Changed,
1769 Expired,
1771 Closed,
1772 }
1773
1774 let linger = state.read().linger;
1775 let mut deadline = kio::time::Deadline::new();
1780
1781 loop {
1782 let empty = {
1783 let s = state.read();
1784 !s.closed && s.routes.is_empty()
1785 };
1786 deadline.set(match (empty, deadline.deadline()) {
1787 (true, None) => web_async::time::Instant::now().checked_add(linger),
1790 (true, at) => at,
1791 (false, _) => None,
1792 });
1793
1794 let step = {
1795 kio::wait(|waiter| {
1796 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
1797 return Poll::Ready(Step::Serve(name, resume));
1798 }
1799 match state.poll(waiter, |s| {
1802 if s.closed || s.routes.is_empty() != empty {
1803 Poll::Ready(())
1804 } else {
1805 Poll::Pending
1806 }
1807 }) {
1808 Poll::Ready(Ok(guard)) => {
1809 return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
1810 }
1811 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1812 Poll::Pending => {}
1813 }
1814 deadline.poll(waiter).map(|_| Step::Expired)
1815 })
1816 .await
1817 };
1818
1819 match step {
1820 Step::Serve(name, resume) => {
1821 web_async::spawn(serve_track(state.clone(), name, resume));
1824 }
1825 Step::Changed => {}
1826 Step::Expired => {
1827 let close = {
1831 let Ok(mut s) = state.write() else { break };
1832 if !s.closed && s.routes.is_empty() {
1833 s.closed = true;
1834 true
1835 } else {
1836 false
1837 }
1838 };
1839 if close {
1840 break;
1841 }
1842 }
1843 Step::Closed => break,
1844 }
1845 }
1846
1847 broadcast.abort_spliced(Error::Dropped);
1849
1850 broadcast.finish();
1852
1853 node.lock().remove(&state, &rest);
1856}
1857
1858async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
1871 enum Step {
1872 Closed,
1873 Splice(u64, broadcast::Consumer),
1874 Complete,
1875 Failed(Error),
1876 NoRoute,
1881 Idle,
1883 Demand,
1885 }
1886
1887 let mut serving: Option<(u64, track::Consumer)> = None;
1889 let mut spliced_edge: Option<u64> = None;
1895 let mut refused: HashSet<u64> = HashSet::new();
1900 let mut refusal: Option<Error> = None;
1901 let mut dead: HashSet<u64> = HashSet::new();
1906 let mut idle_since: Option<web_async::time::Instant> = None;
1908 let mut deadline = kio::time::Deadline::new();
1909
1910 loop {
1911 let serving_id = serving.as_ref().map(|(id, _)| *id);
1912
1913 {
1920 let s = state.read();
1921 refused.retain(|id| s.routes.iter().any(|r| r.id == *id));
1922 dead.retain(|id| s.routes.iter().any(|r| r.id == *id));
1923 let exhausted = !s.routes.is_empty()
1924 && s.serve_route(|id| refused.contains(&id) || dead.contains(&id))
1925 .is_none();
1926 if exhausted && dead.is_empty() {
1927 drop(s);
1928 let err = refusal.take().unwrap_or(Error::NotFound);
1929 tracing::debug!(name = %name, %err, "every source refused track; aborting");
1930 let _ = resume.abort(err);
1931 return;
1932 }
1933 }
1934
1935 let used = resume.is_used();
1947 idle_since = match (resume.is_spliced(), used) {
1948 (true, false) => idle_since.or_else(|| Some(web_async::time::Instant::now())),
1949 _ => None,
1950 };
1951 deadline.set(idle_since.and_then(|at| at.checked_add(TRACK_IDLE_LINGER)));
1952
1953 let step = {
1954 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
1955 kio::wait(|waiter| {
1956 match state.poll(waiter, |s| {
1962 let gone = serving_id.is_some_and(|id| !s.routes.iter().any(|r| r.id == id));
1963 if s.closed
1964 || (used && (gone || matches!(s.serve_route(skip), Some(next) if Some(next) != serving_id)))
1965 {
1966 Poll::Ready(())
1967 } else {
1968 Poll::Pending
1969 }
1970 }) {
1971 Poll::Ready(Ok(guard)) => {
1972 if guard.closed {
1973 return Poll::Ready(Step::Closed);
1974 }
1975 let Some(next) = guard.serve_route(skip) else {
1976 return Poll::Ready(Step::NoRoute);
1977 };
1978 let source = guard
1979 .routes
1980 .iter()
1981 .find(|r| r.id == next)
1982 .expect("servable source in table")
1983 .source
1984 .clone();
1985 return Poll::Ready(Step::Splice(next, source));
1986 }
1987 Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
1988 Poll::Pending => {}
1989 }
1990
1991 let edge = match used {
1996 true => resume.poll_unused(waiter),
1997 false => resume.poll_used(waiter),
1998 };
1999 if edge.is_ready() {
2000 return Poll::Ready(Step::Demand);
2001 }
2002
2003 if let Some((_, track)) = &serving
2006 && let Poll::Ready(result) = track.poll_complete(waiter)
2007 {
2008 return Poll::Ready(match result {
2009 Ok(()) => Step::Complete,
2010 Err(err) => Step::Failed(err),
2011 });
2012 }
2013
2014 deadline.poll(waiter).map(|_| Step::Idle)
2015 })
2016 .await
2017 };
2018
2019 match step {
2020 Step::Closed => return,
2022 Step::Complete => {
2023 let _ = resume.finish();
2024 return;
2025 }
2026 Step::Failed(err) => {
2027 if resume.latest() == spliced_edge
2035 && let Some(id) = serving_id
2036 {
2037 let closing = state
2038 .read()
2039 .routes
2040 .iter()
2041 .find(|r| r.id == id)
2042 .is_some_and(|r| r.source.is_closing());
2043 if closing {
2044 dead.insert(id);
2045 } else {
2046 refused.insert(id);
2047 refusal = Some(err);
2048 }
2049 }
2050 serving = None;
2051 }
2052 Step::Demand => {}
2054 Step::NoRoute => serving = None,
2060 Step::Idle => {
2061 if resume.release().is_err() {
2066 return;
2068 }
2069 serving = None;
2070 }
2071 Step::Splice(id, source) => {
2072 let attempt = match source.track(&name) {
2076 Ok(track) => {
2077 let query = track.info().into_inner();
2080 let skip = |id: u64| refused.contains(&id) || dead.contains(&id);
2081 let info = kio::wait(|waiter| {
2082 if let Poll::Ready(result) = query.poll(waiter) {
2083 return Poll::Ready(Some(result));
2084 }
2085 match state.poll(waiter, |s| {
2086 if s.closed || s.serve_route(skip) != Some(id) {
2087 Poll::Ready(())
2088 } else {
2089 Poll::Pending
2090 }
2091 }) {
2092 Poll::Ready(_) => Poll::Ready(None),
2093 Poll::Pending => Poll::Pending,
2094 }
2095 })
2096 .await;
2097 match info {
2098 None => continue,
2100 Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
2103 Poll::Ready(Err(err)) => Err(err),
2104 _ => Ok(track),
2105 },
2106 Some(Err(err)) => Err(err),
2107 }
2108 }
2109 Err(err) => Err(err),
2110 };
2111
2112 match attempt {
2113 Ok(track) => {
2114 if let Err(err) = resume.takeover(&track) {
2115 let _ = resume.abort(err);
2120 return;
2121 }
2122 dead.clear();
2123 spliced_edge = resume.latest();
2126 serving = Some((id, track));
2127 }
2128 Err(_) if source.is_closing() => {
2132 dead.insert(id);
2133 serving = None;
2134 }
2135 Err(err) => {
2141 tracing::debug!(name = %name, source = id, %err, "source refused track");
2142 refused.insert(id);
2143 refusal = Some(err);
2144 serving = None;
2145 }
2146 }
2147 }
2148 }
2149 }
2150}
2151
2152#[derive(Default)]
2158struct OriginDynamicState {
2159 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
2162
2163 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
2169}
2170
2171#[derive(Default)]
2178struct PendingBroadcast {
2179 resolved: Option<Result<broadcast::Consumer, Error>>,
2180}
2181
2182pub struct Dynamic {
2193 info: Origin,
2194 root: PathOwned,
2195 state: kio::Shared<OriginDynamicState>,
2196}
2197
2198impl Clone for Dynamic {
2199 fn clone(&self) -> Self {
2200 self.state.lock().requests.add_handler();
2204
2205 Self {
2206 info: self.info,
2207 root: self.root.clone(),
2208 state: self.state.clone(),
2209 }
2210 }
2211}
2212
2213impl Dynamic {
2214 fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
2215 state.lock().requests.add_handler();
2216
2217 Self { info, root, state }
2218 }
2219
2220 pub fn info(&self) -> &Origin {
2222 &self.info
2223 }
2224
2225 pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2227 let mut state = ready!(self.state.poll(waiter, |state| {
2228 if state.requests.has_queued() {
2229 Poll::Ready(())
2230 } else {
2231 Poll::Pending
2232 }
2233 }));
2234
2235 let path = state.requests.pop().expect("predicate guaranteed a request");
2236 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2242 Poll::Ready(Ok(Request {
2243 path,
2244 producer,
2245 state: self.state.clone(),
2246 }))
2247 }
2248
2249 pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
2252 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2253 }
2254
2255 pub fn root(&self) -> &Path<'_> {
2257 &self.root
2258 }
2259}
2260
2261impl Drop for Dynamic {
2262 fn drop(&mut self) {
2263 let mut state = self.state.lock();
2266 if state.requests.remove_handler() {
2267 state.requests.drain_queued();
2271 }
2272 }
2273}
2274
2275pub struct Request {
2282 path: PathOwned,
2284
2285 producer: kio::Producer<PendingBroadcast>,
2288
2289 state: kio::Shared<OriginDynamicState>,
2291}
2292
2293impl Request {
2294 pub fn path(&self) -> &Path<'_> {
2296 &self.path
2297 }
2298
2299 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2305 let broadcast = broadcast.consume();
2306
2307 let resolved = {
2313 let mut state = self.state.lock();
2314 let existing = state.served.insert(self.path.clone(), broadcast.weak());
2315 state
2316 .requests
2317 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2318 existing.map(|weak| weak.consume()).unwrap_or(broadcast)
2319 };
2320
2321 if let Ok(mut pending) = self.producer.write() {
2322 pending.resolved = Some(Ok(resolved));
2323 }
2324 }
2326
2327 pub fn reject(self, err: Error) {
2329 self.state
2330 .lock()
2331 .requests
2332 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2333 if let Ok(mut state) = self.producer.write() {
2334 state.resolved = Some(Err(err));
2335 }
2336 }
2337}
2338
2339impl Drop for Request {
2340 fn drop(&mut self) {
2341 self.state
2349 .lock()
2350 .requests
2351 .remove_if(&self.path, |producer| producer.same_channel(&self.producer));
2352 }
2353}
2354
2355pub struct Requesting {
2362 inner: RequestState,
2363 stats: stats::Scope,
2366}
2367
2368enum RequestState {
2369 Ready(broadcast::Consumer),
2371 Failed(Error),
2374 Pending(kio::Consumer<PendingBroadcast>),
2376}
2377
2378impl Requesting {
2379 fn ready(broadcast: broadcast::Consumer) -> Self {
2380 Self {
2381 inner: RequestState::Ready(broadcast),
2382 stats: stats::Scope::default(),
2383 }
2384 }
2385
2386 fn failed(error: Error) -> Self {
2387 Self {
2388 inner: RequestState::Failed(error),
2389 stats: stats::Scope::default(),
2390 }
2391 }
2392
2393 fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2394 Self {
2395 inner: RequestState::Pending(consumer),
2396 stats: stats::Scope::default(),
2397 }
2398 }
2399
2400 fn with_stats(mut self, scope: stats::Scope) -> Self {
2401 self.stats = scope;
2402 self
2403 }
2404
2405 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2407 match &self.inner {
2408 RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
2409 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2410 RequestState::Pending(consumer) => Poll::Ready(
2411 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2412 Some(result) => Poll::Ready(result.clone()),
2413 None => Poll::Pending,
2414 })) {
2415 Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
2416 Err(_closed) => Err(Error::Unroutable),
2418 },
2419 ),
2420 }
2421 }
2422}
2423
2424impl kio::Pollable for Requesting {
2425 type Output = Result<broadcast::Consumer, Error>;
2426
2427 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2428 self.poll_ok(waiter)
2429 }
2430}
2431
2432pub trait Consume<T> {
2440 fn consume(&self) -> T;
2442}
2443
2444impl<T, U: Consume<T>> Consume<T> for &U {
2445 fn consume(&self) -> T {
2446 (**self).consume()
2447 }
2448}
2449
2450impl Consume<Consumer> for Producer {
2451 fn consume(&self) -> Consumer {
2452 Consumer::new(
2456 self.info,
2457 self.root.clone(),
2458 self.nodes.clone(),
2459 self.dynamic.clone(),
2460 stats::Session::default(),
2461 )
2462 }
2463}
2464
2465impl Consume<Consumer> for Consumer {
2466 fn consume(&self) -> Consumer {
2467 self.clone()
2468 }
2469}
2470
2471impl Consume<broadcast::Consumer> for broadcast::Producer {
2472 fn consume(&self) -> broadcast::Consumer {
2473 self.consume()
2475 }
2476}
2477
2478impl Consume<broadcast::Consumer> for broadcast::Consumer {
2479 fn consume(&self) -> broadcast::Consumer {
2480 self.clone()
2481 }
2482}
2483
2484impl Consume<track::Consumer> for track::Producer {
2485 fn consume(&self) -> track::Consumer {
2486 self.consume()
2487 }
2488}
2489
2490impl Consume<track::Consumer> for track::Consumer {
2491 fn consume(&self) -> track::Consumer {
2492 self.clone()
2493 }
2494}
2495
2496#[derive(Clone)]
2502pub struct Consumer {
2503 info: Origin,
2505 nodes: OriginNodes,
2506
2507 root: PathOwned,
2509
2510 dynamic: kio::Shared<OriginDynamicState>,
2513
2514 stats: stats::Session,
2518
2519 exclude: Option<Origin>,
2523}
2524
2525impl std::ops::Deref for Consumer {
2526 type Target = Origin;
2527
2528 fn deref(&self) -> &Self::Target {
2529 &self.info
2530 }
2531}
2532
2533impl Consumer {
2534 fn new(
2535 info: Origin,
2536 root: PathOwned,
2537 nodes: OriginNodes,
2538 dynamic: kio::Shared<OriginDynamicState>,
2539 stats: stats::Session,
2540 ) -> Self {
2541 Self {
2542 info,
2543 nodes,
2544 root,
2545 dynamic,
2546 stats,
2547 exclude: None,
2548 }
2549 }
2550
2551 pub(crate) fn excluding(mut self, peer: Origin) -> Self {
2556 self.exclude = Some(peer);
2557 self
2558 }
2559
2560 pub fn with_stats(mut self, session: stats::Session) -> Self {
2564 self.stats = session;
2565 self
2566 }
2567
2568 fn untagged(&self) -> Self {
2572 Self {
2573 stats: stats::Session::default(),
2574 ..self.clone()
2575 }
2576 }
2577
2578 pub(crate) fn empty(&self) -> Self {
2583 Self {
2584 info: self.info,
2585 nodes: OriginNodes { nodes: Vec::new() },
2586 root: self.root.clone(),
2587 dynamic: self.dynamic.clone(),
2588 stats: self.stats.clone(),
2589 exclude: self.exclude,
2590 }
2591 }
2592
2593 pub fn announced(&self) -> AnnounceConsumer {
2600 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
2601 }
2602
2603 pub fn consume(&self) -> Self {
2605 self.clone()
2606 }
2607
2608 fn resolve(&self, path: impl AsPath) -> Resolved {
2617 let path = path.as_path();
2618 let Some((root, rest)) = self.nodes.get(&path) else {
2619 return Resolved::Missing;
2620 };
2621 let state = root.lock();
2622 state.resolve_broadcast(&rest, self.exclude)
2623 }
2624
2625 #[cfg(test)]
2627 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2628 match self.resolve(path) {
2629 Resolved::Found(broadcast) => Some(broadcast),
2630 Resolved::Excluded | Resolved::Missing => None,
2631 }
2632 }
2633
2634 pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
2646 let path = path.as_path();
2647
2648 let consumer = self.scope(std::slice::from_ref(&path))?;
2650
2651 if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
2655 return None;
2656 }
2657
2658 let mut announced = consumer.untagged().announced();
2662 let scope = self.stats.egress(self.root.join(&path).to_owned());
2663 loop {
2664 let OriginAnnounce {
2665 path: announced_path,
2666 broadcast,
2667 } = announced.next().await?;
2668 if announced_path.as_path() == path
2670 && let Some(broadcast) = broadcast
2671 {
2672 return Some(broadcast.with_stats(scope));
2673 }
2674 }
2675 }
2676
2677 pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
2683 let prefixes = PathPrefixes::new(prefixes);
2684 Some(Consumer {
2685 info: self.info,
2686 root: self.root.clone(),
2687 nodes: self.nodes.select(&prefixes)?,
2688 dynamic: self.dynamic.clone(),
2689 stats: self.stats.clone(),
2690 exclude: self.exclude,
2691 })
2692 }
2693
2694 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
2713 let path = path.as_path();
2714
2715 let absolute = self.root.join(&path).to_owned();
2719 let scope = self.stats.egress(&absolute);
2720
2721 match self.resolve(&path) {
2727 Resolved::Found(broadcast) => return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope)),
2728 Resolved::Excluded => return kio::Pending::new(Requesting::failed(Error::Unroutable)),
2729 Resolved::Missing => {}
2730 }
2731
2732 let mut state = self.dynamic.lock();
2733
2734 if let Some(weak) = state.served.get(&absolute) {
2738 return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
2739 }
2740
2741 let consumer = if let Some(producer) = state.requests.join(&absolute) {
2744 producer.consume()
2745 } else {
2746 let producer = kio::Producer::<PendingBroadcast>::default();
2747 let consumer = producer.consume();
2748 if state.requests.insert(absolute, producer).is_err() {
2749 return kio::Pending::new(Requesting::failed(Error::Unroutable));
2750 }
2751 consumer
2752 };
2753
2754 kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
2755 }
2756
2757 pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
2762 let prefix = prefix.as_path();
2763
2764 Some(Self {
2765 info: self.info,
2766 root: self.root.join(&prefix).to_owned(),
2767 nodes: self.nodes.root(&prefix)?,
2768 dynamic: self.dynamic.clone(),
2769 stats: self.stats.clone(),
2770 exclude: self.exclude,
2771 })
2772 }
2773
2774 pub fn root(&self) -> &Path<'_> {
2776 &self.root
2777 }
2778
2779 pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
2782 self.nodes.nodes.iter().map(|(root, _)| root)
2783 }
2784
2785 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2787 self.root.join(path)
2788 }
2789}
2790
2791#[derive(Clone)]
2796pub struct AnnounceProducer {
2797 nodes: OriginNodes,
2798 root: PathOwned,
2799}
2800
2801impl AnnounceProducer {
2802 fn new(root: PathOwned, nodes: OriginNodes) -> Self {
2803 Self { nodes, root }
2804 }
2805
2806 pub fn consume(&self) -> AnnounceConsumer {
2812 AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
2815 }
2816
2817 pub fn root(&self) -> &Path<'_> {
2819 &self.root
2820 }
2821}
2822
2823pub struct AnnounceConsumer {
2828 id: ConsumerId,
2829 nodes: OriginNodes,
2830 root: PathOwned,
2831
2832 state: kio::Producer<OriginConsumerState>,
2835
2836 stats: stats::Session,
2839
2840 guards: HashMap<PathOwned, stats::Announce>,
2844}
2845
2846impl AnnounceConsumer {
2847 fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
2848 let state = kio::Producer::<OriginConsumerState>::default();
2849 let id = ConsumerId::new();
2850
2851 for (_, node) in &nodes.nodes {
2852 let notify = AnnounceConsumerNotify {
2853 root: root.clone(),
2854 state: state.clone(),
2855 };
2856 node.lock().consume(id, notify);
2857 }
2858
2859 Self {
2860 id,
2861 nodes,
2862 root,
2863 state,
2864 stats,
2865 guards: HashMap::new(),
2866 }
2867 }
2868
2869 fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
2875 let OriginAnnounce { path, broadcast } = update;
2876 let absolute = self.root.join(&path).to_owned();
2877 match broadcast {
2878 Some(broadcast) => {
2879 let scope = self.stats.egress(&absolute);
2880 self.guards.entry(absolute).or_insert_with(|| scope.announce());
2881 OriginAnnounce {
2882 path,
2883 broadcast: Some(broadcast.with_stats(scope)),
2884 }
2885 }
2886 None => {
2887 self.guards.remove(&absolute);
2888 OriginAnnounce { path, broadcast: None }
2889 }
2890 }
2891 }
2892
2893 pub async fn next(&mut self) -> Option<OriginAnnounce> {
2900 kio::wait(|waiter| self.poll_next(waiter)).await
2901 }
2902
2903 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
2909 let update = {
2910 let mut state = match ready!(self.state.poll(waiter, |state| {
2911 if state.pending.is_empty() {
2912 Poll::Pending
2913 } else {
2914 Poll::Ready(())
2915 }
2916 })) {
2917 Ok(state) => state,
2918 Err(_) => return Poll::Ready(None),
2920 };
2921 state.take().expect("predicate guaranteed an update")
2922 };
2923 Poll::Ready(Some(self.attribute(update)))
2924 }
2925
2926 pub fn try_next(&mut self) -> Option<OriginAnnounce> {
2931 let update = self.state.write().ok()?.take()?;
2932 Some(self.attribute(update))
2933 }
2934
2935 pub fn is_closed(&self) -> bool {
2937 self.state.write().is_err()
2938 }
2939
2940 pub fn root(&self) -> &Path<'_> {
2942 &self.root
2943 }
2944
2945 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
2947 self.root.join(path)
2948 }
2949}
2950
2951impl Drop for AnnounceConsumer {
2952 fn drop(&mut self) {
2953 for (_, root) in &self.nodes.nodes {
2954 root.lock().unconsume(self.id);
2955 }
2956 }
2957}
2958
2959#[cfg(test)]
2960use futures::FutureExt;
2961
2962#[cfg(test)]
2963#[allow(missing_docs)] impl AnnounceConsumer {
2965 pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2966 let expected = expected.as_path();
2967 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2968 assert_eq!(announce.path, expected, "wrong path");
2969 let announced = announce.broadcast.expect("should be an active announce");
2970 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2971 }
2972
2973 pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2977 let expected = expected.as_path();
2978 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
2979 assert_eq!(announce.path, expected, "wrong path");
2980 announce.broadcast.expect("should be an active announce")
2981 }
2982
2983 pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
2984 let expected = expected.as_path();
2985 let announce = self.try_next().expect("no next");
2986 assert_eq!(announce.path, expected, "wrong path");
2987 let announced = announce.broadcast.expect("should be an active announce");
2988 assert!(announced.is_clone(broadcast), "should be the same broadcast");
2989 }
2990
2991 pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
2993 let expected = expected.as_path();
2994 let announce = self.try_next().expect("no next");
2995 assert_eq!(announce.path, expected, "wrong path");
2996 announce.broadcast.expect("should be an active announce")
2997 }
2998
2999 pub fn assert_next_none(&mut self, expected: impl AsPath) {
3000 let expected = expected.as_path();
3001 let announce = self.next().now_or_never().expect("next blocked").expect("no next");
3002 assert_eq!(announce.path, expected, "wrong path");
3003 assert!(announce.broadcast.is_none(), "should be unannounced");
3004 }
3005
3006 pub fn assert_next_wait(&mut self) {
3007 if let Some(res) = self.next().now_or_never() {
3008 panic!("next should block: got {:?}", res.map(|a| a.path));
3009 }
3010 }
3011
3012 }
3021
3022#[cfg(test)]
3023mod tests {
3024 use crate::coding::Decode;
3025 use crate::group;
3026
3027 use super::*;
3028
3029 fn announce() -> broadcast::Route {
3031 broadcast::Route::new().with_announce(true)
3032 }
3033
3034 fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
3040 let name = Path::new(name);
3041 let peer_key = fnv_key(&name, [peer]);
3042 (100u64..)
3043 .map(|id| Origin::new(id).unwrap())
3044 .find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
3045 .unwrap()
3046 }
3047
3048 fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
3051 let source = broadcast::Info::new().produce().consume();
3052 FrontState {
3053 path: Path::new("test").to_owned(),
3054 self_origin,
3055 publisher: routes.first().and_then(|r| r.hops.iter().next().copied()),
3056 next_route: routes.len() as u64,
3057 excluded: HashMap::new(),
3058 routes: routes
3059 .into_iter()
3060 .enumerate()
3061 .map(|(id, route)| FrontRoute {
3062 id: id as u64,
3063 route,
3064 source: source.clone(),
3065 })
3066 .collect(),
3067 active: Some(0),
3068 linger: Duration::ZERO,
3069 closed: false,
3070 }
3071 }
3072
3073 fn sibling_route(peer: Origin) -> broadcast::Route {
3076 let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
3077 announce().with_hops(hops)
3078 }
3079
3080 fn upstream_route(cost: u64) -> broadcast::Route {
3082 let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
3083 announce().with_hops(hops).with_cost(cost)
3084 }
3085
3086 #[test]
3090 fn test_carrying_gate_keys() {
3091 let peer = Origin::new(3).unwrap();
3092
3093 let mut lost = front_state(
3095 origin_keyed("test", peer, false),
3096 vec![upstream_route(10), sibling_route(peer)],
3097 );
3098 lost.reselect(true);
3099 assert_eq!(
3100 lost.active,
3101 Some(0),
3102 "carrying front re-parented onto a higher-keyed peer"
3103 );
3104 lost.reselect(false);
3105 assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
3106
3107 let mut won = front_state(
3109 origin_keyed("test", peer, true),
3110 vec![upstream_route(10), sibling_route(peer)],
3111 );
3112 won.reselect(true);
3113 assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
3114 }
3115
3116 #[test]
3121 fn test_carrying_gate_symmetric_race() {
3122 let a = Origin::new(1).unwrap();
3123 let b = Origin::new(2).unwrap();
3124
3125 let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
3126 let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
3127 a_view.reselect(true);
3128 b_view.reselect(true);
3129
3130 let a_moved = a_view.active == Some(1);
3131 let b_moved = b_view.active == Some(1);
3132 assert!(
3133 a_moved != b_moved,
3134 "exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
3135 );
3136 }
3137
3138 #[test]
3143 fn test_carrying_switches_to_benign_routes() {
3144 let peer = Origin::new(3).unwrap();
3145 let lost = origin_keyed("test", peer, false);
3146
3147 let mut forwarder = sibling_route(peer).with_cost(4);
3149 forwarder.advertised = 4;
3150 let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
3151 state.reselect(true);
3152 assert_eq!(
3153 state.active,
3154 Some(1),
3155 "a cheaper forwarder path must win while carrying"
3156 );
3157
3158 let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
3160 let mut state = front_state(lost, vec![upstream_route(10), direct]);
3161 state.reselect(true);
3162 assert_eq!(
3163 state.active,
3164 Some(1),
3165 "a direct publisher route must win while carrying"
3166 );
3167
3168 let mut state = front_state(lost, vec![sibling_route(peer), sibling_route(peer)]);
3173 state.reselect(true);
3174 assert_eq!(
3175 state.active,
3176 Some(1),
3177 "a reconnect on an identical chain must win while carrying"
3178 );
3179 }
3180
3181 #[test]
3184 fn test_carrying_gate_ignores_unannounced_incumbent() {
3185 let peer = Origin::new(3).unwrap();
3186 let unannounced = upstream_route(10).with_announce(false);
3187 let mut state = front_state(
3188 origin_keyed("test", peer, false),
3189 vec![unannounced, sibling_route(peer)],
3190 );
3191 state.reselect(true);
3192 assert_eq!(
3193 state.active,
3194 Some(1),
3195 "an unannounced incumbent must always be displaced"
3196 );
3197 }
3198
3199 async fn settle() {
3202 tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
3203 }
3204
3205 async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
3208 let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
3209 .await
3210 .expect("timed out waiting for a track request")
3211 .expect("source closed");
3212 assert_eq!(request.name(), name, "unexpected track dispatched");
3213 request.accept(None)
3214 }
3215
3216 #[tokio::test]
3220 async fn test_stats_tagged_end_to_end() {
3221 use crate::Timestamp;
3222 use crate::stats::{Config, Registry, Tier};
3223 use bytes::Bytes;
3224
3225 tokio::time::pause();
3226
3227 let registry = Registry::new(Config::new());
3228 let ctx = registry.tier(Tier::default()).session("acme");
3229
3230 let origin = Origin::random().produce();
3231 let ingress = origin.clone().with_stats(ctx.clone());
3232 let egress = origin.consume().with_stats(ctx.clone());
3233
3234 let mut announced = egress.announced();
3237
3238 let source = ingress.create_broadcast("demo", announce()).unwrap();
3240 let mut dynamic = source.dynamic();
3241 settle().await;
3242 settle().await;
3243
3244 let update = announced.next().await.unwrap();
3246 assert_eq!(update.path.as_str(), "demo");
3247 let broadcast = update.broadcast.unwrap();
3248
3249 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3251 let mut producer = accept_track(&mut dynamic, "video").await;
3252 settle().await;
3253 let mut sub = subscribing.await.unwrap();
3254
3255 let mut group = producer.append_group().unwrap();
3257 group
3258 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3259 .unwrap();
3260 group
3261 .write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
3262 .unwrap();
3263 group.finish().unwrap();
3264
3265 let mut group_c = sub.recv_group().await.unwrap().unwrap();
3267 let mut frames = 0;
3268 while let Some(frame) = group_c.read_frame().await.unwrap() {
3269 assert_eq!(frame.payload.len(), 5);
3270 frames += 1;
3271 }
3272 assert_eq!(frames, 2);
3273 settle().await;
3274
3275 let report = registry.report();
3276 let entry = report
3277 .traffic
3278 .iter()
3279 .find(|e| e.path.as_str() == "demo")
3280 .expect("demo tracked");
3281 let path_len = "demo".len() as u64;
3282
3283 let egress = &entry.publisher;
3285 assert_eq!(egress.announced, 1, "one egress announce");
3286 assert_eq!(egress.announced_bytes, path_len);
3287 assert_eq!(egress.subscriptions, 1, "one egress subscription");
3288 assert_eq!(egress.broadcasts, 1, "one viewer");
3289 assert_eq!(egress.groups, 1);
3290 assert_eq!(egress.frames, 2);
3291 assert_eq!(egress.bytes, 10);
3292 assert_eq!(egress.fetches, 0);
3293
3294 let ingress = &entry.subscriber;
3296 assert_eq!(ingress.announced, 1, "one ingress announce");
3297 assert_eq!(ingress.announced_bytes, path_len);
3298 assert_eq!(ingress.subscriptions, 1, "one ingress track");
3299 assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
3300 assert_eq!(ingress.groups, 1);
3301 assert_eq!(ingress.frames, 2);
3302 assert_eq!(ingress.bytes, 10);
3303
3304 let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
3306 let _ = fetched;
3307 settle().await;
3308 let report = registry.report();
3309 let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
3310 assert_eq!(entry.publisher.fetches, 1, "one fetch");
3311 assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
3312 assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
3313 assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
3317 }
3318
3319 #[tokio::test]
3324 async fn test_stats_read_frame_counts_once() {
3325 use crate::Timestamp;
3326 use crate::stats::{Config, Registry, Tier};
3327 use bytes::Bytes;
3328
3329 tokio::time::pause();
3330
3331 let registry = Registry::new(Config::new());
3332 let ctx = registry.tier(Tier::default()).session("acme");
3333
3334 let origin = Origin::random().produce();
3335 let ingress = origin.clone().with_stats(ctx.clone());
3336 let egress = origin.consume().with_stats(ctx.clone());
3337
3338 let mut announced = egress.announced();
3339 let source = ingress.create_broadcast("demo", announce()).unwrap();
3340 let mut dynamic = source.dynamic();
3341 settle().await;
3342 settle().await;
3343
3344 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3345 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3346 let mut producer = accept_track(&mut dynamic, "video").await;
3347 settle().await;
3348 let mut sub = subscribing.await.unwrap();
3349
3350 producer
3352 .write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
3353 .unwrap();
3354
3355 let frame = sub.read_frame().await.unwrap().expect("frame");
3356 assert_eq!(frame.payload.len(), 5);
3357 settle().await;
3358
3359 let report = registry.report();
3360 let entry = report
3361 .traffic
3362 .iter()
3363 .find(|e| e.path.as_str() == "demo")
3364 .expect("demo tracked");
3365 assert_eq!(entry.publisher.groups, 1, "one group, counted once");
3366 assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
3367 assert_eq!(
3368 entry.publisher.bytes, 5,
3369 "payload counted once, not zero and not doubled"
3370 );
3371 }
3372
3373 #[tokio::test]
3377 async fn test_stats_datagrams_counted_both_sides() {
3378 use crate::Timestamp;
3379 use crate::stats::{Config, Registry, Tier};
3380
3381 tokio::time::pause();
3382
3383 let registry = Registry::new(Config::new());
3384 let ctx = registry.tier(Tier::default()).session("acme");
3385
3386 let origin = Origin::random().produce();
3387 let ingress = origin.clone().with_stats(ctx.clone());
3388 let egress = origin.consume().with_stats(ctx.clone());
3389
3390 let mut announced = egress.announced();
3391 let source = ingress.create_broadcast("demo", announce()).unwrap();
3392 let mut dynamic = source.dynamic();
3393 settle().await;
3394 settle().await;
3395
3396 let broadcast = announced.next().await.unwrap().broadcast.unwrap();
3397 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3398 let mut producer = accept_track(&mut dynamic, "video").await;
3399 settle().await;
3400 let mut sub = subscribing.await.unwrap();
3401
3402 producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
3403 let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
3404 assert_eq!(&datagram.payload[..], b"hello");
3405 settle().await;
3406
3407 let report = registry.report();
3408 let entry = report
3409 .traffic
3410 .iter()
3411 .find(|e| e.path.as_str() == "demo")
3412 .expect("demo tracked");
3413
3414 for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
3415 assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
3416 assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
3417 assert_eq!(traffic.frames, 1, "{side}: one frame");
3418 assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
3419 }
3420 }
3421
3422 #[test]
3423 fn origin_rejects_reserved_ids() {
3424 assert!(Origin::new(0).is_err());
3425 assert!(Origin::new(1u64 << 62).is_err());
3426 assert_eq!(Origin::new(1).unwrap().id(), 1);
3427
3428 let mut zero = [0u8].as_slice();
3429 assert_eq!(
3430 Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
3431 Origin::UNKNOWN
3432 );
3433 }
3434
3435 #[test]
3436 fn origin_list_push_fails_at_limit() {
3437 let mut list = OriginList::new();
3438 for _ in 0..MAX_HOPS {
3439 list.push(Origin::random()).unwrap();
3440 }
3441 assert_eq!(list.len(), MAX_HOPS);
3442 assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
3443 }
3444
3445 #[test]
3446 fn origin_list_replace_first() {
3447 let mut list = OriginList::new();
3448 for _ in 0..3 {
3449 list.push(Origin::UNKNOWN).unwrap();
3450 }
3451
3452 assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
3454 assert_eq!(
3455 list.as_slice(),
3456 &[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
3457 );
3458
3459 assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
3461 assert_eq!(list.len(), 3);
3462 }
3463
3464 #[test]
3465 fn origin_list_try_from_vec_enforces_limit() {
3466 let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
3467 assert!(OriginList::try_from(under).is_ok());
3468
3469 let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
3470 assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
3471 }
3472
3473 #[tokio::test]
3474 async fn test_announce() {
3475 tokio::time::pause();
3476
3477 let origin = Origin::random().produce();
3478
3479 let mut consumer1 = origin.consume().announced();
3480 consumer1.assert_next_wait();
3481
3482 let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
3484 settle().await;
3485
3486 consumer1.assert_next_some("test1");
3487 consumer1.assert_next_wait();
3488
3489 let mut consumer2 = origin.consume().announced();
3492
3493 let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
3495 settle().await;
3496
3497 consumer1.assert_next_some("test2");
3498 consumer1.assert_next_wait();
3499
3500 consumer2.assert_next_some("test1");
3501 consumer2.assert_next_some("test2");
3502 consumer2.assert_next_wait();
3503
3504 broadcast1.finish();
3506 settle().await;
3507
3508 consumer1.assert_next_none("test1");
3510 consumer2.assert_next_none("test1");
3511 consumer1.assert_next_wait();
3512 consumer2.assert_next_wait();
3513
3514 let mut consumer3 = origin.consume().announced();
3516 consumer3.assert_next_some("test2");
3517 consumer3.assert_next_wait();
3518
3519 broadcast2.finish();
3520 settle().await;
3521
3522 consumer1.assert_next_none("test2");
3523 consumer2.assert_next_none("test2");
3524 consumer3.assert_next_none("test2");
3525 }
3526
3527 #[tokio::test]
3531 async fn test_duplicate() {
3532 tokio::time::pause();
3533
3534 let origin = Origin::random().produce();
3535 let consumer = origin.consume();
3536 let mut announced = consumer.announced();
3537
3538 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
3539 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
3540 let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
3541 settle().await;
3542 assert!(consumer.get_broadcast("test").is_some());
3543
3544 announced.assert_next_some("test");
3545 announced.assert_next_wait();
3546
3547 broadcast2.finish();
3549 settle().await;
3550 assert!(consumer.get_broadcast("test").is_some());
3551 announced.assert_next_wait();
3552
3553 broadcast1.finish();
3555 settle().await;
3556 assert!(consumer.get_broadcast("test").is_some());
3557 announced.assert_next_wait();
3558
3559 broadcast3.finish();
3561 settle().await;
3562 assert!(consumer.get_broadcast("test").is_none());
3563
3564 announced.assert_next_none("test");
3565 announced.assert_next_wait();
3566 }
3567
3568 #[tokio::test]
3571 async fn test_route_failover() {
3572 tokio::time::pause();
3573
3574 let origin = Origin::random().produce();
3575 let consumer = origin.consume();
3576 let mut announced = consumer.announced();
3577
3578 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3581 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3582
3583 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3585 let mut dynamic_a = source_a.dynamic();
3586 settle().await;
3587 settle().await;
3588 let broadcast = consumer.request_broadcast("test").await.unwrap();
3589 announced.assert_next_some("test");
3590
3591 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3593 let mut dynamic_b = source_b.dynamic();
3594 settle().await;
3595 settle().await;
3596 announced.assert_next_wait();
3597
3598 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3600 let mut producer = accept_track(&mut dynamic_a, "video").await;
3601 settle().await;
3602 dynamic_b.assert_no_request();
3603
3604 let mut sub = subscribing.await.unwrap();
3605 sub.assert_no_group();
3608 assert_eq!(producer.subscription().unwrap().group_start, None);
3609
3610 producer.append_group().unwrap();
3611 producer.append_group().unwrap();
3612 assert_eq!(sub.assert_group().sequence, 0);
3613 assert_eq!(sub.assert_group().sequence, 1);
3614
3615 producer.abort(Error::Dropped).unwrap();
3619 source_a.abort(Error::Dropped).unwrap();
3620 drop(dynamic_a);
3621 settle().await;
3622 announced.assert_next_wait();
3623
3624 let mut producer = accept_track(&mut dynamic_b, "video").await;
3627 settle().await;
3628 sub.assert_no_group();
3629 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
3630 producer.create_group(group::Info { sequence: 1 }).unwrap();
3631 producer.create_group(group::Info { sequence: 2 }).unwrap();
3632 assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
3633 sub.assert_not_closed();
3634 }
3635
3636 #[tokio::test]
3639 async fn test_broadcast_route_watch() {
3640 let mut producer = broadcast::Info::new().produce();
3641 let mut consumer = producer.consume();
3642
3643 assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
3645
3646 producer.set_route(broadcast::Route::default()).unwrap();
3648 assert!(consumer.route_changed().now_or_never().is_none());
3649
3650 let mut hops = OriginList::new();
3651 hops.push(Origin::new(7).unwrap()).unwrap();
3652 let route = broadcast::Route::new().with_hops(hops).with_cost(3);
3653 producer.set_route(route.clone()).unwrap();
3654 assert_eq!(consumer.route_changed().await.unwrap(), route);
3655
3656 let mut fresh = producer.consume();
3658 assert_eq!(fresh.route_changed().await.unwrap(), route);
3659
3660 drop(producer);
3661 assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
3662 }
3663
3664 #[tokio::test]
3668 async fn test_route_cost_update() {
3669 tokio::time::pause();
3670
3671 let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
3675 let consumer = origin.consume();
3676 let mut announced = consumer.announced();
3677
3678 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3681 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3682
3683 let mut source_a = origin
3685 .create_broadcast("test", announce().with_hops(hops_a.clone()))
3686 .unwrap();
3687 let mut dynamic_a = source_a.dynamic();
3688 settle().await;
3689 let broadcast = consumer.request_broadcast("test").await.unwrap();
3690 announced.assert_next_some("test");
3691
3692 let mut watch = broadcast.clone();
3693 assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
3694
3695 let mut source_b = origin
3696 .create_broadcast("test", announce().with_hops(hops_b.clone()))
3697 .unwrap();
3698 let mut dynamic_b = source_b.dynamic();
3699 settle().await;
3700 assert!(
3701 watch.route_changed().now_or_never().is_none(),
3702 "a losing standby must not change the advertised route"
3703 );
3704
3705 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3707 let mut producer = accept_track(&mut dynamic_a, "video").await;
3708 settle().await;
3709 let mut sub = subscribing.await.unwrap();
3710 producer.append_group().unwrap();
3711 assert_eq!(sub.assert_group().sequence, 0);
3712
3713 source_a
3716 .set_route(announce().with_hops(hops_a.clone()).with_cost(10))
3717 .unwrap();
3718 settle().await;
3719 assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
3720 announced.assert_next_wait();
3721
3722 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3723 settle().await;
3724 sub.assert_no_group();
3727 assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
3728 producer_b.create_group(group::Info { sequence: 1 }).unwrap();
3729 assert_eq!(sub.assert_group().sequence, 1);
3730 sub.assert_not_closed();
3731
3732 source_b
3734 .set_route(announce().with_hops(hops_b.clone()).with_cost(5))
3735 .unwrap();
3736 settle().await;
3737 let advertised = watch.route_changed().await.unwrap();
3738 assert_eq!(advertised.hops, hops_b);
3739 assert_eq!(advertised.cost, 5);
3740 announced.assert_next_wait();
3741 }
3742
3743 #[tokio::test]
3746 async fn test_completed_track_survives_route_churn() {
3747 tokio::time::pause();
3748
3749 let origin = Origin::random().produce();
3750 let consumer = origin.consume();
3751
3752 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3754 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3755
3756 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
3757 let mut dynamic_a = source_a.dynamic();
3758 settle().await;
3759 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
3760 let mut dynamic_b = source_b.dynamic();
3761 settle().await;
3762 settle().await;
3763 let broadcast = consumer.request_broadcast("test").await.unwrap();
3764
3765 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3767 let mut producer = accept_track(&mut dynamic_a, "video").await;
3768 settle().await;
3769 let mut sub = subscribing.await.unwrap();
3770 producer.append_group().unwrap();
3771 assert_eq!(sub.assert_group().sequence, 0);
3772 producer.finish().unwrap();
3773 drop(producer);
3774 settle().await;
3775 sub.assert_closed();
3776
3777 source_a.abort(Error::Dropped).unwrap();
3779 drop(dynamic_a);
3780 settle().await;
3781 dynamic_b.assert_no_request();
3782
3783 let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
3785 late.assert_closed();
3786 }
3787
3788 #[tokio::test]
3792 async fn test_refused_track_aborts_instantly() {
3793 tokio::time::pause();
3794
3795 let origin = Origin::random().produce();
3796 let consumer = origin.consume();
3797
3798 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3799 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3800 let mut dynamic = source.dynamic();
3801 settle().await;
3802 settle().await;
3803 let broadcast = consumer.request_broadcast("test").await.unwrap();
3804
3805 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3806 let request = dynamic.requested_track().await.unwrap();
3807 request.reject(Error::NotFound);
3808 settle().await;
3809
3810 assert!(matches!(subscribing.await, Err(Error::NotFound)));
3812 dynamic.assert_no_request();
3813 }
3814
3815 #[tokio::test]
3820 async fn test_stale_rejection_does_not_abort_a_handover() {
3821 tokio::time::pause();
3822
3823 let origin = Origin::random().produce();
3824 let consumer = origin.consume();
3825
3826 let publisher = Origin::new(1).unwrap();
3827 let peer = Origin::new(5).unwrap();
3828 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
3829 let local = OriginList::try_from(vec![publisher]).unwrap();
3830
3831 let source_remote = origin
3833 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
3834 .unwrap();
3835 let mut dynamic_remote = source_remote.dynamic();
3836 settle().await;
3837 settle().await;
3838 let broadcast = consumer.request_broadcast("test").await.unwrap();
3839 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3840 let request_remote = dynamic_remote.requested_track().await.unwrap();
3841
3842 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
3845 let mut dynamic_local = source_local.dynamic();
3846 request_remote.reject(Error::NotFound);
3847 settle().await;
3848
3849 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
3851 settle().await;
3852 let mut sub = subscribing
3853 .await
3854 .expect("the handover must win over the stale rejection");
3855 producer_local.append_group().unwrap();
3856 assert_eq!(sub.assert_group().sequence, 0);
3857 sub.assert_not_closed();
3858 }
3859
3860 #[tokio::test]
3864 async fn test_route_handover() {
3865 tokio::time::pause();
3866
3867 let origin = Origin::random().produce();
3868 let consumer = origin.consume();
3869 let mut announced = consumer.announced();
3870
3871 let hops_long = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
3873 let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3874
3875 let source_a = origin
3876 .create_broadcast("test", announce().with_hops(hops_long))
3877 .unwrap();
3878 let mut dynamic_a = source_a.dynamic();
3879 settle().await;
3880 settle().await;
3881 let broadcast = consumer.request_broadcast("test").await.unwrap();
3882 announced.assert_next_some("test");
3883
3884 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3885 let mut producer_a = accept_track(&mut dynamic_a, "video").await;
3886 settle().await;
3887 let mut sub = subscribing.await.unwrap();
3888 producer_a.append_group().unwrap();
3889 producer_a.append_group().unwrap();
3890 assert_eq!(sub.assert_group().sequence, 0);
3891 assert_eq!(sub.assert_group().sequence, 1);
3892
3893 let source_b = origin
3896 .create_broadcast("test", announce().with_hops(hops_short))
3897 .unwrap();
3898 let mut dynamic_b = source_b.dynamic();
3899 settle().await;
3900 settle().await;
3901 announced.assert_next_wait();
3902
3903 let mut producer_b = accept_track(&mut dynamic_b, "video").await;
3904 settle().await;
3905
3906 sub.assert_no_group();
3909 assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
3910 assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
3911
3912 producer_a.create_group(group::Info { sequence: 2 }).unwrap();
3914 producer_b.create_group(group::Info { sequence: 2 }).unwrap();
3915 producer_b.create_group(group::Info { sequence: 3 }).unwrap();
3916 assert_eq!(sub.assert_group().sequence, 2);
3917 assert_eq!(sub.assert_group().sequence, 3);
3918 sub.assert_no_group();
3919 sub.assert_not_closed();
3920 }
3921
3922 #[tokio::test(start_paused = true)]
3925 async fn test_route_unannounce_immediate() {
3926 let origin = Origin::random().produce();
3927 let consumer = origin.consume();
3928 let mut announced = consumer.announced();
3929
3930 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3931 let mut source = origin
3932 .create_broadcast("test", announce().with_hops(hops.clone()))
3933 .unwrap();
3934 settle().await;
3935 let broadcast = consumer.request_broadcast("test").await.unwrap();
3936 announced.assert_next_some("test");
3937
3938 source.finish();
3941 settle().await;
3942 announced.assert_next_none("test");
3943
3944 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3946 settle().await;
3947 let fresh = consumer.request_broadcast("test").await.unwrap();
3948 announced.assert_next_some("test");
3949 assert!(
3950 !fresh.is_clone(&broadcast),
3951 "re-create must not splice the old broadcast"
3952 );
3953 }
3954
3955 #[tokio::test(start_paused = true)]
3960 async fn test_route_detach_immediate() {
3961 let origin = Origin::random().produce();
3962 let consumer = origin.consume();
3963 let mut announced = consumer.announced();
3964
3965 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
3966 let source = origin
3967 .create_broadcast("test", announce().with_hops(hops.clone()))
3968 .unwrap();
3969 let mut dynamic = source.dynamic();
3970 settle().await;
3971 settle().await;
3972 let broadcast = consumer.request_broadcast("test").await.unwrap();
3973 announced.assert_next_some("test");
3974
3975 let subscribing = broadcast.track("video").unwrap().subscribe(None);
3976 let producer = accept_track(&mut dynamic, "video").await;
3977 settle().await;
3978 let mut sub = subscribing.await.unwrap();
3979
3980 drop(producer);
3982 source.abort(Error::Dropped).unwrap();
3983 drop(dynamic);
3984
3985 settle().await;
3986 announced.assert_next_none("test");
3987 sub.assert_error();
3988
3989 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
3992 settle().await;
3993 settle().await;
3994 let fresh = consumer.request_broadcast("test").await.unwrap();
3995 announced.assert_next_some("test");
3996 assert!(
3997 !fresh.is_clone(&broadcast),
3998 "re-create must not splice the old broadcast"
3999 );
4000 }
4001
4002 #[tokio::test(start_paused = true)]
4007 async fn test_idle_track_releases_without_respinning() {
4008 let origin = Info::new(Origin::random()).produce();
4009 let consumer = origin.consume();
4010
4011 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4012 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4013 let mut dynamic = source.dynamic();
4014 settle().await;
4015 let broadcast = consumer.request_broadcast("test").await.unwrap();
4016
4017 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4018 let producer = accept_track(&mut dynamic, "video").await;
4019 settle().await;
4020 let sub = subscribing.await.unwrap();
4021
4022 drop(sub);
4025 tokio::time::sleep(TRACK_IDLE_LINGER / 2).await;
4026 settle().await;
4027 assert!(
4028 producer.poll_unused(&kio::Waiter::noop()).is_pending(),
4029 "the copy must stay spliced inside the linger",
4030 );
4031
4032 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4035 settle().await;
4036 assert!(
4037 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4038 "an idle copy must be released after the linger",
4039 );
4040
4041 for _ in 0..3 {
4045 tokio::time::sleep(TRACK_IDLE_LINGER).await;
4046 settle().await;
4047 assert!(
4048 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4049 "an unread copy must stay released, not be re-spliced",
4050 );
4051 }
4052 assert!(
4053 dynamic.requested_track().now_or_never().is_none(),
4054 "an unread track must not be re-requested",
4055 );
4056 drop(producer);
4057
4058 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4060 let mut producer = accept_track(&mut dynamic, "video").await;
4061 settle().await;
4062 let mut sub = subscribing.await.unwrap();
4063 producer.append_group().unwrap();
4064 assert_eq!(sub.assert_group().sequence, 0);
4065 }
4066
4067 #[tokio::test(start_paused = true)]
4071 async fn test_back_to_back_fetches_reuse_the_track() {
4072 let origin = Info::new(Origin::random()).produce();
4073 let consumer = origin.consume();
4074
4075 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4076 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4077 let mut dynamic = source.dynamic();
4078 settle().await;
4079 let broadcast = consumer.request_broadcast("test").await.unwrap();
4080
4081 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4083 let mut producer = accept_track(&mut dynamic, "video").await;
4084 producer.append_group().unwrap().finish().unwrap();
4085 settle().await;
4086 let first = fetching.await.expect("first fetch");
4087 drop(first);
4088
4089 settle().await;
4091 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4092 settle().await;
4093 assert!(
4094 dynamic.requested_track().now_or_never().is_none(),
4095 "a fetch inside the linger must reuse the track, not re-request it",
4096 );
4097 drop(fetching.await.expect("second fetch"));
4098
4099 tokio::time::sleep(TRACK_IDLE_LINGER * 2).await;
4101 settle().await;
4102 assert!(
4103 producer.poll_unused(&kio::Waiter::noop()).is_ready(),
4104 "the copy must be released once the fetches stop",
4105 );
4106 drop(producer);
4107
4108 settle().await;
4110 let fetching = broadcast.track("video").unwrap().fetch_group(0, None);
4111 let mut producer = accept_track(&mut dynamic, "video").await;
4112 producer.append_group().unwrap().finish().unwrap();
4113 settle().await;
4114 fetching.await.expect("fetch after the linger");
4115 }
4116
4117 #[tokio::test(start_paused = true)]
4121 async fn test_linger_reconnect_splices() {
4122 let origin = Info::new(Origin::random())
4123 .with_linger(Duration::from_secs(5))
4124 .produce();
4125 let consumer = origin.consume();
4126 let mut announced = consumer.announced();
4127
4128 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4129 let source = origin
4130 .create_broadcast("test", announce().with_hops(hops.clone()))
4131 .unwrap();
4132 let mut dynamic = source.dynamic();
4133 settle().await;
4134 settle().await;
4135 let broadcast = consumer.request_broadcast("test").await.unwrap();
4136 announced.assert_next_some("test");
4137
4138 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4139 let mut producer = accept_track(&mut dynamic, "video").await;
4140 settle().await;
4141 let mut sub = subscribing.await.unwrap();
4142
4143 producer.append_group().unwrap();
4144 producer.append_group().unwrap();
4145 assert_eq!(sub.assert_group().sequence, 0);
4146 assert_eq!(sub.assert_group().sequence, 1);
4147
4148 drop(producer);
4151 source.abort(Error::Dropped).unwrap();
4152 drop(dynamic);
4153 settle().await;
4154
4155 announced.assert_next_wait();
4157 sub.assert_no_group();
4158 sub.assert_not_closed();
4159
4160 let during = consumer.request_broadcast("test").await.unwrap();
4162 assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
4163
4164 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4167 let mut dynamic = source.dynamic();
4168 settle().await;
4169 settle().await;
4170 announced.assert_next_wait();
4171 let again = consumer.request_broadcast("test").await.unwrap();
4172 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4173
4174 let mut producer = accept_track(&mut dynamic, "video").await;
4178 settle().await;
4179 sub.assert_no_group();
4180 assert_eq!(producer.subscription().unwrap().group_start, Some(2));
4181 producer.create_group(group::Info { sequence: 2 }).unwrap();
4182 assert_eq!(sub.assert_group().sequence, 2);
4183 sub.assert_not_closed();
4184 }
4185
4186 #[tokio::test(start_paused = true)]
4189 async fn test_linger_expiry_closes() {
4190 let origin = Info::new(Origin::random())
4191 .with_linger(Duration::from_secs(5))
4192 .produce();
4193 let consumer = origin.consume();
4194 let mut announced = consumer.announced();
4195
4196 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4197 let source = origin
4198 .create_broadcast("test", announce().with_hops(hops.clone()))
4199 .unwrap();
4200 let mut dynamic = source.dynamic();
4201 settle().await;
4202 settle().await;
4203 let broadcast = consumer.request_broadcast("test").await.unwrap();
4204 announced.assert_next_some("test");
4205
4206 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4207 let producer = accept_track(&mut dynamic, "video").await;
4208 settle().await;
4209 let mut sub = subscribing.await.unwrap();
4210
4211 drop(producer);
4212 source.abort(Error::Dropped).unwrap();
4213 drop(dynamic);
4214 settle().await;
4215 announced.assert_next_wait();
4216
4217 tokio::time::sleep(std::time::Duration::from_secs(6)).await;
4219 settle().await;
4220 announced.assert_next_none("test");
4221 sub.assert_error();
4222
4223 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4225 settle().await;
4226 settle().await;
4227 let fresh = consumer.request_broadcast("test").await.unwrap();
4228 announced.assert_next_some("test");
4229 assert!(
4230 !fresh.is_clone(&broadcast),
4231 "a late re-create must not splice the expired broadcast"
4232 );
4233 }
4234
4235 #[tokio::test(start_paused = true)]
4239 async fn test_linger_forever() {
4240 let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
4241 let consumer = origin.consume();
4242 let mut announced = consumer.announced();
4243
4244 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4245 let source = origin
4246 .create_broadcast("test", announce().with_hops(hops.clone()))
4247 .unwrap();
4248 settle().await;
4249 let broadcast = consumer.request_broadcast("test").await.unwrap();
4250 announced.assert_next_some("test");
4251
4252 source.abort(Error::Dropped).unwrap();
4253 settle().await;
4254
4255 tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
4257 announced.assert_next_wait();
4258 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4259 settle().await;
4260 settle().await;
4261 let again = consumer.request_broadcast("test").await.unwrap();
4262 assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
4263 drop(source);
4264 }
4265
4266 #[tokio::test(start_paused = true)]
4280 async fn test_linger_parks_a_live_subscription() {
4281 let origin = Info::new(Origin::random())
4282 .with_linger(Duration::from_secs(5))
4283 .produce();
4284 let consumer = origin.consume();
4285
4286 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4287 let source = origin
4288 .create_broadcast("test", announce().with_hops(hops.clone()))
4289 .unwrap();
4290 let mut dynamic = source.dynamic();
4291 settle().await;
4292 settle().await;
4293 let broadcast = consumer.request_broadcast("test").await.unwrap();
4294
4295 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4298 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4299 settle().await;
4300 let mut sub = subscribing.await.unwrap();
4301 producer.append_group().unwrap();
4302 assert_eq!(sub.assert_group().sequence, 0);
4303
4304 source.abort(Error::Dropped).unwrap();
4309 settle().await;
4310 settle().await;
4311 sub.assert_not_closed();
4312
4313 tokio::time::sleep(Duration::from_secs(4)).await;
4316 settle().await;
4317 sub.assert_not_closed();
4318
4319 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4321 let mut dynamic = source.dynamic();
4322 settle().await;
4323 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4324 settle().await;
4325 producer.create_group(group::Info { sequence: 1 }).unwrap();
4326 assert_eq!(sub.assert_group().sequence, 1);
4327 sub.assert_not_closed();
4328 }
4329
4330 #[tokio::test(start_paused = true)]
4341 async fn test_idle_release_survives_the_route_leaving() {
4342 let origin = Info::new(Origin::random())
4343 .with_linger(Duration::from_secs(600))
4344 .produce();
4345 let consumer = origin.consume();
4346
4347 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4348 let source = origin
4349 .create_broadcast("test", announce().with_hops(hops.clone()))
4350 .unwrap();
4351 let mut dynamic = source.dynamic();
4352 settle().await;
4353 settle().await;
4354 let broadcast = consumer.request_broadcast("test").await.unwrap();
4355
4356 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4357 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4358 settle().await;
4359 let mut sub = subscribing.await.unwrap();
4360 producer.append_group().unwrap();
4361 assert_eq!(sub.assert_group().sequence, 0);
4362
4363 source.abort(Error::Dropped).unwrap();
4366 settle().await;
4367 settle().await;
4368 drop(sub);
4369 drop(producer);
4370 drop(dynamic);
4371 settle().await;
4372 tokio::time::sleep(TRACK_IDLE_LINGER + Duration::from_secs(1)).await;
4373 settle().await;
4374
4375 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4378 let mut dynamic = source.dynamic();
4379 settle().await;
4380 settle().await;
4381 let broadcast = consumer.request_broadcast("test").await.unwrap();
4382 let subscribing = broadcast.track("catalog.json").unwrap().subscribe(None);
4383 let mut producer = accept_track(&mut dynamic, "catalog.json").await;
4384 settle().await;
4385 let mut sub = subscribing.await.unwrap();
4386 producer.append_group().unwrap();
4387 assert_eq!(
4388 sub.assert_group().sequence,
4389 0,
4390 "the reconnect's first group must not be filtered by a stale boundary"
4391 );
4392 }
4393
4394 #[tokio::test(start_paused = true)]
4397 async fn test_linger_skipped_on_finish() {
4398 let origin = Info::new(Origin::random())
4399 .with_linger(Duration::from_secs(5))
4400 .produce();
4401 let consumer = origin.consume();
4402 let mut announced = consumer.announced();
4403
4404 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4405 let mut source = origin
4406 .create_broadcast("test", announce().with_hops(hops.clone()))
4407 .unwrap();
4408 settle().await;
4409 let broadcast = consumer.request_broadcast("test").await.unwrap();
4410 announced.assert_next_some("test");
4411
4412 source.finish();
4415 settle().await;
4416 announced.assert_next_none("test");
4417
4418 let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
4420 settle().await;
4421 let fresh = consumer.request_broadcast("test").await.unwrap();
4422 announced.assert_next_some("test");
4423 assert!(
4424 !fresh.is_clone(&broadcast),
4425 "a finish must not leave a lingering broadcast to splice into"
4426 );
4427 }
4428
4429 #[tokio::test]
4432 async fn test_announce_toggle() {
4433 tokio::time::pause();
4434
4435 let origin = Origin::random().produce();
4436 let consumer = origin.consume();
4437 let mut announced = consumer.announced();
4438
4439 let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4440 settle().await;
4441
4442 announced.assert_next_wait();
4444 let broadcast = consumer
4445 .get_broadcast("test")
4446 .expect("offline broadcast is still routable");
4447 assert!(!broadcast.route().announce);
4448
4449 let requested = consumer.request_broadcast("test").await.unwrap();
4451 assert!(requested.is_clone(&broadcast));
4452
4453 source.set_route(announce()).unwrap();
4455 settle().await;
4456 let face = announced.assert_next_some("test");
4457 assert!(face.is_clone(&broadcast));
4458
4459 let mut fresh = origin.consume().announced();
4461 fresh.assert_next_some("test");
4462 fresh.assert_next_wait();
4463
4464 source.set_route(broadcast::Route::new()).unwrap();
4466 settle().await;
4467 announced.assert_next_none("test");
4468 assert!(consumer.get_broadcast("test").is_some());
4469 let mut fresh = origin.consume().announced();
4470 fresh.assert_next_wait();
4471
4472 source.finish();
4473 settle().await;
4474 assert!(consumer.get_broadcast("test").is_none());
4475 }
4476
4477 #[tokio::test]
4480 async fn test_announce_beats_offline() {
4481 tokio::time::pause();
4482
4483 let origin = Origin::random().produce();
4484 let consumer = origin.consume();
4485 let mut announced = consumer.announced();
4486
4487 let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
4489 settle().await;
4490 announced.assert_next_wait();
4491
4492 let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
4495 settle().await;
4496 announced.assert_next_some("test");
4497 let face = consumer.get_broadcast("test").unwrap();
4498 assert!(face.route().announce);
4499 assert_eq!(face.route().cost, 10);
4500
4501 announced_source.finish();
4504 settle().await;
4505 announced.assert_next_none("test");
4506 assert!(consumer.get_broadcast("test").is_some());
4507 }
4508
4509 #[tokio::test]
4512 async fn test_better_source_no_churn() {
4513 tokio::time::pause();
4514
4515 let origin = Origin::random().produce();
4516 let mut announced = origin.consume().announced();
4517
4518 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap(), Origin::new(3).unwrap()]).unwrap();
4521 let hops_b = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4522 let _a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4523 settle().await;
4524 let face = announced.assert_next_some("test");
4525
4526 let _b = origin
4527 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4528 .unwrap();
4529 settle().await;
4530 announced.assert_next_wait();
4531 let current = origin.consume().get_broadcast("test").unwrap();
4532 assert!(current.is_clone(&face), "the broadcast identity must not change");
4533 assert_eq!(current.route().hops, hops_b);
4535 }
4536
4537 #[tokio::test]
4543 async fn test_publisher_mismatch_replaces() {
4544 tokio::time::pause();
4545
4546 let origin = Origin::random().produce();
4547 let consumer = origin.consume();
4548 let mut announced = consumer.announced();
4549
4550 let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4551 let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4552
4553 let mut source_a = origin
4554 .create_broadcast("test", announce().with_hops(hops_a.clone()))
4555 .unwrap();
4556 settle().await;
4557 let face_a = announced.assert_next_some("test");
4558
4559 let _source_b = origin
4562 .create_broadcast("test", announce().with_hops(hops_b.clone()))
4563 .unwrap();
4564 settle().await;
4565 settle().await;
4566 announced.assert_next_none("test");
4567 let face_b = announced.assert_next_some("test");
4568 assert!(!face_b.is_clone(&face_a), "a replacement, never a splice");
4569 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4570 assert!(face_a.is_closed(), "the displaced front must close");
4574
4575 source_a.finish();
4577 settle().await;
4578 settle().await;
4579 announced.assert_next_wait();
4580 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_b);
4581 }
4582
4583 #[tokio::test]
4588 async fn test_reconnect_wins_over_stale_route() {
4589 tokio::time::pause();
4590
4591 let origin = Origin::random().produce();
4592 let consumer = origin.consume();
4593
4594 let publisher = Origin::new(1).unwrap();
4595 let hops = OriginList::try_from(vec![publisher]).unwrap();
4596
4597 let stale = origin
4600 .create_broadcast("test", announce().with_hops(hops.clone()))
4601 .unwrap();
4602 let mut stale_dynamic = stale.dynamic();
4603 settle().await;
4604
4605 let fresh = origin
4607 .create_broadcast("test", announce().with_hops(hops.clone()))
4608 .unwrap();
4609 let mut fresh_dynamic = fresh.dynamic();
4610 settle().await;
4611 settle().await;
4612
4613 let broadcast = consumer.request_broadcast("test").await.unwrap();
4615 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4616 settle().await;
4617 let _producer = accept_track(&mut fresh_dynamic, "video").await;
4618 settle().await;
4619 subscribing.await.unwrap();
4620 stale_dynamic.assert_no_request();
4621 }
4622
4623 #[tokio::test]
4629 async fn test_carrying_reconnect_switches_immediately() {
4630 tokio::time::pause();
4631
4632 let origin = Origin::random().produce();
4633 let consumer = origin.consume();
4634
4635 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4636
4637 let stale = origin
4638 .create_broadcast("test", announce().with_hops(hops.clone()))
4639 .unwrap();
4640 let mut stale_dynamic = stale.dynamic();
4641 settle().await;
4642
4643 let broadcast = consumer.request_broadcast("test").await.unwrap();
4645 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4646 settle().await;
4647 let _stale_producer = accept_track(&mut stale_dynamic, "video").await;
4648 settle().await;
4649 let _subscription = subscribing.await.unwrap();
4651
4652 let fresh = origin
4654 .create_broadcast("test", announce().with_hops(hops.clone()))
4655 .unwrap();
4656 let mut fresh_dynamic = fresh.dynamic();
4657 settle().await;
4658 settle().await;
4659
4660 let _fresh_producer = accept_track(&mut fresh_dynamic, "video").await;
4663 }
4664
4665 #[tokio::test]
4675 async fn test_offline_mismatch_never_evicts_a_live_front() {
4676 tokio::time::pause();
4677
4678 let origin = Origin::random().produce();
4679 let consumer = origin.consume();
4680 let mut announced = consumer.announced();
4681
4682 let hops_live = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4683 let hops_cache = OriginList::try_from(vec![Origin::new(2).unwrap()]).unwrap();
4684
4685 let mut live = origin
4686 .create_broadcast("test", announce().with_hops(hops_live.clone()))
4687 .unwrap();
4688 let mut live_dynamic = live.dynamic();
4689 settle().await;
4690 let face = announced.assert_next_some("test");
4691
4692 let broadcast = consumer.request_broadcast("test").await.unwrap();
4694 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4695 settle().await;
4696 let _producer = accept_track(&mut live_dynamic, "video").await;
4697 settle().await;
4698 subscribing.await.unwrap();
4699
4700 let cache = origin
4703 .create_broadcast("test", broadcast::Route::new().with_hops(hops_cache.clone()))
4704 .unwrap();
4705 settle().await;
4706 settle().await;
4707 announced.assert_next_wait();
4708 assert!(!face.is_closed(), "the live front must survive");
4709 assert_eq!(consumer.get_broadcast("test").unwrap().route().hops, hops_live);
4710
4711 live.finish();
4714 settle().await;
4715 settle().await;
4716 announced.assert_next_none("test");
4717 let taken = consumer
4718 .get_broadcast("test")
4719 .expect("the parked source must take over");
4720 assert_eq!(taken.route().hops, hops_cache);
4721 announced.assert_next_wait();
4723 drop(cache);
4724 }
4725
4726 #[tokio::test]
4732 async fn test_dispatch_excludes_requester() {
4733 tokio::time::pause();
4734
4735 let origin = Origin::random().produce();
4736 let consumer = origin.consume();
4737
4738 let peer = Origin::new(5).unwrap();
4739 let publisher = Origin::new(1).unwrap();
4740 let tainted = OriginList::try_from(vec![publisher, peer]).unwrap();
4742 let clean = OriginList::try_from(vec![publisher]).unwrap();
4743
4744 let source_a = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
4745 let mut dynamic_a = source_a.dynamic();
4746 settle().await;
4747 let source_b = origin
4748 .create_broadcast("test", announce().with_hops(clean).with_cost(5))
4749 .unwrap();
4750 let mut dynamic_b = source_b.dynamic();
4751 settle().await;
4752 settle().await;
4753
4754 let shared = consumer.request_broadcast("test").await.unwrap();
4757 let subscribing = shared.track("video").unwrap().subscribe(None);
4758 let _producer_a = accept_track(&mut dynamic_a, "video").await;
4759 settle().await;
4760 subscribing.await.unwrap();
4761
4762 let scoped = consumer.clone().excluding(peer);
4767 let pinned = scoped.request_broadcast("test").await.unwrap();
4768 let subscribing = pinned.track("video").unwrap().subscribe(None);
4769 let _producer_b = accept_track(&mut dynamic_b, "video").await;
4770 settle().await;
4771 subscribing.await.unwrap();
4772 dynamic_a.assert_no_request();
4773 }
4774
4775 #[tokio::test]
4780 async fn test_unknown_publishers_do_not_splice() {
4781 tokio::time::pause();
4782
4783 let origin = Origin::random().produce();
4784 let consumer = origin.consume();
4785 let mut announced = consumer.announced();
4786
4787 let unknown_a = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
4788 let unknown_b = OriginList::try_from(vec![Origin::UNKNOWN]).unwrap();
4789
4790 let source_a = origin
4791 .create_broadcast("test", announce().with_hops(unknown_a))
4792 .unwrap();
4793 settle().await;
4794 settle().await;
4795 announced.assert_next_some("test");
4796
4797 let source_b = origin
4801 .create_broadcast("test", announce().with_hops(unknown_b))
4802 .unwrap();
4803 settle().await;
4804 settle().await;
4805 announced.assert_next_none("test");
4806 announced.assert_next_some("test");
4807
4808 drop(source_a);
4809 drop(source_b);
4810 }
4811
4812 #[tokio::test]
4815 async fn test_known_publishers_still_splice() {
4816 tokio::time::pause();
4817
4818 let origin = Origin::random().produce();
4819 let consumer = origin.consume();
4820 let mut announced = consumer.announced();
4821
4822 let publisher = Origin::new(1).unwrap();
4823 let hops_a = OriginList::try_from(vec![publisher]).unwrap();
4824 let hops_b = OriginList::try_from(vec![publisher, Origin::new(3).unwrap()]).unwrap();
4825
4826 let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
4827 settle().await;
4828 settle().await;
4829 announced.assert_next_some("test");
4830
4831 let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
4833 settle().await;
4834 settle().await;
4835 announced.assert_next_wait();
4836
4837 drop(source_a);
4838 drop(source_b);
4839 }
4840
4841 #[tokio::test]
4847 async fn test_standby_join_splices_live_subscriber() {
4848 tokio::time::pause();
4849
4850 let origin = Origin::random().produce();
4851 let consumer = origin.consume();
4852
4853 let publisher = Origin::new(1).unwrap();
4854 let peer = Origin::new(5).unwrap();
4855 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4856 let local = OriginList::try_from(vec![publisher]).unwrap();
4857
4858 let source_remote = origin
4860 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4861 .unwrap();
4862 let mut dynamic_remote = source_remote.dynamic();
4863 settle().await;
4864 settle().await;
4865 let broadcast = consumer.request_broadcast("test").await.unwrap();
4866 let subscribing = broadcast.track("video").unwrap().subscribe(None);
4867 let mut producer_remote = accept_track(&mut dynamic_remote, "video").await;
4868 settle().await;
4869 let mut sub = subscribing.await.unwrap();
4870 producer_remote.append_group().unwrap();
4871 assert_eq!(sub.assert_group().sequence, 0);
4872
4873 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4876 let mut dynamic_local = source_local.dynamic();
4877 settle().await;
4878 let mut producer_local = accept_track(&mut dynamic_local, "video").await;
4879 settle().await;
4880 sub.assert_no_group();
4881 assert_eq!(producer_local.subscription().unwrap().group_start, Some(1));
4882 producer_local.create_group(group::Info { sequence: 1 }).unwrap();
4883 assert_eq!(sub.assert_group().sequence, 1);
4884 sub.assert_not_closed();
4885 }
4886
4887 #[tokio::test]
4894 async fn test_standby_missing_track_keeps_incumbent() {
4895 tokio::time::pause();
4896
4897 let origin = Origin::random().produce();
4898 let consumer = origin.consume();
4899
4900 let publisher = Origin::new(1).unwrap();
4901 let peer = Origin::new(5).unwrap();
4902 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
4903 let local = OriginList::try_from(vec![publisher]).unwrap();
4904
4905 let source_remote = origin
4907 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
4908 .unwrap();
4909 let mut dynamic_remote = source_remote.dynamic();
4910 settle().await;
4911 settle().await;
4912 let broadcast = consumer.request_broadcast("test").await.unwrap();
4913 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4914 let mut producer_remote = accept_track(&mut dynamic_remote, "audio").await;
4915 settle().await;
4916 let mut sub = subscribing.await.unwrap();
4917 producer_remote.append_group().unwrap();
4918 assert_eq!(sub.assert_group().sequence, 0);
4919
4920 let source_local = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
4923 let mut dynamic_local = source_local.dynamic();
4924 settle().await;
4925 let request = dynamic_local.requested_track().await.unwrap();
4926 assert_eq!(request.name(), "audio");
4927 request.reject(Error::NotFound);
4928 settle().await;
4929
4930 producer_remote.append_group().unwrap();
4932 assert_eq!(sub.assert_group().sequence, 1);
4933 sub.assert_not_closed();
4934
4935 source_remote.abort(Error::Dropped).unwrap();
4938 settle().await;
4939 settle().await;
4940 sub.assert_closed();
4941 dynamic_local.assert_no_request();
4942
4943 let retry = broadcast.track("audio").unwrap().subscribe(None);
4945 let mut producer_local = accept_track(&mut dynamic_local, "audio").await;
4946 settle().await;
4947 let mut sub = retry.await.expect("a fresh request must reach the standby");
4948 producer_local.create_group(group::Info { sequence: 2 }).unwrap();
4949 assert_eq!(sub.assert_group().sequence, 2);
4950 }
4951
4952 #[tokio::test]
4957 async fn test_unservable_track_retried_by_a_later_request() {
4958 tokio::time::pause();
4959
4960 let origin = Origin::random().produce();
4961 let consumer = origin.consume();
4962
4963 let source = origin.create_broadcast("test", announce()).unwrap();
4964 let mut dynamic = source.dynamic();
4965 settle().await;
4966 settle().await;
4967 let broadcast = consumer.request_broadcast("test").await.unwrap();
4968
4969 let subscribing = broadcast.track("audio").unwrap().subscribe(None);
4971 let request = dynamic.requested_track().await.unwrap();
4972 request.reject(Error::NotFound);
4973 settle().await;
4974 assert!(matches!(subscribing.await, Err(Error::NotFound)));
4975
4976 let retry = broadcast.track("audio").unwrap().subscribe(None);
4978 let mut producer = accept_track(&mut dynamic, "audio").await;
4979 settle().await;
4980 let mut sub = retry.await.expect("a fresh request must reach the source");
4981 producer.append_group().unwrap();
4982 assert_eq!(sub.assert_group().sequence, 0);
4983 }
4984
4985 #[tokio::test]
4991 async fn test_track_dying_without_progress_aborts() {
4992 tokio::time::pause();
4993
4994 let origin = Origin::random().produce();
4995 let consumer = origin.consume();
4996
4997 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
4998 let source = origin.create_broadcast("test", announce().with_hops(hops)).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("video").unwrap().subscribe(None);
5005 let producer = accept_track(&mut dynamic, "video").await;
5006 settle().await;
5007 let mut sub = subscribing.await.unwrap();
5008
5009 drop(producer);
5012 settle().await;
5013 sub.assert_closed();
5014 dynamic.assert_no_request();
5015
5016 let retry = broadcast.track("video").unwrap().subscribe(None);
5018 let mut producer = accept_track(&mut dynamic, "video").await;
5019 settle().await;
5020 let mut sub = retry.await.expect("a fresh request must reach the source");
5021 producer.append_group().unwrap();
5022 assert_eq!(sub.assert_group().sequence, 0);
5023 }
5024
5025 #[tokio::test]
5030 async fn test_delivered_copy_death_survives_unrelated_wakes() {
5031 tokio::time::pause();
5032
5033 let origin = Origin::random().produce();
5034 let consumer = origin.consume();
5035
5036 let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
5037 let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
5038 let mut dynamic = source.dynamic();
5039 settle().await;
5040 settle().await;
5041 let broadcast = consumer.request_broadcast("test").await.unwrap();
5042
5043 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5044 let mut producer = accept_track(&mut dynamic, "video").await;
5045 settle().await;
5046 let mut sub = subscribing.await.unwrap();
5047 producer.append_group().unwrap();
5048 assert_eq!(sub.assert_group().sequence, 0);
5049
5050 drop(sub);
5053 settle().await;
5054 let resubscribing = broadcast.track("video").unwrap().subscribe(None);
5055 settle().await;
5056 let mut sub = resubscribing.await.unwrap();
5057 assert_eq!(sub.assert_group().sequence, 0, "cached group re-served");
5058
5059 drop(producer);
5062 let mut producer = accept_track(&mut dynamic, "video").await;
5063 settle().await;
5064 producer.create_group(group::Info { sequence: 1 }).unwrap();
5065 assert_eq!(sub.assert_group().sequence, 1);
5066 sub.assert_not_closed();
5067 }
5068
5069 #[tokio::test]
5075 async fn test_per_track_fallback_respects_exclusion() {
5076 tokio::time::pause();
5077
5078 let origin = Origin::random().produce();
5079 let consumer = origin.consume();
5080
5081 let publisher = Origin::new(1).unwrap();
5082 let peer = Origin::new(5).unwrap();
5083 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5084 let local = OriginList::try_from(vec![publisher]).unwrap();
5085
5086 let source_tainted = origin
5088 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5089 .unwrap();
5090 let mut dynamic_tainted = source_tainted.dynamic();
5091 settle().await;
5092 settle().await;
5093
5094 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5097 let mut dynamic_clean = source_clean.dynamic();
5098 settle().await;
5099
5100 let scoped = consumer.clone().excluding(peer);
5101 let broadcast = scoped.request_broadcast("test").await.unwrap();
5102 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5103 settle().await;
5104
5105 let request = dynamic_clean.requested_track().await.unwrap();
5109 request.reject(Error::NotFound);
5110 settle().await;
5111 dynamic_tainted.assert_no_request();
5112 }
5113
5114 #[tokio::test]
5118 async fn test_exclusion_survives_failover_onto_a_tainted_route() {
5119 tokio::time::pause();
5120
5121 let origin = Origin::random().produce();
5122 let consumer = origin.consume();
5123
5124 let publisher = Origin::new(1).unwrap();
5125 let peer = Origin::new(5).unwrap();
5126 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5127 let local = OriginList::try_from(vec![publisher]).unwrap();
5128
5129 let source_tainted = origin
5130 .create_broadcast("test", announce().with_hops(via_peer).with_cost(2))
5131 .unwrap();
5132 let mut dynamic_tainted = source_tainted.dynamic();
5133 settle().await;
5134 settle().await;
5135 let source_clean = origin.create_broadcast("test", announce().with_hops(local)).unwrap();
5136 let mut dynamic_clean = source_clean.dynamic();
5137 settle().await;
5138
5139 let scoped = consumer.clone().excluding(peer);
5140 let broadcast = scoped.request_broadcast("test").await.unwrap();
5141 let _subscribing = broadcast.track("video").unwrap().subscribe(None);
5142 let _clean = accept_track(&mut dynamic_clean, "video").await;
5143 settle().await;
5144
5145 source_clean.abort(Error::Dropped).unwrap();
5147 settle().await;
5148 settle().await;
5149 dynamic_tainted.assert_no_request();
5150
5151 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5154 }
5155
5156 #[tokio::test]
5161 async fn test_exclusion_holds_when_a_tainted_route_attaches_later() {
5162 tokio::time::pause();
5163
5164 let origin = Origin::random().produce();
5165 let consumer = origin.consume();
5166
5167 let publisher = Origin::new(1).unwrap();
5168 let peer = Origin::new(5).unwrap();
5169 let local = OriginList::try_from(vec![publisher]).unwrap();
5170 let via_peer = OriginList::try_from(vec![publisher, peer]).unwrap();
5171
5172 let source_clean = origin
5176 .create_broadcast("test", announce().with_hops(local).with_cost(5))
5177 .unwrap();
5178 let mut dynamic_clean = source_clean.dynamic();
5179 settle().await;
5180 settle().await;
5181 let scoped = consumer.clone().excluding(peer);
5182 let broadcast = scoped.request_broadcast("test").await.unwrap();
5183 let subscribing = broadcast.track("video").unwrap().subscribe(None);
5184 let mut producer_clean = accept_track(&mut dynamic_clean, "video").await;
5185 settle().await;
5186 let mut sub = subscribing.await.unwrap();
5187 producer_clean.append_group().unwrap();
5188 assert_eq!(sub.assert_group().sequence, 0);
5189
5190 let mut tainted = announce().with_hops(via_peer.clone()).with_cost(0);
5197 tainted.advertised = 1;
5198 let mut source_tainted = origin.create_broadcast("test", tainted).unwrap();
5199 let mut dynamic_tainted = source_tainted.dynamic();
5200 settle().await;
5201 settle().await;
5202 dynamic_tainted.assert_no_request();
5203 producer_clean.append_group().unwrap();
5204 assert_eq!(sub.assert_group().sequence, 1);
5205 sub.assert_not_closed();
5206
5207 drop(sub);
5210 drop(broadcast);
5211 drop(scoped);
5212 settle().await;
5213 let mut bumped = announce().with_hops(via_peer).with_cost(1);
5214 bumped.advertised = 1;
5215 source_tainted.set_route(bumped).unwrap();
5216 settle().await;
5217 let plain = consumer.request_broadcast("test").await.unwrap();
5218 let _plain_track = plain.track("video").unwrap().subscribe(None);
5219 settle().await;
5220 settle().await;
5221 assert!(
5222 dynamic_tainted.requested_track().now_or_never().is_some(),
5223 "the front must be free to use the route again once the peer is gone"
5224 );
5225 }
5226
5227 #[tokio::test]
5231 async fn test_excluded_path_never_reaches_the_dynamic_handler() {
5232 tokio::time::pause();
5233
5234 let origin = Origin::random().produce();
5235 let consumer = origin.consume();
5236 let mut dynamic = origin.dynamic();
5237
5238 let peer = Origin::new(5).unwrap();
5239 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5240 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5241 settle().await;
5242 settle().await;
5243
5244 let scoped = consumer.clone().excluding(peer);
5245 assert!(matches!(scoped.request_broadcast("test").await, Err(Error::Unroutable)));
5246 assert!(
5247 dynamic.requested_broadcast().now_or_never().is_none(),
5248 "the dynamic handler was asked to route around the exclusion"
5249 );
5250
5251 let _pending = scoped.request_broadcast("other");
5253 settle().await;
5254 assert!(
5255 dynamic.requested_broadcast().now_or_never().is_some(),
5256 "a genuinely missing path must still fall back"
5257 );
5258 }
5259
5260 #[tokio::test]
5263 async fn test_dispatch_all_tainted_unroutable() {
5264 tokio::time::pause();
5265
5266 let origin = Origin::random().produce();
5267 let consumer = origin.consume();
5268
5269 let peer = Origin::new(5).unwrap();
5270 let tainted = OriginList::try_from(vec![Origin::new(1).unwrap(), peer]).unwrap();
5271 let _source = origin.create_broadcast("test", announce().with_hops(tainted)).unwrap();
5272 settle().await;
5273 settle().await;
5274
5275 let scoped = consumer.clone().excluding(peer);
5276 match scoped.request_broadcast("test").await {
5277 Err(Error::Unroutable) => {}
5278 Err(err) => panic!("expected Unroutable, got {err:?}"),
5279 Ok(_) => panic!("expected Unroutable, got a broadcast"),
5280 }
5281
5282 consumer.request_broadcast("test").await.unwrap();
5284 }
5285
5286 #[tokio::test]
5287 async fn test_duplicate_reverse() {
5288 tokio::time::pause();
5289
5290 let origin = Origin::random().produce();
5291
5292 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
5293 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
5294 settle().await;
5295 assert!(origin.consume().get_broadcast("test").is_some());
5296
5297 broadcast2.finish();
5299 settle().await;
5300 assert!(origin.consume().get_broadcast("test").is_some());
5301
5302 broadcast1.finish();
5303 settle().await;
5304 assert!(origin.consume().get_broadcast("test").is_none());
5305 }
5306
5307 #[tokio::test]
5308 async fn test_deterministic_tiebreak() {
5309 tokio::time::pause();
5310
5311 fn hops(ids: &[u64]) -> OriginList {
5312 OriginList::try_from(
5313 ids.iter()
5314 .copied()
5315 .map(|id| Origin::new(id).unwrap())
5316 .collect::<Vec<_>>(),
5317 )
5318 .unwrap()
5319 }
5320
5321 async fn winner(first: &[u64], second: &[u64]) -> OriginList {
5324 let origin = Origin::random().produce();
5325 let _a = origin
5326 .create_broadcast("test", announce().with_hops(hops(first)))
5327 .unwrap();
5328 let _b = origin
5329 .create_broadcast("test", announce().with_hops(hops(second)))
5330 .unwrap();
5331 settle().await;
5332 origin.consume().get_broadcast("test").unwrap().route().hops
5333 }
5334
5335 let forward = winner(&[5, 20], &[5, 40]).await;
5339 let reverse = winner(&[5, 40], &[5, 20]).await;
5340 assert_eq!(forward, reverse, "tie-break must not depend on publish order");
5341
5342 assert_eq!(winner(&[5, 20], &[5]).await.len(), 1);
5344 assert_eq!(winner(&[5], &[5, 20]).await.len(), 1);
5345 }
5346
5347 #[tokio::test]
5352 async fn test_many_announces() {
5353 let origin = Origin::random().produce();
5354
5355 let mut consumer = origin.consume().announced();
5356 let mut broadcasts = Vec::new();
5358 for i in 0..256 {
5359 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5360 settle().await;
5361 }
5362
5363 for i in 0..256 {
5364 consumer.assert_next_some(format!("test{i:03}"));
5365 }
5366 consumer.assert_next_wait();
5367 }
5368
5369 #[tokio::test]
5370 async fn test_many_announces_try() {
5371 let origin = Origin::random().produce();
5372
5373 let mut consumer = origin.consume().announced();
5374 let mut broadcasts = Vec::new();
5376 for i in 0..256 {
5377 broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
5378 settle().await;
5379 }
5380
5381 for i in 0..256 {
5382 consumer.assert_try_next_some(format!("test{i:03}"));
5383 }
5384 }
5385
5386 #[tokio::test]
5387 async fn test_with_root_basic() {
5388 let origin = Origin::random().produce();
5389
5390 let foo_producer = origin.with_root("foo").expect("should create root");
5392 assert_eq!(foo_producer.root().as_str(), "foo");
5393
5394 let mut consumer = origin.consume().announced();
5395
5396 let _broadcast = foo_producer
5398 .create_broadcast("bar/baz", announce())
5399 .expect("publish allowed");
5400 settle().await;
5401 consumer.assert_next_some("foo/bar/baz");
5403
5404 let mut foo_consumer = foo_producer.consume().announced();
5406 foo_consumer.assert_next_some("bar/baz");
5407 }
5408
5409 #[tokio::test]
5410 async fn test_with_root_nested() {
5411 let origin = Origin::random().produce();
5412
5413 let foo_producer = origin.with_root("foo").expect("should create foo root");
5415 let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
5416 assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
5417
5418 let mut consumer = origin.consume().announced();
5419
5420 let _broadcast = foo_bar_producer
5422 .create_broadcast("baz", announce())
5423 .expect("publish allowed");
5424 settle().await;
5425 consumer.assert_next_some("foo/bar/baz");
5427
5428 let mut foo_bar_consumer = foo_bar_producer.consume().announced();
5430 foo_bar_consumer.assert_next_some("baz");
5431 }
5432
5433 #[tokio::test]
5434 async fn test_publish_scope_allows() {
5435 let origin = Origin::random().produce();
5436
5437 let limited_producer = origin
5439 .scope(&["allowed/path1".into(), "allowed/path2".into()])
5440 .expect("should create limited producer");
5441
5442 let _broadcast = limited_producer
5444 .create_broadcast("allowed/path1", announce())
5445 .expect("publish allowed");
5446 let _keep2 = limited_producer
5447 .create_broadcast("allowed/path1/nested", announce())
5448 .expect("publish allowed");
5449 let _keep3 = limited_producer
5450 .create_broadcast("allowed/path2", announce())
5451 .expect("publish allowed");
5452 settle().await;
5453
5454 assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
5456 assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
5458 }
5459
5460 #[tokio::test]
5461 async fn test_publish_max_parts() {
5462 let origin = Origin::random().produce();
5463
5464 let at_limit = (0..Path::MAX_PARTS)
5465 .map(|i| i.to_string())
5466 .collect::<Vec<_>>()
5467 .join("/");
5468 let _broadcast = origin
5469 .create_broadcast(at_limit.as_str(), announce())
5470 .expect("publish allowed");
5471 settle().await;
5472
5473 let too_deep = format!("{at_limit}/extra");
5474 assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
5475
5476 let rooted = origin.with_root("root").expect("wildcard allows any root");
5478 assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
5479 }
5480
5481 #[tokio::test]
5482 async fn test_publish_scope_empty() {
5483 let origin = Origin::random().produce();
5484
5485 assert!(origin.scope(&[]).is_none());
5487 }
5488
5489 #[tokio::test]
5490 async fn test_consume_scope_filters() {
5491 let origin = Origin::random().produce();
5492
5493 let mut consumer = origin.consume().announced();
5494
5495 let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
5497 let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
5498 let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
5499 settle().await;
5500
5501 let mut limited_consumer = origin
5503 .consume()
5504 .scope(&["allowed".into()])
5505 .expect("should create limited consumer")
5506 .announced();
5507
5508 limited_consumer.assert_next_some("allowed");
5510 limited_consumer.assert_next_some("allowed/nested");
5511 limited_consumer.assert_next_wait(); consumer.assert_next_some("allowed");
5515 consumer.assert_next_some("allowed/nested");
5516 consumer.assert_next_some("notallowed");
5517 }
5518
5519 #[tokio::test]
5520 async fn test_consume_scope_multiple_prefixes() {
5521 let origin = Origin::random().produce();
5522
5523 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5524 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5525 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5526 settle().await;
5527
5528 let mut limited_consumer = origin
5530 .consume()
5531 .scope(&["foo".into(), "bar".into()])
5532 .expect("should create limited consumer")
5533 .announced();
5534
5535 limited_consumer.assert_next_some("bar/test");
5537 limited_consumer.assert_next_some("foo/test");
5538 limited_consumer.assert_next_wait(); }
5540
5541 #[tokio::test]
5542 async fn test_with_root_and_publish_scope() {
5543 let origin = Origin::random().produce();
5544
5545 let foo_producer = origin.with_root("foo").expect("should create foo root");
5547
5548 let limited_producer = foo_producer
5550 .scope(&["bar".into(), "goop/pee".into()])
5551 .expect("should create limited producer");
5552
5553 let mut consumer = origin.consume().announced();
5554
5555 let _broadcast = limited_producer
5557 .create_broadcast("bar", announce())
5558 .expect("publish allowed");
5559 let _keep2 = limited_producer
5560 .create_broadcast("bar/nested", announce())
5561 .expect("publish allowed");
5562 let _keep3 = limited_producer
5563 .create_broadcast("goop/pee", announce())
5564 .expect("publish allowed");
5565 let _keep4 = limited_producer
5566 .create_broadcast("goop/pee/nested", announce())
5567 .expect("publish allowed");
5568 settle().await;
5569
5570 assert!(limited_producer.create_broadcast("baz", announce()).is_err());
5572 assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
5574
5575 consumer.assert_next_some("foo/bar");
5577 consumer.assert_next_some("foo/bar/nested");
5578 consumer.assert_next_some("foo/goop/pee");
5579 consumer.assert_next_some("foo/goop/pee/nested");
5580 }
5581
5582 #[tokio::test]
5583 async fn test_with_root_and_consume_scope() {
5584 let origin = Origin::random().produce();
5585
5586 let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
5588 let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
5589 let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
5590 settle().await;
5591
5592 let foo_producer = origin.with_root("foo").expect("should create foo root");
5594
5595 let mut limited_consumer = foo_producer
5597 .consume()
5598 .scope(&["bar".into(), "goop/pee".into()])
5599 .expect("should create limited consumer")
5600 .announced();
5601
5602 limited_consumer.assert_next_some("bar/test");
5604 limited_consumer.assert_next_some("goop/pee/test");
5605 limited_consumer.assert_next_wait(); }
5607
5608 #[tokio::test]
5609 async fn test_with_root_unauthorized() {
5610 let origin = Origin::random().produce();
5611
5612 let limited_producer = origin
5614 .scope(&["allowed".into()])
5615 .expect("should create limited producer");
5616
5617 assert!(limited_producer.with_root("notallowed").is_none());
5619
5620 let allowed_root = limited_producer
5622 .with_root("allowed")
5623 .expect("should create allowed root");
5624 assert_eq!(allowed_root.root().as_str(), "allowed");
5625 }
5626
5627 #[tokio::test]
5628 async fn test_wildcard_permission() {
5629 let origin = Origin::random().produce();
5630
5631 let root_producer = origin.clone();
5633
5634 let _broadcast = root_producer
5636 .create_broadcast("any/path", announce())
5637 .expect("publish allowed");
5638 let _keep2 = root_producer
5639 .create_broadcast("other/path", announce())
5640 .expect("publish allowed");
5641 settle().await;
5642
5643 let foo_producer = root_producer.with_root("foo").expect("should create any root");
5645 assert_eq!(foo_producer.root().as_str(), "foo");
5646 }
5647
5648 #[tokio::test]
5649 async fn test_consume_broadcast_with_permissions() {
5650 let origin = Origin::random().produce();
5651
5652 let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
5653 let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
5654 settle().await;
5655
5656 let limited_consumer = origin
5658 .consume()
5659 .scope(&["allowed".into()])
5660 .expect("should create limited consumer");
5661
5662 let result = limited_consumer.get_broadcast("allowed/test");
5664 assert!(result.is_some());
5665 assert!(
5666 result
5667 .unwrap()
5668 .is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
5669 );
5670
5671 assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
5673
5674 let consumer = origin.consume();
5676 assert!(consumer.get_broadcast("allowed/test").is_some());
5677 assert!(consumer.get_broadcast("notallowed/test").is_some());
5678 }
5679
5680 #[tokio::test]
5681 async fn test_nested_paths_with_permissions() {
5682 let origin = Origin::random().produce();
5683
5684 let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
5686
5687 let _broadcast = limited_producer
5689 .create_broadcast("a/b/c", announce())
5690 .expect("publish allowed");
5691 let _keep2 = limited_producer
5692 .create_broadcast("a/b/c/d", announce())
5693 .expect("publish allowed");
5694 let _keep3 = limited_producer
5695 .create_broadcast("a/b/c/d/e", announce())
5696 .expect("publish allowed");
5697 settle().await;
5698
5699 assert!(limited_producer.create_broadcast("a", announce()).is_err());
5701 assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
5702 assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
5703 }
5704
5705 #[tokio::test]
5706 async fn test_multiple_consumers_with_different_permissions() {
5707 let origin = Origin::random().produce();
5708
5709 let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
5711 let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
5712 let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
5713 settle().await;
5714
5715 let mut foo_consumer = origin
5717 .consume()
5718 .scope(&["foo".into()])
5719 .expect("should create foo consumer")
5720 .announced();
5721
5722 let mut bar_consumer = origin
5723 .consume()
5724 .scope(&["bar".into()])
5725 .expect("should create bar consumer")
5726 .announced();
5727
5728 let mut foobar_consumer = origin
5729 .consume()
5730 .scope(&["foo".into(), "bar".into()])
5731 .expect("should create foobar consumer")
5732 .announced();
5733
5734 foo_consumer.assert_next_some("foo/test");
5736 foo_consumer.assert_next_wait();
5737
5738 bar_consumer.assert_next_some("bar/test");
5739 bar_consumer.assert_next_wait();
5740
5741 foobar_consumer.assert_next_some("bar/test");
5742 foobar_consumer.assert_next_some("foo/test");
5743 foobar_consumer.assert_next_wait();
5744 }
5745
5746 #[tokio::test]
5747 async fn test_select_with_empty_prefix() {
5748 let origin = Origin::random().produce();
5749
5750 let demo_producer = origin.with_root("demo").expect("should create demo root");
5752 let limited_producer = demo_producer
5753 .scope(&["worm-node".into(), "foobar".into()])
5754 .expect("should create limited producer");
5755
5756 let _broadcast1 = limited_producer
5758 .create_broadcast("worm-node/test", announce())
5759 .expect("publish allowed");
5760 let _broadcast2 = limited_producer
5761 .create_broadcast("foobar/test", announce())
5762 .expect("publish allowed");
5763 settle().await;
5764
5765 let mut consumer = limited_producer
5767 .consume()
5768 .scope(&["".into()])
5769 .expect("should create consumer with empty prefix")
5770 .announced();
5771
5772 let a1 = consumer.try_next().expect("expected first announcement");
5774 let a2 = consumer.try_next().expect("expected second announcement");
5775 consumer.assert_next_wait();
5776
5777 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
5778 paths.sort();
5779 assert_eq!(paths, ["foobar/test", "worm-node/test"]);
5780 }
5781
5782 #[tokio::test]
5783 async fn test_select_narrowing_scope() {
5784 let origin = Origin::random().produce();
5785
5786 let demo_producer = origin.with_root("demo").expect("should create demo root");
5788 let limited_producer = demo_producer
5789 .scope(&["worm-node".into(), "foobar".into()])
5790 .expect("should create limited producer");
5791
5792 let _broadcast1 = limited_producer
5794 .create_broadcast("worm-node", announce())
5795 .expect("publish allowed");
5796 let _broadcast2 = limited_producer
5797 .create_broadcast("worm-node/foo", announce())
5798 .expect("publish allowed");
5799 let _broadcast3 = limited_producer
5800 .create_broadcast("foobar/bar", announce())
5801 .expect("publish allowed");
5802 settle().await;
5803
5804 let mut worm_consumer = limited_producer
5806 .consume()
5807 .scope(&["worm-node".into()])
5808 .expect("should create worm-node consumer")
5809 .announced();
5810
5811 worm_consumer.assert_next_some("worm-node");
5813 worm_consumer.assert_next_some("worm-node/foo");
5814 worm_consumer.assert_next_wait(); let mut foo_consumer = limited_producer
5818 .consume()
5819 .scope(&["worm-node/foo".into()])
5820 .expect("should create worm-node/foo consumer")
5821 .announced();
5822
5823 foo_consumer.assert_next_some("worm-node/foo");
5824 foo_consumer.assert_next_wait(); }
5826
5827 #[tokio::test]
5828 async fn test_select_multiple_roots_with_empty_prefix() {
5829 let origin = Origin::random().produce();
5830
5831 let limited_producer = origin
5833 .scope(&["app1".into(), "app2".into(), "shared".into()])
5834 .expect("should create limited producer");
5835
5836 let _broadcast1 = limited_producer
5838 .create_broadcast("app1/data", announce())
5839 .expect("publish allowed");
5840 let _broadcast2 = limited_producer
5841 .create_broadcast("app2/config", announce())
5842 .expect("publish allowed");
5843 let _broadcast3 = limited_producer
5844 .create_broadcast("shared/resource", announce())
5845 .expect("publish allowed");
5846 settle().await;
5847
5848 let mut consumer = limited_producer
5850 .consume()
5851 .scope(&["".into()])
5852 .expect("should create consumer with empty prefix")
5853 .announced();
5854
5855 consumer.assert_next_some("app1/data");
5857 consumer.assert_next_some("app2/config");
5858 consumer.assert_next_some("shared/resource");
5859 consumer.assert_next_wait();
5860 }
5861
5862 #[tokio::test]
5863 async fn test_publish_scope_with_empty_prefix() {
5864 let origin = Origin::random().produce();
5865
5866 let limited_producer = origin
5868 .scope(&["services/api".into(), "services/web".into()])
5869 .expect("should create limited producer");
5870
5871 let same_producer = limited_producer
5873 .scope(&["".into()])
5874 .expect("should create producer with empty prefix");
5875
5876 let _broadcast = same_producer
5878 .create_broadcast("services/api", announce())
5879 .expect("publish allowed");
5880 let _keep2 = same_producer
5881 .create_broadcast("services/web", announce())
5882 .expect("publish allowed");
5883 assert!(same_producer.create_broadcast("services/db", announce()).is_err());
5884 assert!(same_producer.create_broadcast("other", announce()).is_err());
5885 }
5886
5887 #[tokio::test]
5888 async fn test_select_narrowing_to_deeper_path() {
5889 let origin = Origin::random().produce();
5890
5891 let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
5893
5894 let _broadcast1 = limited_producer
5896 .create_broadcast("org/team1/project1", announce())
5897 .expect("publish allowed");
5898 let _broadcast2 = limited_producer
5899 .create_broadcast("org/team1/project2", announce())
5900 .expect("publish allowed");
5901 let _broadcast3 = limited_producer
5902 .create_broadcast("org/team2/project1", announce())
5903 .expect("publish allowed");
5904 settle().await;
5905
5906 let mut team2_consumer = limited_producer
5908 .consume()
5909 .scope(&["org/team2".into()])
5910 .expect("should create team2 consumer")
5911 .announced();
5912
5913 team2_consumer.assert_next_some("org/team2/project1");
5914 team2_consumer.assert_next_wait(); let mut project1_consumer = limited_producer
5918 .consume()
5919 .scope(&["org/team1/project1".into()])
5920 .expect("should create project1 consumer")
5921 .announced();
5922
5923 project1_consumer.assert_next_some("org/team1/project1");
5925 project1_consumer.assert_next_wait();
5926 }
5927
5928 #[tokio::test]
5929 async fn test_select_with_non_matching_prefix() {
5930 let origin = Origin::random().produce();
5931
5932 let limited_producer = origin
5934 .scope(&["allowed/path".into()])
5935 .expect("should create limited producer");
5936
5937 assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
5939
5940 assert!(limited_producer.scope(&["other/path".into()]).is_none());
5942 }
5943
5944 #[tokio::test]
5947 async fn test_with_root_trailing_slash_consumer() {
5948 let origin = Origin::random().produce();
5949
5950 let prefix = "some_prefix/".to_string();
5952 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5953
5954 let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5955 settle().await;
5956 consumer.assert_next_some("test");
5957 }
5958
5959 #[tokio::test]
5961 async fn test_with_root_trailing_slash_producer() {
5962 let origin = Origin::random().produce();
5963
5964 let prefix = "some_prefix/".to_string();
5966 let rooted = origin.with_root(prefix).unwrap();
5967
5968 let _b = rooted.create_broadcast("test", announce()).unwrap();
5969 settle().await;
5970
5971 let mut consumer = rooted.consume().announced();
5972 consumer.assert_next_some("test");
5973 }
5974
5975 #[tokio::test]
5977 async fn test_with_root_trailing_slash_unannounce() {
5978 tokio::time::pause();
5979
5980 let origin = Origin::random().produce();
5981
5982 let prefix = "some_prefix/".to_string();
5983 let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
5984
5985 let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
5986 settle().await;
5987 consumer.assert_next_some("test");
5988
5989 b.finish();
5991 settle().await;
5992
5993 consumer.assert_next_none("test");
5995 }
5996
5997 #[tokio::test]
5998 async fn test_select_maintains_access_with_wider_prefix() {
5999 let origin = Origin::random().produce();
6000
6001 let demo_producer = origin.with_root("demo").expect("should create demo root");
6003 let user_producer = demo_producer
6004 .scope(&["worm-node".into(), "foobar".into()])
6005 .expect("should create user producer");
6006
6007 let _broadcast1 = user_producer
6009 .create_broadcast("worm-node/data", announce())
6010 .expect("publish allowed");
6011 let _broadcast2 = user_producer
6012 .create_broadcast("foobar", announce())
6013 .expect("publish allowed");
6014 settle().await;
6015
6016 let mut consumer = user_producer
6018 .consume()
6019 .scope(&["".into()])
6020 .expect("scope with empty prefix should not fail when user has specific permissions")
6021 .announced();
6022
6023 let a1 = consumer.try_next().expect("expected first announcement");
6025 let a2 = consumer.try_next().expect("expected second announcement");
6026 consumer.assert_next_wait();
6027
6028 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6029 paths.sort();
6030 assert_eq!(paths, ["foobar", "worm-node/data"]);
6031
6032 let mut narrow_consumer = user_producer
6034 .consume()
6035 .scope(&["worm-node".into()])
6036 .expect("should be able to narrow scope to worm-node")
6037 .announced();
6038
6039 narrow_consumer.assert_next_some("worm-node/data");
6040 narrow_consumer.assert_next_wait(); }
6042
6043 #[tokio::test]
6044 async fn test_duplicate_prefixes_deduped() {
6045 let origin = Origin::random().produce();
6046
6047 let producer = origin
6049 .scope(&["demo".into(), "demo".into()])
6050 .expect("should create producer");
6051
6052 let _broadcast = producer
6053 .create_broadcast("demo/stream", announce())
6054 .expect("publish allowed");
6055 settle().await;
6056
6057 let mut consumer = producer.consume().announced();
6058 consumer.assert_next_some("demo/stream");
6059 consumer.assert_next_wait();
6060 }
6061
6062 #[tokio::test]
6063 async fn test_overlapping_prefixes_deduped() {
6064 let origin = Origin::random().produce();
6065
6066 let producer = origin
6068 .scope(&["demo".into(), "demo/foo".into()])
6069 .expect("should create producer");
6070
6071 let _broadcast = producer
6073 .create_broadcast("demo/bar/stream", announce())
6074 .expect("publish allowed");
6075 settle().await;
6076
6077 let mut consumer = producer.consume().announced();
6078 consumer.assert_next_some("demo/bar/stream");
6079 consumer.assert_next_wait();
6080 }
6081
6082 #[tokio::test]
6083 async fn test_overlapping_prefixes_no_duplicate_announcements() {
6084 let origin = Origin::random().produce();
6085
6086 let producer = origin
6088 .scope(&["demo".into(), "demo/foo".into()])
6089 .expect("should create producer");
6090
6091 let _broadcast = producer
6092 .create_broadcast("demo/foo/stream", announce())
6093 .expect("publish allowed");
6094 settle().await;
6095
6096 let mut consumer = producer.consume().announced();
6097 consumer.assert_next_some("demo/foo/stream");
6099 consumer.assert_next_wait();
6100 }
6101
6102 #[tokio::test]
6103 async fn test_allowed_returns_deduped_prefixes() {
6104 let origin = Origin::random().produce();
6105
6106 let producer = origin
6107 .scope(&["demo".into(), "demo/foo".into(), "anon".into()])
6108 .expect("should create producer");
6109
6110 let allowed: Vec<_> = producer.allowed().collect();
6111 assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
6112 }
6113
6114 #[tokio::test]
6115 async fn test_announced_broadcast_already_announced() {
6116 let origin = Origin::random().produce();
6117
6118 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6119 settle().await;
6120
6121 let consumer = origin.consume();
6122 let result = consumer.announced_broadcast("test").await.expect("should find it");
6123 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6124 }
6125
6126 #[tokio::test]
6127 async fn test_announced_broadcast_delayed() {
6128 tokio::time::pause();
6129
6130 let origin = Origin::random().produce();
6131
6132 let consumer = origin.consume();
6133
6134 let wait = tokio::spawn({
6136 let consumer = consumer.clone();
6137 async move { consumer.announced_broadcast("test").await }
6138 });
6139
6140 tokio::task::yield_now().await;
6142
6143 let _broadcast = origin.create_broadcast("test", announce()).unwrap();
6144 settle().await;
6145
6146 let result = wait.await.unwrap().expect("should find it");
6147 assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
6148 }
6149
6150 #[tokio::test]
6151 async fn test_announced_broadcast_ignores_unrelated_paths() {
6152 tokio::time::pause();
6153
6154 let origin = Origin::random().produce();
6155
6156 let consumer = origin.consume();
6157
6158 let wait = tokio::spawn({
6159 let consumer = consumer.clone();
6160 async move { consumer.announced_broadcast("target").await }
6161 });
6162
6163 tokio::task::yield_now().await;
6164
6165 let _other = origin.create_broadcast("other", announce()).unwrap();
6167 settle().await;
6168 tokio::task::yield_now().await;
6169 assert!(!wait.is_finished(), "must not resolve on unrelated path");
6170
6171 let _target = origin.create_broadcast("target", announce()).unwrap();
6172 settle().await;
6173 let result = wait.await.unwrap().expect("should find target");
6174 assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
6175 }
6176
6177 #[tokio::test]
6178 async fn test_announced_broadcast_skips_nested_paths() {
6179 tokio::time::pause();
6180
6181 let origin = Origin::random().produce();
6182
6183 let consumer = origin.consume();
6184
6185 let wait = tokio::spawn({
6186 let consumer = consumer.clone();
6187 async move { consumer.announced_broadcast("foo").await }
6188 });
6189
6190 tokio::task::yield_now().await;
6191
6192 let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
6194 settle().await;
6195 tokio::task::yield_now().await;
6196 assert!(!wait.is_finished(), "must not resolve on a nested path");
6197
6198 let _exact = origin.create_broadcast("foo", announce()).unwrap();
6199 settle().await;
6200 let result = wait.await.unwrap().expect("should find foo exactly");
6201 assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
6202 }
6203
6204 #[tokio::test]
6205 async fn test_announced_broadcast_disallowed() {
6206 let origin = Origin::random().produce();
6207 let limited = origin
6208 .consume()
6209 .scope(&["allowed".into()])
6210 .expect("should create limited");
6211
6212 assert!(limited.announced_broadcast("notallowed").await.is_none());
6214 }
6215
6216 #[tokio::test]
6217 async fn test_announced_broadcast_scope_too_narrow() {
6218 let origin = Origin::random().produce();
6221 let limited = origin
6222 .consume()
6223 .scope(&["foo/specific".into()])
6224 .expect("should create limited");
6225
6226 let result = limited
6228 .announced_broadcast("foo")
6229 .now_or_never()
6230 .expect("must not block");
6231 assert!(result.is_none());
6232 }
6233
6234 #[tokio::test]
6238 async fn test_coalesce_announce_then_unannounce() {
6239 tokio::time::pause();
6241
6242 let origin = Origin::random().produce();
6243 let mut announced = origin.consume().announced();
6244
6245 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6246 settle().await;
6247 broadcast.finish();
6248
6249 settle().await;
6250
6251 announced.assert_next_wait();
6252 }
6253
6254 #[tokio::test]
6255 async fn test_coalesce_announce_unannounce_announce() {
6256 tokio::time::pause();
6259
6260 let origin = Origin::random().produce();
6261 let mut announced = origin.consume().announced();
6262
6263 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6264 settle().await;
6265 broadcast1.finish();
6266 settle().await;
6267 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6268 settle().await;
6269
6270 announced.assert_next_some("test");
6271 announced.assert_next_wait();
6272 }
6273
6274 #[tokio::test]
6275 async fn test_coalesce_unannounce_announce_preserved() {
6276 tokio::time::pause();
6279
6280 let origin = Origin::random().produce();
6281 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6282 settle().await;
6283
6284 let mut announced = origin.consume().announced();
6285 announced.assert_next_some("test");
6286
6287 broadcast1.finish();
6289 settle().await;
6290
6291 let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6292 settle().await;
6293
6294 announced.assert_next_none("test");
6296 announced.assert_next_some("test");
6297 announced.assert_next_wait();
6298 }
6299
6300 #[tokio::test]
6301 async fn test_coalesce_unannounce_announce_unannounce() {
6302 tokio::time::pause();
6305
6306 let origin = Origin::random().produce();
6307 let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
6308 settle().await;
6309
6310 let mut announced = origin.consume().announced();
6311 announced.assert_next_some("test");
6312
6313 broadcast1.finish();
6314 settle().await;
6315
6316 let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
6317 settle().await;
6318 broadcast2.finish();
6319 settle().await;
6320
6321 announced.assert_next_none("test");
6322 announced.assert_next_wait();
6323 }
6324
6325 #[tokio::test]
6326 async fn test_coalesce_churn_bounded() {
6327 tokio::time::pause();
6332
6333 let origin = Origin::random().produce();
6334 let mut announced = origin.consume().announced();
6335
6336 for _ in 0..1000 {
6337 let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
6338 settle().await;
6339 broadcast.finish();
6340 }
6341 settle().await;
6342
6343 let mut collected = Vec::new();
6344 while let Some(update) = announced.try_next() {
6345 collected.push(update);
6346 }
6347 assert!(
6348 collected.len() <= 1,
6349 "expected at most one pending update, got {}",
6350 collected.len()
6351 );
6352 assert!(
6353 collected.iter().all(|a| a.path == Path::new("test")),
6354 "unexpected path in pending updates",
6355 );
6356 }
6357
6358 #[tokio::test]
6362 async fn test_consumer_clone_is_side_effect_free() {
6363 let origin = Origin::random().produce();
6364
6365 let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
6366 let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
6367 settle().await;
6368
6369 let consumer = origin.consume();
6370 let mut announced = consumer.announced();
6371
6372 for _ in 0..16 {
6375 let cloned = consumer.clone();
6376 assert!(cloned.get_broadcast("test1").is_some());
6377 assert!(cloned.get_broadcast("test2").is_some());
6378 }
6379
6380 let a1 = announced.try_next().expect("first announcement");
6383 let a2 = announced.try_next().expect("second announcement");
6384 announced.assert_next_wait();
6385
6386 let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
6387 paths.sort();
6388 assert_eq!(paths, ["test1", "test2"]);
6389
6390 let mut fresh = consumer.announced();
6392 let b1 = fresh.try_next().expect("backlog: first");
6393 let b2 = fresh.try_next().expect("backlog: second");
6394 fresh.assert_next_wait();
6395
6396 let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
6397 paths.sort();
6398 assert_eq!(paths, ["test1", "test2"]);
6399 }
6400
6401 #[tokio::test]
6403 async fn dynamic_request_unroutable_without_handler() {
6404 let origin = Origin::random().produce();
6405 let consumer = origin.consume();
6406 assert!(matches!(
6407 consumer.request_broadcast("missing").await,
6408 Err(Error::Unroutable)
6409 ));
6410 }
6411
6412 #[tokio::test(start_paused = true)]
6415 async fn dynamic_request_served_not_announced() {
6416 let origin = Origin::random().produce();
6417 let mut dynamic = origin.dynamic();
6418 let consumer = origin.consume();
6419
6420 let mut announced = origin.consume().announced();
6422 announced.assert_next_wait();
6423
6424 let served = broadcast::Info::new().produce();
6425 let request_fut = consumer.request_broadcast("fallback");
6428
6429 let mut served_dynamic = served.dynamic();
6431
6432 let request = dynamic.requested_broadcast().await.unwrap();
6433 assert_eq!(request.path(), &Path::new("fallback"));
6434 request.accept(&served);
6435
6436 let broadcast = request_fut.await.unwrap();
6437 assert!(broadcast.is_clone(&served.consume()));
6438
6439 let track_fut = broadcast.track("video").unwrap().subscribe(None);
6441 let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
6442 let mut track = track_fut.await.unwrap();
6443 producer.append_group().unwrap();
6444 track.assert_group();
6445
6446 announced.assert_next_wait();
6448 }
6449
6450 #[tokio::test(start_paused = true)]
6452 async fn dynamic_request_coalesces() {
6453 let origin = Origin::random().produce();
6454 let mut dynamic = origin.dynamic();
6455 let consumer = origin.consume();
6456
6457 let f1 = consumer.request_broadcast("dup");
6459 let f2 = consumer.request_broadcast("dup");
6460
6461 let request = dynamic.requested_broadcast().await.unwrap();
6463 assert_eq!(request.path(), &Path::new("dup"));
6464 assert!(
6465 dynamic.requested_broadcast().now_or_never().is_none(),
6466 "a coalesced request must not be served twice"
6467 );
6468
6469 let served = broadcast::Info::new().produce();
6471 request.accept(&served);
6472 assert!(f1.await.unwrap().is_clone(&served.consume()));
6473 assert!(f2.await.unwrap().is_clone(&served.consume()));
6474 }
6475
6476 #[tokio::test(start_paused = true)]
6479 async fn dynamic_request_dedups_served() {
6480 let origin = Origin::random().produce();
6481 let mut dynamic = origin.dynamic();
6482 let consumer = origin.consume();
6483
6484 let request_fut = consumer.request_broadcast("fallback");
6485 let request = dynamic.requested_broadcast().await.unwrap();
6486 let served = broadcast::Info::new().produce();
6487 request.accept(&served);
6488 let first = request_fut.await.unwrap();
6489 assert!(first.is_clone(&served.consume()));
6490
6491 let second = consumer.request_broadcast("fallback").await.unwrap();
6493 assert!(second.is_clone(&served.consume()));
6494
6495 assert!(
6497 dynamic.requested_broadcast().now_or_never().is_none(),
6498 "a still-live served broadcast must not be re-requested from the handler"
6499 );
6500 }
6501
6502 #[tokio::test(start_paused = true)]
6504 async fn dynamic_request_reserves_after_close() {
6505 let origin = Origin::random().produce();
6506 let mut dynamic = origin.dynamic();
6507 let consumer = origin.consume();
6508
6509 let request_fut = consumer.request_broadcast("fallback");
6510 let request = dynamic.requested_broadcast().await.unwrap();
6511 let served = broadcast::Info::new().produce();
6512 request.accept(&served);
6513 request_fut.await.unwrap();
6514
6515 drop(served);
6517
6518 let request_fut = consumer.request_broadcast("fallback");
6520 let request = dynamic.requested_broadcast().await.unwrap();
6521 assert_eq!(request.path(), &Path::new("fallback"));
6522 let served = broadcast::Info::new().produce();
6523 request.accept(&served);
6524 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6525 }
6526
6527 #[tokio::test(start_paused = true)]
6530 async fn dynamic_request_served_cache_bounded() {
6531 let origin = Origin::random().produce();
6532 let mut dynamic = origin.dynamic();
6533 let consumer = origin.consume();
6534
6535 for i in 0..100 {
6536 let path = format!("one-shot/{i}");
6537 let request_fut = consumer.request_broadcast(&path);
6538 let request = dynamic.requested_broadcast().await.unwrap();
6539 let served = broadcast::Info::new().produce();
6540 request.accept(&served);
6541 request_fut.await.unwrap();
6542 drop(served);
6544 }
6545
6546 assert!(
6549 origin.dynamic.read().served.len() <= 4,
6550 "stale served entries must be reclaimed, not accumulate per distinct path: {}",
6551 origin.dynamic.read().served.len()
6552 );
6553 }
6554
6555 #[tokio::test(start_paused = true)]
6558 async fn dynamic_request_coalesces_after_handoff() {
6559 let origin = Origin::random().produce();
6560 let mut dynamic = origin.dynamic();
6561 let consumer = origin.consume();
6562
6563 let f1 = consumer.request_broadcast("fallback");
6564 let request = dynamic.requested_broadcast().await.unwrap();
6566
6567 let f2 = consumer.request_broadcast("fallback");
6569 assert!(
6570 dynamic.requested_broadcast().now_or_never().is_none(),
6571 "a repeat request during hand-off must coalesce, not re-queue"
6572 );
6573
6574 let served = broadcast::Info::new().produce();
6576 request.accept(&served);
6577 assert!(f1.await.unwrap().is_clone(&served.consume()));
6578 assert!(f2.await.unwrap().is_clone(&served.consume()));
6579 }
6580
6581 #[tokio::test(start_paused = true)]
6583 async fn dynamic_request_dropped_after_handoff() {
6584 let origin = Origin::random().produce();
6585 let mut dynamic = origin.dynamic();
6586 let consumer = origin.consume();
6587
6588 let f1 = consumer.request_broadcast("fallback");
6589 let request = dynamic.requested_broadcast().await.unwrap();
6590 let f2 = consumer.request_broadcast("fallback");
6591
6592 drop(request);
6594 assert!(matches!(f1.await, Err(Error::Unroutable)));
6595 assert!(matches!(f2.await, Err(Error::Unroutable)));
6596 }
6597
6598 #[tokio::test(start_paused = true)]
6600 async fn dynamic_request_rejected() {
6601 let origin = Origin::random().produce();
6602 let mut dynamic = origin.dynamic();
6603 let consumer = origin.consume();
6604
6605 let request_fut = consumer.request_broadcast("fallback");
6606
6607 let request = dynamic.requested_broadcast().await.unwrap();
6608 request.reject(Error::Cancel);
6609
6610 assert!(matches!(request_fut.await, Err(Error::Cancel)));
6611 }
6612
6613 #[tokio::test(start_paused = true)]
6617 async fn dynamic_request_rerequest_after_reject() {
6618 let origin = Origin::random().produce();
6619 let mut dynamic = origin.dynamic();
6620 let consumer = origin.consume();
6621
6622 let f1 = consumer.request_broadcast("fallback");
6623 dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
6624 assert!(matches!(f1.await, Err(Error::Unroutable)));
6625
6626 let served = broadcast::Info::new().produce();
6627 let f2 = consumer.request_broadcast("fallback");
6629 let request = dynamic.requested_broadcast().await.unwrap();
6630 assert_eq!(request.path(), &Path::new("fallback"));
6631 request.accept(&served);
6632 assert!(f2.await.unwrap().is_clone(&served.consume()));
6633 }
6634
6635 #[tokio::test(start_paused = true)]
6638 async fn dynamic_request_handler_dropped() {
6639 let origin = Origin::random().produce();
6640 let dynamic = origin.dynamic();
6641 let consumer = origin.consume();
6642
6643 let request_fut = consumer.request_broadcast("fallback");
6644 drop(dynamic);
6645 assert!(matches!(request_fut.await, Err(Error::Unroutable)));
6646
6647 assert!(matches!(
6649 consumer.request_broadcast("again").await,
6650 Err(Error::Unroutable)
6651 ));
6652 }
6653
6654 #[tokio::test(start_paused = true)]
6658 async fn dynamic_request_accept_after_handler_dropped() {
6659 let origin = Origin::random().produce();
6660 let mut dynamic = origin.dynamic();
6661 let consumer = origin.consume();
6662
6663 let request_fut = consumer.request_broadcast("fallback");
6664
6665 let request = dynamic.requested_broadcast().await.unwrap();
6667 drop(dynamic);
6668
6669 let served = broadcast::Info::new().produce();
6670 request.accept(&served);
6672 assert!(request_fut.await.unwrap().is_clone(&served.consume()));
6673 }
6674
6675 #[tokio::test(start_paused = true)]
6677 async fn dynamic_request_prefers_announced() {
6678 let origin = Origin::random().produce();
6679 let mut dynamic = origin.dynamic();
6680 let consumer = origin.consume();
6681
6682 let _broadcast = origin.create_broadcast("live", announce()).unwrap();
6683 settle().await;
6684
6685 let got = consumer.request_broadcast("live").await.unwrap();
6686 assert!(
6687 got.is_clone(&consumer.get_broadcast("live").unwrap()),
6688 "should return the published broadcast"
6689 );
6690 assert!(
6691 dynamic.requested_broadcast().now_or_never().is_none(),
6692 "a published path must not queue a fallback request"
6693 );
6694 }
6695
6696 #[tokio::test(start_paused = true)]
6698 async fn dynamic_clone_keeps_alive() {
6699 let origin = Origin::random().produce();
6700 let dynamic = origin.dynamic();
6701 let consumer = origin.consume();
6702
6703 drop(dynamic.clone());
6704
6705 let request_fut = consumer.request_broadcast("fallback");
6708 assert!(
6709 request_fut.now_or_never().is_none(),
6710 "request should stay pending until served"
6711 );
6712 }
6713}