1use crate::{broadcast, cache, group, stats, track};
2use kio::Pollable;
3use std::{
4 cmp::Reverse,
5 collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque},
6 fmt,
7 sync::Arc,
8 sync::atomic::{AtomicU64, Ordering},
9 task::{Poll, ready},
10 time::Duration,
11};
12
13use rand::RngExt;
14
15use super::{
16 Requests, WeakCache, WeakEntry,
17 front::{Action, Candidate, Event, Front, Pin, Refusal},
18};
19use crate::{
20 AsPath, Error, InvalidPattern, Path, PathOwned, Pattern, Patterns,
21 coding::{BoundsExceeded, Decode, DecodeError, Encode, EncodeError},
22 path::Segment,
23 runtime::{Instant, Timers},
24 time::Clock,
25 util::{Keepalive, TaskSet, Tasks, TasksWeak},
26};
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
38pub struct Hop {
39 id: u64,
41}
42
43impl Hop {
44 pub const UNKNOWN: Self = Self { id: 0 };
51
52 pub fn new(id: u64) -> Result<Self, InvalidHop> {
58 if id == 0 || id >= 1u64 << 62 {
59 return Err(InvalidHop::Range);
60 }
61 Ok(Self { id })
62 }
63
64 pub fn random() -> Self {
71 let mut rng = rand::rng();
72 let id = rng.random_range(1..(1u64 << 53));
73 Self { id }
74 }
75
76 pub fn id(self) -> u64 {
78 self.id
79 }
80
81 pub(crate) fn from_wire(id: u64) -> Result<Self, DecodeError> {
83 if id >= 1u64 << 62 {
84 return Err(DecodeError::InvalidValue);
85 }
86 Ok(Self { id })
87 }
88}
89
90#[derive(Clone, Debug)]
96#[non_exhaustive]
97pub struct Config {
98 pub hop: Hop,
101
102 pub pool: cache::Pool,
109
110 pub cache_duration: Duration,
118
119 pub default_max_age: Duration,
129
130 pub update_hold: Duration,
141}
142
143pub const DEFAULT_UPDATE_HOLD: Duration = Duration::from_millis(300);
146
147impl Default for Config {
148 fn default() -> Self {
150 let pool = cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY));
151 Self {
152 hop: Hop::random(),
153 pool,
154 cache_duration: Duration::MAX,
155 default_max_age: track::DEFAULT_MAX_AGE,
156 update_hold: DEFAULT_UPDATE_HOLD,
157 }
158 }
159}
160
161impl Config {
162 pub fn new(hop: Hop) -> Self {
164 Self { hop, ..Self::default() }
165 }
166}
167
168impl From<Hop> for Config {
169 fn from(hop: Hop) -> Self {
171 Self::new(hop)
172 }
173}
174
175impl TryFrom<u64> for Hop {
176 type Error = InvalidHop;
177
178 fn try_from(id: u64) -> Result<Self, Self::Error> {
179 Self::new(id)
180 }
181}
182
183impl fmt::Display for Hop {
184 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
185 self.id.fmt(f)
186 }
187}
188
189impl<V: Copy> Encode<V> for Hop
190where
191 u64: Encode<V>,
192{
193 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
194 self.id.encode(w, version)
195 }
196}
197
198impl<V: Copy> Decode<V> for Hop
199where
200 u64: Decode<V>,
201{
202 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
203 Self::from_wire(u64::decode(r, version)?)
204 }
205}
206
207pub(crate) const MAX_HOPS: usize = 32;
213
214#[derive(Debug, Clone, Default, PartialEq, Eq)]
222pub struct Hops(Vec<Hop>);
223
224#[derive(Debug, Clone, Copy, PartialEq, Eq)]
226#[non_exhaustive]
227pub enum InvalidHop {
228 Range,
231
232 TooMany,
235
236 Duplicate,
240}
241
242impl fmt::Display for InvalidHop {
243 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
244 match self {
245 Self::Range => write!(f, "local hop id must be non-zero and below 2^62"),
246 Self::TooMany => write!(f, "too many hops (max {MAX_HOPS})"),
247 Self::Duplicate => write!(f, "hop already in the chain"),
248 }
249 }
250}
251
252impl std::error::Error for InvalidHop {}
253
254impl From<InvalidHop> for DecodeError {
255 fn from(err: InvalidHop) -> Self {
256 match err {
257 InvalidHop::TooMany => DecodeError::BoundsExceeded,
258 InvalidHop::Range | InvalidHop::Duplicate => DecodeError::InvalidValue,
259 }
260 }
261}
262
263impl Hops {
264 pub fn new() -> Self {
266 Self(Vec::new())
267 }
268
269 pub fn push(&mut self, hop: Hop) -> Result<(), InvalidHop> {
275 if self.0.len() >= MAX_HOPS {
276 return Err(InvalidHop::TooMany);
277 }
278 if hop != Hop::UNKNOWN && self.0.contains(&hop) {
279 return Err(InvalidHop::Duplicate);
280 }
281 self.0.push(hop);
282 Ok(())
283 }
284
285 pub(crate) fn stamp(&mut self, stamp: Hop) -> Result<(), InvalidHop> {
294 match self.0.first() {
295 None => {
296 self.push(stamp)?;
297 self.push(Hop::UNKNOWN)
298 }
299 Some(first) if *first == Hop::UNKNOWN => {
300 if self.0.len() >= MAX_HOPS {
301 return Err(InvalidHop::TooMany);
302 }
303 if self.0.contains(&stamp) {
304 return Err(InvalidHop::Duplicate);
305 }
306 self.0.insert(0, stamp);
307 Ok(())
308 }
309 Some(_) => Ok(()),
310 }
311 }
312
313 pub fn contains(&self, hop: &Hop) -> bool {
315 self.0.contains(hop)
316 }
317
318 pub fn len(&self) -> usize {
320 self.0.len()
321 }
322
323 pub fn is_empty(&self) -> bool {
325 self.0.is_empty()
326 }
327
328 pub fn iter(&self) -> std::slice::Iter<'_, Hop> {
330 self.0.iter()
331 }
332
333 pub fn as_slice(&self) -> &[Hop] {
335 &self.0
336 }
337}
338
339impl TryFrom<Vec<Hop>> for Hops {
340 type Error = InvalidHop;
341
342 fn try_from(v: Vec<Hop>) -> Result<Self, Self::Error> {
343 if v.len() > MAX_HOPS {
344 return Err(InvalidHop::TooMany);
345 }
346 for (i, hop) in v.iter().enumerate() {
348 if *hop != Hop::UNKNOWN && v[i + 1..].contains(hop) {
349 return Err(InvalidHop::Duplicate);
350 }
351 }
352 Ok(Self(v))
353 }
354}
355
356impl<'a> IntoIterator for &'a Hops {
357 type Item = &'a Hop;
358 type IntoIter = std::slice::Iter<'a, Hop>;
359
360 fn into_iter(self) -> Self::IntoIter {
361 self.iter()
362 }
363}
364
365impl<V: Copy> Encode<V> for Hops
366where
367 u64: Encode<V>,
368 Hop: Encode<V>,
369{
370 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
371 (self.0.len() as u64).encode(w, version)?;
372 for origin in &self.0 {
373 origin.encode(w, version)?;
374 }
375 Ok(())
376 }
377}
378
379impl<V: Copy> Decode<V> for Hops
380where
381 u64: Decode<V>,
382 Hop: Decode<V>,
383{
384 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
385 let count = u64::decode(r, version)? as usize;
386 if count > MAX_HOPS {
387 return Err(DecodeError::BoundsExceeded);
388 }
389 let mut list = Self(Vec::with_capacity(count));
392 for _ in 0..count {
393 list.push(Hop::decode(r, version)?)?;
394 }
395 Ok(list)
396 }
397}
398
399const MAX_COST: u64 = (1 << 62) - 1;
407
408#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)]
417pub struct Cost {
418 pub warm: u64,
427
428 pub cold: u64,
435}
436
437impl Cost {
438 pub const fn new(cost: u64) -> Self {
441 Self { warm: cost, cold: cost }
442 }
443
444 pub const MAX: Self = Self::new(MAX_COST);
451
452 pub const DRAIN: Self = Self::MAX;
455
456 pub(crate) const UNKNOWN: Self = Self {
460 warm: 0,
461 cold: MAX_COST,
462 };
463
464 pub(crate) fn charged(self, link_cost: u64) -> Self {
467 Self {
468 warm: self.warm.saturating_add(link_cost).min(MAX_COST),
469 cold: self.cold.saturating_add(link_cost).min(MAX_COST),
470 }
471 }
472
473 pub(crate) fn clamped(self) -> Self {
476 Self {
477 warm: self.warm.min(MAX_COST),
478 cold: self.cold.min(MAX_COST),
479 }
480 }
481}
482
483impl From<u64> for Cost {
484 fn from(cost: u64) -> Self {
485 Self::new(cost)
486 }
487}
488
489#[derive(Clone, Debug, PartialEq, Eq)]
500#[non_exhaustive]
501pub struct Route {
502 pub hops: Hops,
508
509 pub cost: Cost,
514
515 pub(crate) via: Hop,
521
522 pub(crate) source: Source,
524}
525
526impl Default for Route {
527 fn default() -> Self {
528 Self {
529 hops: Hops::new(),
530 cost: Cost::default(),
531 via: Hop::UNKNOWN,
532 source: Source::Local,
533 }
534 }
535}
536
537#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash)]
543pub enum Source {
544 #[default]
546 Local,
547 Peer(Hop),
550}
551
552impl Route {
553 pub fn with_hops(mut self, hops: Hops) -> Self {
555 self.hops = hops;
556 self
557 }
558
559 pub fn with_cost(mut self, cost: impl Into<Cost>) -> Self {
564 self.cost = cost.into();
565 self
566 }
567
568 pub(crate) fn with_via(mut self, via: Hop) -> Self {
573 self.via = via;
574 self
575 }
576
577 pub fn is_anonymous(&self) -> bool {
585 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
586 }
587
588 pub fn source(&self) -> Source {
593 self.source
594 }
595}
596
597static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
598
599#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
600struct ConsumerId(u64);
601
602impl ConsumerId {
603 fn new() -> Self {
604 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
605 }
606}
607
608fn fnv_key(name: &str, origins: impl IntoIterator<Item = Hop>) -> u64 {
617 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
619
620 let mut hash = SEED;
621 for &byte in name.as_bytes() {
622 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
623 }
624 for origin in origins {
625 for &byte in &origin.id().to_le_bytes() {
626 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
627 }
628 }
629
630 hash
631}
632
633fn route_order(prefix: &Path, entry: &RouteEntry) -> (bool, Cost, bool, usize, u64, Reverse<u64>) {
641 (
642 entry.is_anonymous(),
643 entry.cost,
644 !entry.local,
645 entry.hops.len(),
646 fnv_key(prefix.as_str(), entry.hops.iter().copied()),
647 Reverse(entry.id),
648 )
649}
650
651type RouteMeta = (Hops, Cost, Source);
653
654#[derive(Clone, Copy)]
661enum Held {
662 Anchoring(Instant),
663 Until(Instant),
664}
665
666type AnnounceMeta = (RouteMeta, Option<Vec<Pattern>>);
673
674enum PendingUpdate {
675 Announce(AnnounceMeta),
676 Unannounce(AnnounceMeta),
677 UnannounceAnnounce { old: AnnounceMeta, new: AnnounceMeta },
678}
679
680#[derive(Default)]
685struct OriginConsumerState {
686 pending: BTreeMap<PathOwned, PendingUpdate>,
687 delivered: BTreeSet<PathOwned>,
692 ended: bool,
695}
696
697impl OriginConsumerState {
698 fn apply_announce(&mut self, prefix: PathOwned, meta: RouteMeta, captures: Option<Vec<Pattern>>) {
699 let meta = (meta, captures);
700 let new = match self.pending.remove(&prefix) {
701 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(meta),
703 Some(PendingUpdate::Unannounce(old) | PendingUpdate::UnannounceAnnounce { old, .. }) => {
705 PendingUpdate::UnannounceAnnounce { old, new: meta }
706 }
707 };
708 self.pending.insert(prefix, new);
709 }
710
711 fn apply_unannounce(&mut self, prefix: PathOwned, last: RouteMeta, captures: Option<Vec<Pattern>>) {
712 let last = (last, captures);
713 match self.pending.remove(&prefix) {
714 Some(PendingUpdate::Announce(_)) if !self.delivered.contains(&prefix) => {}
717 None | Some(PendingUpdate::Announce(_) | PendingUpdate::Unannounce(_)) => {
720 self.pending.insert(prefix, PendingUpdate::Unannounce(last));
721 }
722 Some(PendingUpdate::UnannounceAnnounce { old, .. }) => {
725 self.pending.insert(prefix, PendingUpdate::Unannounce(old));
726 }
727 }
728 }
729
730 fn is_update(&self, prefix: &PathOwned, update: &PendingUpdate) -> bool {
733 matches!(update, PendingUpdate::Announce(_)) && self.delivered.contains(prefix)
734 }
735
736 fn scan(
739 &self,
740 held: &mut HashMap<PathOwned, Held>,
741 now: Option<Instant>,
742 hold: Duration,
743 ) -> Result<PathOwned, Option<Instant>> {
744 let mut wake = None;
745 for (prefix, update) in &self.pending {
746 let holds = !hold.is_zero() && !self.ended && self.is_update(prefix, update);
747 let Some(now) = now.filter(|_| holds) else {
748 return Ok(prefix.clone());
749 };
750 let until = match held.get(prefix).copied() {
753 None => {
754 held.insert(prefix.clone(), Held::Anchoring(now));
755 now + Duration::from_nanos(1)
756 }
757 Some(Held::Anchoring(seen)) if now <= seen => seen + Duration::from_nanos(1),
758 Some(Held::Anchoring(_)) => {
759 held.insert(prefix.clone(), Held::Until(now + hold));
760 now + hold
761 }
762 Some(Held::Until(until)) if until <= now => return Ok(prefix.clone()),
763 Some(Held::Until(until)) => until,
764 };
765 wake = Some(wake.map_or(until, |w: Instant| w.min(until)));
766 }
767 Err(wake)
768 }
769
770 fn take_prefix(&mut self, prefix: PathOwned) -> Option<AnnounceUpdate> {
772 let ((meta, captures), kind) = match self.pending.remove(&prefix).unwrap() {
773 PendingUpdate::Announce(meta) => {
774 let kind = match self.delivered.insert(prefix.clone()) {
776 true => AnnounceKind::Announced,
777 false => AnnounceKind::Updated,
778 };
779 (meta, kind)
780 }
781 PendingUpdate::Unannounce(meta) => {
782 self.delivered.remove(&prefix);
783 (meta, AnnounceKind::Retracted)
784 }
785 PendingUpdate::UnannounceAnnounce { old, new } => {
786 self.delivered.remove(&prefix);
789 self.pending.insert(prefix.clone(), PendingUpdate::Announce(new));
790 (old, AnnounceKind::Retracted)
791 }
792 };
793 Some(AnnounceUpdate {
794 prefix,
795 captures,
796 route: Route {
797 hops: meta.0,
798 cost: meta.1,
799 via: Hop::UNKNOWN,
800 source: meta.2,
801 },
802 kind,
803 })
804 }
805}
806
807struct RouteEntry {
809 id: u64,
810 prefix: PathOwned,
811 scope: Patterns,
814 hops: Hops,
815 cost: Cost,
816 via: Hop,
820 local: bool,
823 peer: bool,
826 server: Option<kio::Shared<ServeState>>,
830 source: Option<broadcast::Consumer>,
834 advertised: bool,
838 stale: bool,
842 claim: Pattern,
848}
849
850impl RouteEntry {
851 fn live(&self) -> bool {
853 self.advertised && !self.stale
854 }
855
856 fn is_anonymous(&self) -> bool {
857 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
858 }
859
860 fn entered(&self) -> Source {
862 match self.peer {
863 true => Source::Peer(self.via),
864 false => Source::Local,
865 }
866 }
867
868 fn serves(&self, path: &Path) -> bool {
872 self.server.is_some() || (self.source.is_some() && self.prefix == *path)
873 }
874
875 fn qualifies(&self, pin: Pin) -> bool {
877 match pin {
878 Pin::Any => true,
879 Pin::Local => self.local,
880 Pin::Publisher(first) => self.hops.iter().next() == Some(&first),
881 Pin::Route(id) => self.id == id && self.hops.iter().next().is_none_or(|first| *first == Hop::UNKNOWN),
883 }
884 }
885
886 fn visible_to(&self, exclude: Option<Hop>) -> bool {
891 match exclude {
892 Some(peer) if peer != Hop::UNKNOWN => self.via != peer && !self.hops.contains(&peer),
893 _ => true,
894 }
895 }
896
897 fn overlaps(&self, allowed: &Patterns) -> bool {
899 self.scope.iter().any(|scope| {
900 scope
901 .intersect(&self.claim)
902 .is_ok_and(|scoped| scoped.iter().any(|restriction| allowed.overlaps(restriction)))
903 })
904 }
905}
906
907fn prefix_claim(prefix: &Path) -> Result<Pattern, InvalidPattern> {
909 if prefix.parts().count() == Path::MAX_PARTS {
910 Pattern::literal(prefix.as_str())
911 } else {
912 Pattern::subtree(prefix.as_str())
913 }
914}
915
916#[derive(Default)]
921struct ServeState {
922 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
925
926 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
932
933 closed: bool,
936}
937
938type FrontKey = (PathOwned, Horizon);
943
944#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash)]
946struct Horizon {
947 exclude: Option<Hop>,
951 local: bool,
953}
954
955impl Horizon {
956 fn admits(&self, entry: &RouteEntry) -> bool {
958 !(self.local && entry.peer) && entry.visible_to(self.exclude)
959 }
960}
961
962#[derive(Clone)]
965struct RemoteFront {
966 request: kio::Producer<PendingBroadcast>,
970 broadcast: broadcast::WeakConsumer,
973 pin: kio::Lock<Pin>,
976}
977
978type CursorRoute = (u64, RouteMeta, bool, Option<Vec<Pattern>>);
980
981impl WeakEntry for RemoteFront {
982 fn is_closed(&self) -> bool {
983 self.broadcast.is_closed()
984 }
985
986 fn same_channel(&self, other: &Self) -> bool {
987 self.broadcast.same_channel(&other.broadcast)
988 }
989}
990
991struct TableCursor {
994 root: PathOwned,
996 allowed: Patterns,
998 heads: Vec<PathOwned>,
1002 horizon: Horizon,
1004 hidden: Hidden,
1006 state: kio::Producer<OriginConsumerState>,
1009 under: PathOwned,
1012 holes: Vec<PathOwned>,
1015 named: Option<(Mount, Patterns)>,
1018 current: HashMap<PathOwned, CursorRoute>,
1023}
1024
1025impl TableCursor {
1026 fn presented(&self, prefix: &Path, claim: &Pattern) -> Option<PathOwned> {
1032 if !self.allowed.overlaps(claim) {
1033 return None;
1034 }
1035
1036 if let Some(relative) = prefix.strip_prefix(&self.root) {
1037 return Some(relative.to_owned());
1038 }
1039 self.root.has_prefix(prefix).then(PathOwned::default)
1040 }
1041
1042 fn captures(&self, prefix: &Path) -> Option<Vec<Pattern>> {
1045 let (literal, allowed) = match &self.named {
1046 Some((mount, allowed)) => (Pattern::literal(mount.name(prefix)?.as_str()).ok()?, allowed),
1047 None => (Pattern::literal(prefix.as_str()).ok()?, &self.allowed),
1048 };
1049 allowed
1050 .iter()
1051 .filter_map(|allowed| {
1052 allowed
1053 .captures(&literal)
1054 .map(|captures| (allowed.specificity(), captures))
1055 })
1056 .max_by_key(|(specificity, _)| *specificity)
1057 .map(|(_, captures)| captures)
1058 }
1059
1060 fn visible(&self, entry: &RouteEntry) -> bool {
1064 entry.live()
1065 && self.horizon.admits(entry)
1066 && entry.overlaps(&self.allowed)
1067 && self.discovers(&entry.prefix)
1068 && !self.holes.iter().any(|hole| entry.prefix.has_prefix(hole))
1069 && self.named.as_ref().is_none_or(|(mount, _)| mount.names(&entry.prefix))
1070 }
1071
1072 fn update(&mut self, presented: &PathOwned, best: Option<&RouteEntry>) {
1074 match best {
1075 Some(entry) => {
1076 let meta = (entry.hops.clone(), entry.cost, entry.entered());
1077 let served = entry.server.is_some();
1078 let captures = self.captures(&entry.prefix);
1079 let previous = self
1080 .current
1081 .insert(presented.clone(), (entry.id, meta.clone(), served, captures.clone()));
1082 match previous {
1083 Some((_, prev, prev_served, prev_captures))
1090 if prev == meta && prev_served == served && prev_captures == captures => {}
1091 Some((_, prev, _, prev_captures)) if prev_captures != captures => {
1094 if let Ok(mut state) = self.state.write() {
1095 state.apply_unannounce(self.under.join(presented), prev, prev_captures);
1096 state.apply_announce(self.under.join(presented), meta, captures);
1097 }
1098 }
1099 _ => {
1100 if let Ok(mut state) = self.state.write() {
1101 state.apply_announce(self.under.join(presented), meta, captures);
1102 }
1103 }
1104 }
1105 }
1106 None => {
1107 if let Some((_, last, _, captures)) = self.current.remove(presented)
1108 && let Ok(mut state) = self.state.write()
1109 {
1110 state.apply_unannounce(self.under.join(presented), last, captures);
1111 }
1112 }
1113 }
1114 }
1115
1116 fn discovers(&self, prefix: &Path) -> bool {
1118 self.hidden.discovers(&self.heads, prefix)
1119 }
1120}
1121
1122#[derive(Clone)]
1125struct OriginScope {
1126 allowed: Patterns,
1129 mounts: Arc<[Mount]>,
1131}
1132
1133#[derive(Clone, Debug)]
1136struct Mount {
1137 at: PathOwned,
1138 target: PathOwned,
1139}
1140
1141impl Mount {
1142 fn resolve(&self, path: &Path) -> Option<PathOwned> {
1145 let resolved = self.target.join(path.strip_prefix(&self.at)?);
1146 (resolved.parts().count() <= Path::MAX_PARTS).then_some(resolved)
1147 }
1148
1149 fn name(&self, path: &Path) -> Option<PathOwned> {
1152 Some(self.at.join(path.strip_prefix(&self.target)?))
1153 }
1154
1155 fn names(&self, path: &Path) -> bool {
1159 path.strip_prefix(&self.target)
1160 .is_none_or(|rest| self.at.parts().count() + rest.parts().count() <= Path::MAX_PARTS)
1161 }
1162
1163 fn translate(&self, patterns: &Patterns) -> Patterns {
1170 let target = self.target.as_str();
1171 patterns
1172 .rebase(self.at.as_str())
1173 .iter()
1174 .filter_map(|member| {
1175 member
1176 .rooted(target)
1177 .or_else(|_| {
1178 let segments = member
1179 .segments()
1180 .iter()
1181 .filter(|segment| **segment != Segment::Globstar);
1182 Pattern::new(segments.cloned())?.rooted(target)
1183 })
1184 .ok()
1185 })
1186 .collect()
1187 }
1188
1189 fn translate_head(&self, head: &Path) -> Option<PathOwned> {
1192 match head.strip_prefix(&self.at) {
1193 Some(rest) => Some(self.target.join(rest)),
1194 None => self.at.has_prefix(head).then(|| self.target.clone()),
1195 }
1196 }
1197}
1198
1199impl OriginScope {
1200 fn empty() -> Self {
1202 Self {
1203 allowed: Patterns::new(),
1204 mounts: Arc::from([]),
1205 }
1206 }
1207
1208 fn narrow(&self, patterns: &Patterns) -> Option<Self> {
1210 let allowed = self.allowed.intersect(patterns).ok()?;
1211 if allowed.is_empty() {
1212 None
1213 } else {
1214 Some(Self {
1215 allowed,
1216 mounts: self.mounts.clone(),
1217 })
1218 }
1219 }
1220
1221 fn mount(&self, path: &Path) -> Option<&Mount> {
1223 self.mounts.iter().find(|mount| path.has_prefix(&mount.at))
1224 }
1225
1226 fn resolve<'a>(&self, path: &'a Path<'a>) -> Option<Path<'a>> {
1229 match self.mount(path) {
1230 Some(mount) => mount.resolve(path),
1231 None => Some(path.borrow()),
1232 }
1233 }
1234
1235 fn publishes(&self, prefix: &Path) -> bool {
1238 self.mount(prefix).is_none()
1239 }
1240
1241 fn permits(&self, path: &Path) -> bool {
1243 self.allowed.matches(path.as_str())
1244 }
1245
1246 fn relative(&self, root: &Path) -> Patterns {
1248 self.allowed.rebase(root.as_str())
1249 }
1250}
1251
1252impl Default for OriginScope {
1253 fn default() -> Self {
1254 Self {
1255 allowed: Patterns::from(Pattern::all()),
1256 mounts: Arc::from([]),
1257 }
1258 }
1259}
1260
1261pub(crate) fn interest_prefixes(allowed: &Patterns) -> Vec<PathOwned> {
1264 let mut heads: Vec<PathOwned> = allowed
1265 .iter()
1266 .map(|pattern| Path::new(pattern.head()).to_owned())
1267 .collect();
1268 heads.sort();
1269 heads.dedup();
1270 let covered = heads.clone();
1271 heads.retain(|head| !covered.iter().any(|other| other != head && head.has_prefix(other)));
1272 heads
1273}
1274
1275#[derive(Clone, Debug, Default, PartialEq, Eq)]
1282pub(crate) struct Hidden {
1283 include: bool,
1285 from: Option<PathOwned>,
1287 beyond: Option<Vec<PathOwned>>,
1290}
1291
1292impl Hidden {
1293 fn discovers(&self, heads: &[PathOwned], prefix: &Path) -> bool {
1295 (self.include || !hides(self.from.as_ref().map(std::slice::from_ref).unwrap_or(heads), prefix))
1296 && self.beyond.as_ref().is_none_or(|outer| hides(outer, prefix))
1297 }
1298
1299 fn translate(&self, mount: &Mount) -> Self {
1303 let beyond = self.beyond.as_ref().and_then(|outer| {
1304 (!hides(outer, &mount.at)).then(|| outer.iter().filter_map(|head| mount.translate_head(head)).collect())
1305 });
1306 Self {
1307 include: self.include,
1308 from: self.from.as_ref().and_then(|head| mount.translate_head(head)),
1309 beyond,
1310 }
1311 }
1312}
1313
1314fn hides(heads: &[PathOwned], prefix: &Path) -> bool {
1317 heads
1318 .iter()
1319 .any(|head| prefix.strip_prefix(head).is_some_and(|below| below.is_hidden()))
1320}
1321
1322#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1324pub enum AnnounceKind {
1325 Announced,
1327 Updated,
1329 Retracted,
1331}
1332
1333impl AnnounceKind {
1334 pub fn is_active(self) -> bool {
1336 !matches!(self, Self::Retracted)
1337 }
1338}
1339
1340#[derive(Clone, Debug)]
1348pub struct AnnounceUpdate {
1349 pub prefix: PathOwned,
1351 pub captures: Option<Vec<Pattern>>,
1354 pub route: Route,
1357 pub kind: AnnounceKind,
1359}
1360
1361#[derive(Clone)]
1363pub struct Producer {
1364 hop: Hop,
1367
1368 scope: OriginScope,
1370
1371 root: PathOwned,
1373
1374 shared: kio::Shared<OriginState>,
1377
1378 pool: cache::Pool,
1381
1382 cache_duration: Duration,
1385
1386 default_max_age: Duration,
1389
1390 stats: stats::Session,
1394
1395 peer: bool,
1398
1399 tasks: Tasks,
1403
1404 timers: Clock,
1406}
1407
1408impl Producer {
1409 pub fn new(config: Config) -> (Self, Driver) {
1416 let (tasks, set) = TaskSet::new();
1417 let scope = OriginScope::default();
1418 let shared = kio::Shared::new(OriginState::new(config.update_hold));
1419 let timers = Clock::default();
1420 let pool = config.pool.clone();
1421 let producer = Self {
1422 hop: config.hop,
1423 scope: scope.clone(),
1424 root: PathOwned::default(),
1425 shared: shared.clone(),
1426 pool: config.pool,
1427 cache_duration: config.cache_duration,
1428 default_max_age: config.default_max_age,
1429 stats: stats::Session::default(),
1430 peer: false,
1431 tasks,
1432 timers: timers.clone(),
1433 };
1434 let driver = Driver {
1435 state: DriverState {
1436 set,
1437 shared,
1438 done: false,
1439 },
1440 timers,
1441 pool,
1442 };
1443 (producer, driver)
1444 }
1445
1446 pub fn with_stats(mut self, session: stats::Session) -> Self {
1450 self.stats = session;
1451 self
1452 }
1453
1454 pub fn peer(mut self) -> Self {
1462 self.peer = true;
1463 self
1464 }
1465
1466 pub fn config(&self) -> Config {
1468 Config {
1469 hop: self.hop,
1470 pool: self.pool.clone(),
1471 cache_duration: self.cache_duration,
1472 default_max_age: self.default_max_age,
1473 update_hold: self.shared.lock().update_hold,
1474 }
1475 }
1476
1477 pub fn hop(&self) -> Hop {
1479 self.hop
1480 }
1481
1482 pub(crate) fn default_max_age(&self) -> Duration {
1485 self.default_max_age
1486 }
1487
1488 pub(crate) fn empty(hop: Hop) -> Self {
1493 let (tasks, _) = TaskSet::new();
1496 Self {
1497 hop,
1498 scope: OriginScope::empty(),
1499 root: PathOwned::default(),
1500 shared: kio::Shared::new(OriginState::new(DEFAULT_UPDATE_HOLD)),
1501 pool: cache::Pool::default(),
1502 cache_duration: Duration::MAX,
1503 default_max_age: track::DEFAULT_MAX_AGE,
1504 stats: stats::Session::default(),
1505 peer: false,
1506 tasks,
1507 timers: Clock::default(),
1508 }
1509 }
1510
1511 pub fn create_broadcast(&self, path: impl AsPath) -> Result<broadcast::Producer, Error> {
1542 let path = path.as_path();
1543
1544 let full = self.root.join(&path).to_owned();
1545 if !self.scope.permits(&full) || !self.scope.publishes(&full) {
1546 return Err(Error::Unauthorized);
1547 }
1548 if full.parts().count() > Path::MAX_PARTS {
1552 return Err(BoundsExceeded.into());
1553 }
1554 let claim = prefix_claim(&full)?;
1557
1558 let ingress = self.stats.ingress(&full);
1560
1561 let announcing = Announcing {
1566 hop: self.hop,
1567 shared: self.shared.clone(),
1568 requested: full.clone(),
1569 prefixes: vec![(full.clone(), claim)],
1570 scope: self.scope.allowed.clone(),
1571 local: true,
1572 peer: self.peer,
1573 stats: self.stats.clone(),
1574 };
1575 let info = broadcast::Info {
1576 pool: self.pool.clone(),
1577 cache_duration: self.cache_duration,
1578 path: full,
1579 };
1580 let source = info.produce().with_stats(ingress.clone());
1581 let entry = announcing.announce(
1582 Route::default(),
1583 Serving {
1584 server: None,
1585 source: Some(source.consume()),
1586 advertised: false,
1587 },
1588 )?;
1589 Ok(source.with_announcer(Announcer {
1590 entry,
1591 ingress,
1592 _keepalive: self.tasks.keepalive(),
1593 }))
1594 }
1595
1596 pub fn publish(&self, path: impl AsPath, route: Route) -> Result<broadcast::Producer, Error> {
1598 let broadcast = self.create_broadcast(path)?;
1599 broadcast.announce(route)?;
1600 Ok(broadcast)
1601 }
1602
1603 pub(crate) fn create_source(&self, path: impl AsPath) -> broadcast::Producer {
1609 let path = path.as_path();
1610 let full = self.root.join(&path).to_owned();
1611 let ingress = self.stats.ingress(&full);
1612 broadcast::Info {
1613 pool: self.pool.clone(),
1614 cache_duration: self.cache_duration,
1615 path: full,
1616 }
1617 .produce()
1618 .with_stats(ingress)
1619 }
1620
1621 #[cfg(test)]
1630 pub(crate) fn announce(&self, prefix: impl AsPath, route: Route) -> Result<AnnounceProducer, Error> {
1631 Announcing::new(self, prefix)?.announce(
1632 route,
1633 Serving {
1634 server: None,
1635 source: None,
1636 advertised: true,
1637 },
1638 )
1639 }
1640
1641 pub fn dynamic(&self, prefix: impl AsPath, route: Route) -> Result<Dynamic, Error> {
1663 let announcing = Announcing::new(self, prefix)?;
1664 let serve = kio::Shared::<ServeState>::default();
1665 serve.lock().requests.add_handler();
1666 let announcement = announcing.announce(
1667 route,
1668 Serving {
1669 server: Some(serve.clone()),
1670 source: None,
1671 advertised: true,
1672 },
1673 )?;
1674 Ok(Dynamic {
1675 announcement,
1676 state: serve,
1677 _keepalive: self.tasks.keepalive(),
1678 })
1679 }
1680
1681 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Producer, Error> {
1688 let root = self.root.join(root).to_owned();
1689 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
1690 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
1691 Ok(Producer {
1692 hop: self.hop,
1693 scope,
1694 root,
1695 shared: self.shared.clone(),
1696 pool: self.pool.clone(),
1697 cache_duration: self.cache_duration,
1698 default_max_age: self.default_max_age,
1699 stats: self.stats.clone(),
1700 peer: self.peer,
1701 tasks: self.tasks.clone(),
1702 timers: self.timers.clone(),
1703 })
1704 }
1705
1706 pub fn mount(&self, at: impl AsPath, target: impl AsPath) -> Result<Producer, Error> {
1725 let at = self.root.join(at).to_owned();
1726 let target = self.root.join(target).to_owned();
1727 if [&at, &target]
1728 .into_iter()
1729 .any(|path| path.parts().count() > Path::MAX_PARTS)
1730 {
1731 return Err(BoundsExceeded.into());
1732 }
1733 Pattern::literal(at.as_str())?;
1735 if !self
1736 .scope
1737 .allowed
1738 .covers(&Patterns::from(Pattern::subtree(target.as_str())?))
1739 {
1740 return Err(Error::Unauthorized);
1741 }
1742 let overlaps = |a: &Path, b: &Path| a.has_prefix(b) || b.has_prefix(a);
1746 if overlaps(&at, &target)
1747 || self
1748 .scope
1749 .mounts
1750 .iter()
1751 .any(|mount| overlaps(&at, &mount.at) || overlaps(&target, &mount.at) || overlaps(&at, &mount.target))
1752 {
1753 return Err(Error::Duplicate);
1754 }
1755 let mounts = self
1756 .scope
1757 .mounts
1758 .iter()
1759 .cloned()
1760 .chain([Mount { at, target }])
1761 .collect();
1762 Ok(Producer {
1763 scope: OriginScope {
1764 allowed: self.scope.allowed.clone(),
1765 mounts,
1766 },
1767 ..self.clone()
1768 })
1769 }
1770
1771 pub fn consume(&self) -> Consumer {
1776 Consumer::from_producer(self, stats::Session::default())
1779 }
1780
1781 pub fn root(&self) -> &Path<'_> {
1783 &self.root
1784 }
1785
1786 pub fn allowed(&self) -> Patterns {
1788 self.scope.relative(&self.root)
1789 }
1790
1791 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1793 self.root.join(path)
1794 }
1795}
1796
1797struct Announcing {
1801 hop: Hop,
1802 shared: kio::Shared<OriginState>,
1803 requested: PathOwned,
1805 prefixes: Vec<(PathOwned, Pattern)>,
1809 scope: Patterns,
1811 local: bool,
1812 peer: bool,
1814 stats: stats::Session,
1815}
1816
1817impl Announcing {
1818 fn new(producer: &Producer, prefix: impl AsPath) -> Result<Self, Error> {
1820 let requested = producer.root.join(prefix.as_path()).to_owned();
1821 if requested.parts().count() > Path::MAX_PARTS {
1822 return Err(BoundsExceeded.into());
1823 }
1824 let claim = prefix_claim(&requested)?;
1825 if !producer.scope.allowed.overlaps(&claim) || !producer.scope.publishes(&requested) {
1826 return Err(Error::Unauthorized);
1827 }
1828 Ok(Self {
1829 hop: producer.hop,
1830 shared: producer.shared.clone(),
1831 requested: requested.clone(),
1832 prefixes: vec![(requested, claim)],
1833 scope: producer.scope.allowed.clone(),
1834 local: false,
1835 peer: producer.peer,
1836 stats: producer.stats.clone(),
1837 })
1838 }
1839
1840 fn announce(&self, route: Route, serving: Serving) -> Result<AnnounceProducer, Error> {
1841 debug_assert!(
1842 !route.hops.contains(&self.hop),
1843 "announce called with a looping hop chain",
1844 );
1845
1846 let via = route.via;
1847
1848 let mut shared = self.shared.lock();
1849 if shared.closed {
1850 return Err(Error::Closed);
1851 }
1852
1853 let mut entries = Vec::with_capacity(self.prefixes.len());
1854 for (prefix, claim) in &self.prefixes {
1855 shared.reannounced(prefix, &route.hops);
1856 let id = shared.next_route;
1857 shared.next_route += 1;
1858 let stale = shared.withdrawn_through(prefix, &route.hops);
1859 shared.routes.insert(RouteEntry {
1860 id,
1861 prefix: prefix.clone(),
1862 scope: self.scope.clone(),
1863 hops: route.hops.clone(),
1864 cost: route.cost,
1865 via,
1866 local: self.local,
1867 peer: self.peer,
1868 server: serving.server.clone(),
1869 source: serving.source.clone(),
1870 advertised: serving.advertised,
1871 stale,
1872 claim: claim.clone(),
1873 });
1874 shared.sync_route(prefix, claim);
1875 entries.push((prefix.clone(), id));
1876 }
1877 drop(shared);
1878
1879 let guard = serving
1881 .advertised
1882 .then(|| self.stats.ingress(&self.requested).announce());
1883
1884 Ok(AnnounceProducer {
1885 shared: self.shared.clone(),
1886 entries,
1887 guard,
1888 })
1889 }
1890}
1891
1892struct Serving {
1894 server: Option<kio::Shared<ServeState>>,
1895 source: Option<broadcast::Consumer>,
1896 advertised: bool,
1897}
1898
1899pub(crate) struct Announcer {
1906 entry: AnnounceProducer,
1907 ingress: stats::Scope,
1909 _keepalive: Keepalive,
1913}
1914
1915impl Announcer {
1916 pub(crate) fn announce(&mut self, route: Route) -> Result<(), Error> {
1918 self.entry.update(route)?;
1919 if self.entry.guard.is_none() {
1920 self.entry.guard = Some(self.ingress.announce());
1921 }
1922 Ok(())
1923 }
1924
1925 pub(crate) fn withdraw(&mut self) {
1927 self.entry.withdraw();
1928 self.entry.guard = None;
1929 }
1930}
1931
1932#[must_use = "dropping an announcement retracts the route"]
1939pub(crate) struct AnnounceProducer {
1940 shared: kio::Shared<OriginState>,
1941 entries: Vec<(PathOwned, u64)>,
1944 guard: Option<stats::Announce>,
1946}
1947
1948impl AnnounceProducer {
1949 pub fn update(&self, route: Route) -> Result<(), Error> {
1957 let mut shared = self.shared.lock();
1958 if shared.closed {
1959 return Err(Error::Closed);
1960 }
1961 for (prefix, id) in &self.entries {
1962 shared.reannounced(prefix, &route.hops);
1963 let stale = shared.withdrawn_through(prefix, &route.hops);
1964 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1966 return Err(Error::Closed);
1967 };
1968 if entry.hops.iter().next() != route.hops.iter().next()
1971 && let Some(server) = &entry.server
1972 {
1973 drop(std::mem::take(&mut server.lock().served));
1974 }
1975 entry.hops = route.hops.clone();
1976 entry.stale = stale;
1977 entry.cost = route.cost;
1978 entry.via = route.via;
1979 entry.advertised = true;
1980 let claim = entry.claim.clone();
1981 shared.sync_route(prefix, &claim);
1982 shared.prune_withdrawn(prefix);
1983 }
1984 Ok(())
1985 }
1986
1987 fn withdraw(&self) {
1992 let mut shared = self.shared.lock();
1993 for (prefix, id) in &self.entries {
1994 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1995 continue;
1996 };
1997 if !entry.advertised {
1998 continue;
1999 }
2000 entry.advertised = false;
2001 let claim = entry.claim.clone();
2002 shared.sync_route(prefix, &claim);
2003 }
2004 }
2005
2006 fn retract(&self, withdrawn: bool) {
2009 let mut shared = self.shared.lock();
2010 for (prefix, id) in &self.entries {
2011 let Some(entry) = shared.routes.remove(prefix, *id) else {
2012 continue;
2013 };
2014 if let Some(server) = &entry.server {
2017 let mut server = server.lock();
2018 server.closed = true;
2019 for producer in server.requests.drain_all() {
2020 if let Ok(mut request) = producer.write() {
2021 request.resolved.get_or_insert(Err(Error::Unroutable));
2022 }
2023 }
2024 }
2025 if withdrawn
2030 && let Some(&peer) = entry.hops.iter().last()
2031 && peer != Hop::UNKNOWN
2032 && !shared
2033 .routes
2034 .at(prefix)
2035 .any(|other| other.live() && other.hops.iter().last() == Some(&peer))
2036 {
2037 shared.withdrawn.entry(prefix.clone()).or_default().insert(peer);
2038 shared.restale(prefix);
2039 }
2040 shared.sync_route(&entry.prefix, &entry.claim);
2041 shared.prune_withdrawn(&entry.prefix);
2042 }
2043 }
2044}
2045
2046impl Drop for AnnounceProducer {
2047 fn drop(&mut self) {
2048 self.retract(false);
2049 }
2050}
2051
2052#[must_use = "poll the driver or the origin makes no progress"]
2067pub struct Driver {
2068 state: DriverState,
2069 timers: Clock,
2071 pool: cache::Pool,
2074}
2075
2076struct DriverState {
2078 set: TaskSet,
2080 shared: kio::Shared<OriginState>,
2083 done: bool,
2085}
2086
2087impl Driver {
2088 pub fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
2095 self.timers.advance(now);
2096 self.timers.register_driver(waiter);
2097 let result = self.state.poll(waiter);
2098 let gc = self.pool.gc(now);
2099 if result.is_ready() {
2100 return Err(Error::Closed);
2101 }
2102 Ok(self.timers.timeout().into_iter().chain(gc).min())
2103 }
2104}
2105
2106impl crate::time::Driver for Driver {
2107 fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
2108 self.poll(now, waiter)
2109 }
2110}
2111
2112impl DriverState {
2113 fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
2114 if !self.done {
2118 ready!(self.set.poll(waiter));
2119 self.done = true;
2120 }
2121 Poll::Ready(())
2122 }
2123
2124 fn teardown(&mut self) {
2128 drop(std::mem::replace(&mut self.set, TaskSet::owned()));
2131
2132 let (servers, cursors, fronts) = {
2137 let mut shared = self.shared.lock();
2138 shared.closed = true;
2139 shared.routes.poke_all();
2141 let servers: Vec<_> = shared
2142 .routes
2143 .entries()
2144 .filter_map(|entry| entry.server.clone())
2145 .collect();
2146 let cursors: Vec<_> = shared.cursors.values().map(|cursor| cursor.state.clone()).collect();
2147 let fronts: Vec<_> = shared.fronts.values().map(|front| front.request.clone()).collect();
2148 (servers, cursors, fronts)
2149 };
2150
2151 for producer in fronts {
2154 if let Ok(mut request) = producer.write() {
2155 request.resolved.get_or_insert(Err(Error::Dropped));
2156 }
2157 }
2158 for server in servers {
2162 let mut server = server.lock();
2163 server.closed = true;
2164 for producer in server.requests.drain_all() {
2165 if let Ok(mut request) = producer.write() {
2166 request.resolved.get_or_insert(Err(Error::Dropped));
2167 }
2168 }
2169 }
2170
2171 for state in cursors {
2174 if let Ok(mut state) = state.write() {
2175 state.ended = true;
2176 }
2177 }
2178 }
2179}
2180
2181impl Drop for DriverState {
2182 fn drop(&mut self) {
2183 self.teardown();
2184 }
2185}
2186
2187const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
2201
2202struct WarmCopy {
2207 track: track::Producer,
2208 _dynamic: track::Dynamic,
2209 edge: Option<WarmGroup>,
2212}
2213
2214impl Drop for WarmCopy {
2215 fn drop(&mut self) {
2216 let _ = self.track.finish();
2217 }
2218}
2219
2220struct WarmGroup(group::Producer);
2227
2228impl Drop for WarmGroup {
2229 fn drop(&mut self) {
2230 if !self.0.is_finished() {
2231 let _ = self.0.clone().abort(Error::Cancel);
2232 }
2233 }
2234}
2235
2236struct HeldGroup(group::Producer);
2241
2242impl Drop for HeldGroup {
2243 fn drop(&mut self) {
2244 if !self.0.is_finished() {
2245 let _ = self.0.abort_if_last(Error::Cancel);
2246 }
2247 }
2248}
2249
2250fn warm_copy(source: &track::Consumer, head: Option<&WarmGroup>, held: Option<&group::Producer>) -> Option<WarmCopy> {
2258 let info = source.cached_info()?;
2259 let mut track = track::Producer::new(Arc::new(source.broadcast().clone()), source.name(), info);
2260 let head = head.map(|head| &head.0);
2261 let mut groups = source.cached_groups();
2262 if let Some(held) = held
2263 && !groups.iter().any(|(group, _)| group.sequence == held.sequence)
2264 {
2265 groups.push((held.clone(), true));
2266 }
2267 let latest = groups
2269 .iter()
2270 .filter(|(_, visible)| *visible)
2271 .map(|(group, _)| group.sequence)
2272 .max();
2273 let mut edge = None;
2274 let mut first = None;
2275
2276 if let Some(head) = head
2280 && !groups.iter().any(|(group, _)| group.sequence == head.sequence)
2281 {
2282 let is_latest = latest.is_none_or(|latest| head.sequence >= latest);
2283 if head.is_finished() {
2284 let _ = track.adopt_group(head.clone(), true);
2285 first = Some(head.sequence);
2286 if is_latest {
2287 edge = Some(WarmGroup(head.clone()));
2288 }
2289 } else if is_latest {
2290 edge = warm_rebuild(&track, head, None);
2291 first = edge.as_ref().map(|_| head.sequence);
2292 }
2293 }
2294
2295 for (group, visible) in groups {
2296 let finished = group.is_finished();
2297 let whole = group.live_first_frame() == Some(0);
2298 let is_latest = visible && Some(group.sequence) == latest;
2299 let warm = if finished && whole {
2300 let _ = track.adopt_group(group.clone(), visible);
2301 Some(WarmGroup(group))
2302 } else if (finished || is_latest) && (whole || head.is_some_and(|head| head.sequence == group.sequence)) {
2303 warm_rebuild(&track, &group, head)
2309 } else {
2310 None
2312 };
2313 if visible && let Some(warm) = &warm {
2314 first = Some(first.map_or(warm.0.sequence, |first: u64| first.min(warm.0.sequence)));
2315 }
2316 if is_latest {
2317 edge = warm;
2318 }
2319 }
2320 track.start_at(first).ok()?;
2324 let dynamic = track.dynamic();
2325 Some(WarmCopy {
2326 track,
2327 _dynamic: dynamic,
2328 edge,
2329 })
2330}
2331
2332fn warm_rebuild(track: &track::Producer, live: &group::Producer, head: Option<&group::Producer>) -> Option<WarmGroup> {
2335 let mut start = live.live_first_frame()? as u64;
2337 let mut tail = live.consume();
2338 tail.start_at(start);
2339 let mut frames = Vec::new();
2340 if start > 0
2341 && let Some(head) = head.filter(|head| head.sequence == live.sequence)
2342 && let Some(head_start) = head.live_first_frame()
2343 {
2344 let mut head = head.consume();
2345 head.start_at(head_start as u64);
2346 while head.index() < start {
2347 match head.poll_read_frame(&kio::Waiter::noop()) {
2348 Poll::Ready(Ok(Some(frame))) => frames.push(frame),
2349 _ => break,
2350 }
2351 }
2352 match head.index() == start {
2353 true => start = head_start as u64,
2354 false => frames.clear(),
2356 }
2357 }
2358 while let Poll::Ready(Ok(Some(frame))) = tail.poll_read_frame(&kio::Waiter::noop()) {
2359 frames.push(frame);
2360 }
2361 if frames.is_empty() {
2362 return None;
2363 }
2364
2365 let rebuilt = WarmGroup(
2367 track
2368 .create_group(group::Info {
2369 sequence: live.sequence,
2370 })
2371 .ok()?,
2372 );
2373 let mut writer = rebuilt.0.clone();
2374 if start > 0 {
2375 writer.start_at(start).ok()?;
2376 }
2377 for frame in frames {
2378 writer.write_frame(frame.timestamp, frame.payload).ok()?;
2379 }
2380 if live.is_finished() {
2381 writer.finish().ok()?;
2382 }
2383 Some(rebuilt)
2384}
2385
2386struct FrontTask {
2388 shared: kio::Shared<OriginState>,
2390 broadcast: broadcast::Producer,
2392 path: PathOwned,
2394 horizon: Horizon,
2396 watch: Watch,
2398 request: kio::Producer<PendingBroadcast>,
2400 pin: kio::Lock<Pin>,
2402 timers: Clock,
2403}
2404
2405struct TrackIo {
2408 resume: super::resume::Producer,
2409 staged: Option<(u64, track::Consumer)>,
2411 query: Option<(u64, track::Consumer, track::Querying)>,
2413 copy: Option<(u64, track::Consumer)>,
2415 edge: Option<track::Position>,
2420 warm: Option<WarmCopy>,
2423 head: Option<WarmGroup>,
2426 held: Option<HeldGroup>,
2431 dead: Option<track::Consumer>,
2434 used: bool,
2436}
2437
2438impl TrackIo {
2439 fn end(&mut self) {
2442 self.staged = None;
2443 self.query = None;
2444 self.copy = None;
2445 self.warm = None;
2446 self.head = None;
2447 self.held = None;
2448 self.dead = None;
2449 }
2450
2451 fn bury(&mut self) {
2454 let Some((_, copy)) = self.copy.take() else { return };
2455 let after = self.held.as_ref().map(|held| held.0.sequence);
2456 if let Poll::Ready(group) = copy.poll_latest_group(after, &kio::Waiter::noop()) {
2457 self.held = Some(HeldGroup(group));
2458 }
2459 self.dead = Some(copy);
2460 }
2461
2462 fn retire(&mut self, copy: &track::Consumer) -> Result<(), Error> {
2467 let Some(held) = self.held.take() else {
2468 return Ok(());
2469 };
2470 if held.0.is_finished() {
2471 return Ok(());
2472 }
2473 let Some(warm) = warm_copy(copy, self.head.as_ref(), Some(&held.0)) else {
2474 return Ok(());
2475 };
2476 self.head = None;
2477 self.resume.park(&warm.track)?;
2478 self.warm = Some(warm);
2479 Ok(())
2480 }
2481}
2482
2483async fn run_front(task: FrontTask) {
2487 let FrontTask {
2488 shared,
2489 broadcast,
2490 path,
2491 horizon,
2492 watch,
2493 request,
2494 pin,
2495 timers,
2496 } = task;
2497
2498 enum Step {
2500 Assigned(Arc<str>, super::resume::Producer),
2501 Resolved(u64, Result<broadcast::Consumer, Error>),
2502 SourceClosed(u64),
2503 Info(Arc<str>, u64, Result<track::Info, Error>),
2504 Ended(Arc<str>, u64, Result<(), Error>),
2505 Held(Arc<str>, group::Producer),
2506 Demand(Arc<str>),
2507 Deadline,
2508 Table,
2509 }
2510
2511 let mut front = Front::new(TRACK_IDLE_LINGER);
2512 let mut sources: HashMap<u64, broadcast::Consumer> = HashMap::new();
2513 let mut next_source = 0u64;
2514 let mut upstream: Option<(u64, kio::Consumer<PendingBroadcast>)> = None;
2516 let mut tracks: HashMap<Arc<str>, TrackIo> = HashMap::new();
2517 let mut deadline = crate::runtime::Deadline::new(&timers);
2518 let mut seen = 0;
2520 let mut events: VecDeque<Event> = VecDeque::new();
2521
2522 let select = |front: &mut Front, sources: &HashMap<u64, broadcast::Consumer>, seen: &mut u64| -> Event {
2525 let table = shared.read();
2526 if table.closed {
2527 return Event::Closed;
2528 }
2529 *seen = watch.seen();
2531 front.retain_routes(|route| table.routes.covers(&path.as_path(), route));
2532 let best = table
2533 .best_route(&path.as_path(), horizon, front.pin(), front.refused_routes())
2534 .map(|entry| Candidate {
2535 route: entry.id,
2536 first: entry.hops.iter().next().copied(),
2537 local: entry.local,
2538 });
2539 let serving_closing = front
2540 .serving()
2541 .and_then(|id| sources.get(&id))
2542 .is_some_and(|source| source.is_closing());
2543 Event::Selected { best, serving_closing }
2544 };
2545
2546 events.push_back(select(&mut front, &sources, &mut seen));
2547
2548 loop {
2549 while let Some(event) = events.pop_front() {
2550 for action in front.step(event) {
2551 match action {
2552 Action::Reselect => events.push_back(select(&mut front, &sources, &mut seen)),
2553 Action::Request { route } => {
2554 let found = {
2556 let table = shared.read();
2557 table
2558 .routes
2559 .covering(&path.as_path())
2560 .find(|entry| entry.id == route && entry.live())
2561 .map(|entry| {
2562 (
2563 Candidate {
2564 route,
2565 first: entry.hops.iter().next().copied(),
2566 local: entry.local,
2567 },
2568 entry.source.clone(),
2569 entry.server.clone(),
2570 )
2571 })
2572 };
2573 let Some((candidate, source, server)) = found else {
2574 events.push_back(Event::Resolved {
2575 route,
2576 result: Err(Refusal {
2577 err: Error::Unroutable,
2578 standing: false,
2579 }),
2580 });
2581 continue;
2582 };
2583 front.identify(candidate);
2584 *pin.lock() = front.pin();
2585 if let Some(source) = source {
2586 let id = next_source;
2587 next_source += 1;
2588 sources.insert(id, source);
2589 events.push_back(Event::Resolved { route, result: Ok(id) });
2590 continue;
2591 }
2592 let Some(server) = server else {
2593 events.push_back(Event::Resolved {
2594 route,
2595 result: Err(Refusal {
2596 err: Error::Unroutable,
2597 standing: true,
2598 }),
2599 });
2600 continue;
2601 };
2602 let mut serve = server.lock();
2603 if serve.closed {
2604 drop(serve);
2607 events.push_back(Event::Resolved {
2608 route,
2609 result: Err(Refusal {
2610 err: Error::Unroutable,
2611 standing: true,
2612 }),
2613 });
2614 continue;
2615 }
2616 if let Some(weak) = serve.served.get(&path) {
2619 drop(serve);
2620 let id = next_source;
2621 next_source += 1;
2622 sources.insert(id, weak.consume());
2623 events.push_back(Event::Resolved { route, result: Ok(id) });
2624 continue;
2625 }
2626 let pending = match serve.requests.join(&path) {
2627 Some(producer) => producer.consume(),
2628 None => {
2629 let producer = kio::Producer::<PendingBroadcast>::default();
2630 let consumer = producer.consume();
2631 match serve.requests.insert(path.clone(), producer) {
2632 Ok(()) => consumer,
2633 Err(_) => {
2636 drop(serve);
2637 events.push_back(Event::Resolved {
2638 route,
2639 result: Err(Refusal {
2640 err: Error::Unroutable,
2641 standing: true,
2642 }),
2643 });
2644 continue;
2645 }
2646 }
2647 }
2648 };
2649 upstream = Some((route, pending));
2650 }
2651 Action::Detach { source } => {
2652 sources.remove(&source);
2653 for io in tracks.values_mut() {
2656 if io.copy.as_ref().is_some_and(|(s, _)| *s == source) {
2657 io.bury();
2658 }
2659 if io.query.as_ref().is_some_and(|(s, ..)| *s == source) {
2660 io.query = None;
2661 }
2662 if io.staged.as_ref().is_some_and(|(s, _)| *s == source) {
2663 io.staged = None;
2664 }
2665 }
2666 }
2667 Action::Resolve => {
2668 if let Ok(mut pending) = request.write() {
2669 pending.resolved.get_or_insert(Ok(broadcast.consume()));
2670 }
2671 }
2672 Action::Query { track: name, source } => {
2673 let Some(io) = tracks.get_mut(&name) else { continue };
2674 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2675 match sources.get(&source).map(|s| s.track(&name)) {
2676 Some(Ok(copy)) => {
2677 let query = copy.query().into_inner();
2680 io.query = Some((source, copy, query));
2681 }
2682 Some(Err(err)) => events.push_back(Event::TrackInfo {
2683 track: name,
2684 source,
2685 closing,
2686 result: Err(err),
2687 }),
2688 None => {}
2689 }
2690 }
2691 Action::Splice { track: name, source } => {
2692 let Some(io) = tracks.get_mut(&name) else { continue };
2693 let Some((staged, copy)) = io.staged.take() else {
2694 continue;
2695 };
2696 if staged != source {
2697 continue;
2698 }
2699 let outgoing = io.copy.take().map(|(_, copy)| copy).or_else(|| io.dead.take());
2700 let retired = match outgoing {
2701 Some(outgoing) => io.retire(&outgoing),
2702 None => Ok(()),
2703 };
2704 if let Err(err) = retired.and_then(|()| io.resume.takeover(©)) {
2705 let _ = io.resume.abort(err);
2709 tracks.remove(&name);
2710 continue;
2711 }
2712 if let Some(head) = io.warm.take().and_then(|mut warm| warm.edge.take()) {
2713 io.head = Some(head);
2714 }
2715 io.edge = io.resume.resume_position();
2718 io.held = None;
2719 io.copy = Some((source, copy));
2720 }
2721 Action::Park { track: name } => {
2722 let Some(io) = tracks.get_mut(&name) else { continue };
2723 let Some((_, copy)) = io.copy.take() else { continue };
2724 let warm = warm_copy(©, io.head.as_ref(), io.held.as_ref().map(|held| &held.0));
2728 drop(copy);
2729 io.head = None;
2730 io.held = None;
2731 io.dead = None;
2732 let parked = match &warm {
2733 Some(warm) => io.resume.park(&warm.track),
2734 None => io.resume.release(),
2735 };
2736 if parked.is_err() {
2737 tracks.remove(&name);
2738 continue;
2739 }
2740 io.warm = warm;
2741 }
2742 Action::Forget { track: name } => {
2743 if let Some(io) = tracks.get_mut(&name)
2748 && !broadcast.forget_spliced(&name, &io.resume)
2749 {
2750 if !io.used {
2751 io.used = true;
2752 events.push_back(Event::Used { track: name });
2753 }
2754 continue;
2755 }
2756 tracks.remove(&name);
2757 events.push_back(Event::Forgotten { track: name });
2758 }
2759 Action::Finish { track: name } => {
2760 if let Some(io) = tracks.get_mut(&name) {
2761 io.end();
2762 let _ = io.resume.finish();
2763 }
2764 }
2765 Action::Abort { track: name, err } => {
2766 if let Some(io) = tracks.get_mut(&name) {
2767 tracing::debug!(name = %name, %err, "aborting track");
2768 io.end();
2769 let _ = io.resume.abort(err);
2770 }
2771 }
2772 Action::Arm { at } => deadline.set(at),
2773 Action::End { err } => {
2774 if let Ok(mut pending) = request.write() {
2775 pending.resolved.get_or_insert(Err(err.clone()));
2776 }
2777 broadcast.close();
2784 broadcast.release_spliced(err.clone());
2785 for (_, mut io) in tracks.drain() {
2786 let waiting = io.staged.take().map(|(_, copy)| copy);
2790 let waiting = waiting.or_else(|| io.query.take().map(|(_, copy, _)| copy));
2791 if let Some(copy) = waiting
2792 && io.resume.is_used()
2793 {
2794 if io.resume.takeover(©).is_err() {
2795 continue;
2796 }
2797 io.warm = None;
2798 }
2799 if !io.resume.is_used() || !io.resume.is_spliced() || io.warm.is_some() {
2801 let _ = io.resume.abort(err.clone());
2802 }
2803 }
2804 return;
2805 }
2806 }
2807 }
2808 }
2809
2810 let step = kio::wait(|waiter| {
2811 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
2812 return Poll::Ready(Step::Assigned(name, resume));
2813 }
2814 if let Some((route, pending)) = &upstream
2815 && let Poll::Ready(result) = pending.poll(waiter, |p| match &p.resolved {
2816 Some(result) => Poll::Ready(result.clone()),
2817 None => Poll::Pending,
2818 }) {
2819 return Poll::Ready(Step::Resolved(
2820 *route,
2821 match result {
2822 Ok(resolved) => resolved,
2823 Err(_closed) => Err(Error::Unroutable),
2826 },
2827 ));
2828 }
2829 if let Some(id) = front.serving()
2830 && let Some(source) = sources.get(&id)
2831 && source.poll_closed(waiter).is_ready()
2832 {
2833 return Poll::Ready(Step::SourceClosed(id));
2834 }
2835 for (name, io) in &tracks {
2836 if let Some((source, _, query)) = &io.query
2837 && let Poll::Ready(result) = query.poll(waiter)
2838 {
2839 return Poll::Ready(Step::Info(name.clone(), *source, result));
2840 }
2841 if let Some((source, copy)) = &io.copy
2842 && let Poll::Ready(result) = copy.poll_complete(waiter)
2843 {
2844 return Poll::Ready(Step::Ended(name.clone(), *source, result));
2845 }
2846 if let Some((_, copy)) = &io.copy
2847 && let Poll::Ready(group) =
2848 copy.poll_latest_group(io.held.as_ref().map(|held| held.0.sequence), waiter)
2849 {
2850 return Poll::Ready(Step::Held(name.clone(), group));
2851 }
2852 let edge = match io.used {
2854 true => io.resume.poll_unused(waiter),
2855 false => io.resume.poll_used(waiter),
2856 };
2857 if edge.is_ready() {
2858 return Poll::Ready(Step::Demand(name.clone()));
2859 }
2860 }
2861 if deadline.poll(waiter).is_ready() {
2862 return Poll::Ready(Step::Deadline);
2863 }
2864 watch.poll_changed(waiter, seen).map(|()| Step::Table)
2865 })
2866 .await;
2867
2868 let event = match step {
2869 Step::Assigned(name, resume) => {
2870 tracks.insert(
2871 name.clone(),
2872 TrackIo {
2873 resume,
2874 staged: None,
2875 query: None,
2876 copy: None,
2877 edge: None,
2878 warm: None,
2879 head: None,
2880 held: None,
2881 dead: None,
2882 used: false,
2883 },
2884 );
2885 Event::TrackAssigned {
2886 track: name,
2887 now: timers.now(),
2888 }
2889 }
2890 Step::Resolved(route, result) => {
2891 upstream = None;
2892 match result {
2893 Ok(source) => {
2894 let id = next_source;
2895 next_source += 1;
2896 sources.insert(id, source);
2897 Event::Resolved { route, result: Ok(id) }
2898 }
2899 Err(err) => {
2900 let standing =
2904 !matches!(err, Error::Unroutable) || shared.read().routes.covers(&path.as_path(), route);
2905 Event::Resolved {
2906 route,
2907 result: Err(Refusal { err, standing }),
2908 }
2909 }
2910 }
2911 }
2912 Step::SourceClosed(source) => Event::SourceClosed { source },
2913 Step::Info(name, source, result) => {
2914 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2915 let Some(io) = tracks.get_mut(&name) else { continue };
2916 let Some((_, copy, _)) = io.query.take() else { continue };
2917 let result = match result {
2920 Ok(info) => match copy.poll_complete(&kio::Waiter::noop()) {
2921 Poll::Ready(Err(err)) => Err(err),
2922 _ => Ok(info),
2923 },
2924 Err(err) => Err(err),
2925 };
2926 if result.is_ok() && io.used {
2929 io.staged = Some((source, copy));
2930 }
2931 Event::TrackInfo {
2932 track: name,
2933 source,
2934 closing,
2935 result,
2936 }
2937 }
2938 Step::Ended(name, source, result) => {
2939 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2940 let Some(io) = tracks.get_mut(&name) else { continue };
2941 io.bury();
2942 let delivered = io.resume.resume_position() != io.edge;
2943 Event::TrackEnded {
2944 track: name,
2945 source,
2946 closing,
2947 result,
2948 delivered,
2949 }
2950 }
2951 Step::Held(name, group) => {
2952 if let Some(io) = tracks.get_mut(&name) {
2953 io.held = Some(HeldGroup(group));
2954 }
2955 continue;
2956 }
2957 Step::Demand(name) => {
2958 let Some(io) = tracks.get_mut(&name) else { continue };
2959 io.used = io.resume.is_used();
2960 if !io.used {
2961 io.query = None;
2964 io.staged = None;
2965 }
2966 match io.used {
2967 true => Event::Used { track: name },
2968 false => Event::Unused {
2969 track: name,
2970 now: timers.now(),
2971 },
2972 }
2973 }
2974 Step::Deadline => {
2975 deadline.set(None);
2978 Event::Deadline { now: timers.now() }
2979 }
2980 Step::Table => select(&mut front, &sources, &mut seen),
2981 };
2982 events.push_back(event);
2983 }
2984}
2985
2986#[derive(Default)]
2991struct RouteTable {
2992 root: RouteNode,
2993}
2994
2995#[derive(Default)]
2998struct RouteNode {
2999 entries: Vec<RouteEntry>,
3001 cursors: Vec<ConsumerId>,
3003 cursors_below: usize,
3006 watches: Vec<(u64, kio::Producer<Watched>)>,
3009 watches_below: usize,
3012 children: HashMap<String, RouteNode>,
3013}
3014
3015#[derive(Default)]
3018struct Watched {
3019 generation: u64,
3020}
3021
3022struct Watch {
3028 shared: kio::Shared<OriginState>,
3029 path: PathOwned,
3030 id: u64,
3031 signal: kio::Consumer<Watched>,
3032}
3033
3034impl Watch {
3035 fn seen(&self) -> u64 {
3039 self.signal.read().generation
3040 }
3041
3042 fn poll_changed(&self, waiter: &kio::Waiter, seen: u64) -> Poll<()> {
3044 self.signal
3045 .poll(waiter, |watched| match watched.generation != seen {
3046 true => Poll::Ready(()),
3047 false => Poll::Pending,
3048 })
3049 .map(|_| ())
3050 }
3051}
3052
3053impl Drop for Watch {
3054 fn drop(&mut self) {
3055 self.shared.lock().routes.remove_watch(&self.path, self.id);
3056 }
3057}
3058
3059#[derive(Clone, Copy)]
3061struct Below {
3062 cursors: usize,
3063 watches: usize,
3064}
3065
3066impl Below {
3067 const NONE: Self = Self { cursors: 0, watches: 0 };
3068 const CURSOR: Self = Self { cursors: 1, watches: 0 };
3069 const WATCH: Self = Self { cursors: 0, watches: 1 };
3070}
3071
3072impl RouteNode {
3073 fn is_empty(&self) -> bool {
3075 self.entries.is_empty() && self.cursors.is_empty() && self.watches.is_empty() && self.children.is_empty()
3076 }
3077
3078 fn find<'a>(&self, mut parts: impl Iterator<Item = &'a str>) -> Option<&Self> {
3080 match parts.next() {
3081 None => Some(self),
3082 Some(part) => self.children.get(part)?.find(parts),
3083 }
3084 }
3085
3086 fn reach<'a>(&mut self, mut parts: impl Iterator<Item = &'a str>, below: Below) -> &mut Self {
3089 self.cursors_below += below.cursors;
3090 self.watches_below += below.watches;
3091 match parts.next() {
3092 None => self,
3093 Some(part) => self.children.entry(part.to_string()).or_default().reach(parts, below),
3094 }
3095 }
3096
3097 fn edit<'a, R>(
3101 &mut self,
3102 mut parts: impl Iterator<Item = &'a str>,
3103 below: Below,
3104 f: impl FnOnce(&mut Self) -> R,
3105 ) -> Option<R> {
3106 let result = match parts.next() {
3107 None => f(self),
3108 Some(part) => {
3109 let child = self.children.get_mut(part)?;
3110 let result = child.edit(parts, below, f)?;
3111 if child.is_empty() {
3112 self.children.remove(part);
3113 }
3114 result
3115 }
3116 };
3117 self.cursors_below -= below.cursors;
3118 self.watches_below -= below.watches;
3119 Some(result)
3120 }
3121
3122 fn poke(&self) {
3124 for (_, watch) in &self.watches {
3125 if let Ok(mut watched) = watch.write() {
3126 watched.generation += 1;
3127 }
3128 }
3129 }
3130
3131 fn poke_below(&self) {
3134 if self.watches_below == 0 {
3135 return;
3136 }
3137 self.poke();
3138 for child in self.children.values() {
3139 child.poke_below();
3140 }
3141 }
3142
3143 fn walk<'a>(&'a self, visit: &mut impl FnMut(&'a Self)) {
3145 visit(self);
3146 for child in self.children.values() {
3147 child.walk(visit);
3148 }
3149 }
3150
3151 fn collect_cursors(&self, out: &mut Vec<ConsumerId>) {
3153 if self.cursors_below == 0 {
3154 return;
3155 }
3156 out.extend(&self.cursors);
3157 for child in self.children.values() {
3158 child.collect_cursors(out);
3159 }
3160 }
3161}
3162
3163impl RouteTable {
3164 fn split(&self, path: &Path) -> (Vec<&RouteNode>, Option<&RouteNode>) {
3167 let mut above = Vec::new();
3168 let mut node = &self.root;
3169 for part in path.parts() {
3170 above.push(node);
3171 match node.children.get(part) {
3172 Some(child) => node = child,
3173 None => return (above, None),
3174 }
3175 }
3176 (above, Some(node))
3177 }
3178
3179 fn covering(&self, path: &Path) -> impl Iterator<Item = &RouteEntry> {
3181 let (above, at) = self.split(path);
3182 above.into_iter().chain(at).flat_map(|node| node.entries.iter())
3183 }
3184
3185 fn covers(&self, path: &Path, id: u64) -> bool {
3187 self.covering(path).any(|entry| entry.id == id && entry.live())
3188 }
3189
3190 fn at(&self, prefix: &Path) -> impl Iterator<Item = &RouteEntry> {
3192 self.root
3193 .find(prefix.parts())
3194 .into_iter()
3195 .flat_map(|node| node.entries.iter())
3196 }
3197
3198 fn at_mut(&mut self, prefix: &Path) -> impl Iterator<Item = &mut RouteEntry> {
3200 let mut node = Some(&mut self.root);
3201 for part in prefix.parts() {
3202 node = node.and_then(|node| node.children.get_mut(part));
3203 }
3204 node.into_iter().flat_map(|node| node.entries.iter_mut())
3205 }
3206
3207 fn entries(&self) -> impl Iterator<Item = &RouteEntry> {
3209 let mut nodes = Vec::new();
3210 self.root.walk(&mut |node| nodes.push(node));
3211 nodes.into_iter().flat_map(|node| node.entries.iter())
3212 }
3213
3214 fn insert(&mut self, entry: RouteEntry) {
3216 let node = self.root.reach(entry.prefix.parts(), Below::NONE);
3217 node.entries.push(entry);
3218 }
3219
3220 fn entry_mut(&mut self, prefix: &Path, id: u64) -> Option<&mut RouteEntry> {
3222 let mut node = &mut self.root;
3223 for part in prefix.parts() {
3224 node = node.children.get_mut(part)?;
3225 }
3226 node.entries.iter_mut().find(|entry| entry.id == id)
3227 }
3228
3229 fn remove(&mut self, prefix: &Path, id: u64) -> Option<RouteEntry> {
3231 self.root
3232 .edit(prefix.parts(), Below::NONE, |node| {
3233 let index = node.entries.iter().position(|entry| entry.id == id)?;
3234 Some(node.entries.swap_remove(index))
3235 })
3236 .flatten()
3237 }
3238
3239 fn add_cursor(&mut self, head: &Path, id: ConsumerId) {
3241 self.root.reach(head.parts(), Below::CURSOR).cursors.push(id);
3242 }
3243
3244 fn remove_cursor(&mut self, head: &Path, id: ConsumerId) {
3247 self.root.edit(head.parts(), Below::CURSOR, |node| {
3248 node.cursors.retain(|cursor| *cursor != id)
3249 });
3250 }
3251
3252 fn add_watch(&mut self, path: &Path, id: u64) -> kio::Consumer<Watched> {
3254 let producer = kio::Producer::<Watched>::default();
3255 let consumer = producer.consume();
3256 self.root.reach(path.parts(), Below::WATCH).watches.push((id, producer));
3257 consumer
3258 }
3259
3260 fn remove_watch(&mut self, path: &Path, id: u64) {
3263 self.root.edit(path.parts(), Below::WATCH, |node| {
3264 node.watches.retain(|(watch, _)| *watch != id)
3265 });
3266 }
3267
3268 fn poke_below(&self, prefix: &Path) {
3270 if let (_, Some(node)) = self.split(prefix) {
3271 node.poke_below();
3272 }
3273 }
3274
3275 fn poke_all(&self) {
3277 self.root.walk(&mut |node| node.poke());
3278 }
3279
3280 fn cursors_touching(&self, prefix: &Path) -> Vec<ConsumerId> {
3284 let (above, at) = self.split(prefix);
3285 let mut cursors: Vec<ConsumerId> = above.iter().flat_map(|node| node.cursors.iter().copied()).collect();
3286 if let Some(node) = at {
3287 node.collect_cursors(&mut cursors);
3288 }
3289 cursors.sort_unstable();
3291 cursors.dedup();
3292 cursors
3293 }
3294}
3295
3296struct OriginState {
3303 update_hold: Duration,
3306
3307 routes: RouteTable,
3310 next_route: u64,
3311 next_watch: u64,
3312
3313 cursors: HashMap<ConsumerId, TableCursor>,
3317
3318 fronts: WeakCache<FrontKey, RemoteFront>,
3326
3327 withdrawn: HashMap<PathOwned, HashSet<Hop>>,
3330
3331 closed: bool,
3334}
3335
3336impl OriginState {
3337 fn new(update_hold: Duration) -> Self {
3338 Self {
3339 update_hold,
3340 routes: RouteTable::default(),
3341 next_route: 0,
3342 next_watch: 0,
3343 cursors: HashMap::new(),
3344 fronts: WeakCache::default(),
3345 withdrawn: HashMap::new(),
3346 closed: false,
3347 }
3348 }
3349
3350 fn withdrawn_through(&self, prefix: &Path, hops: &Hops) -> bool {
3352 if self.withdrawn.is_empty() {
3353 return false;
3354 }
3355 self.withdrawn
3356 .get(prefix)
3357 .is_some_and(|peers| hops.iter().any(|hop| peers.contains(hop)))
3358 }
3359
3360 fn restale(&mut self, prefix: &Path) {
3362 let peers = self.withdrawn.get(prefix);
3363 let mut changed = None;
3364 for entry in self.routes.at_mut(prefix) {
3365 let stale = peers.is_some_and(|peers| entry.hops.iter().any(|hop| peers.contains(hop)));
3366 if entry.stale != stale {
3367 entry.stale = stale;
3368 changed = Some(entry.claim.clone());
3369 }
3370 }
3371 if let Some(claim) = changed {
3372 self.sync_route(prefix, &claim);
3373 }
3374 }
3375
3376 fn reannounced(&mut self, prefix: &PathOwned, hops: &Hops) {
3380 if let Some(sender) = hops.iter().last()
3381 && self.withdrawn.get_mut(prefix).is_some_and(|peers| peers.remove(sender))
3382 {
3383 self.restale(prefix);
3384 self.prune_withdrawn(prefix);
3385 }
3386 }
3387
3388 fn prune_withdrawn(&mut self, prefix: &PathOwned) {
3390 if self.withdrawn.is_empty() {
3391 return;
3392 }
3393 let Some(peers) = self.withdrawn.get_mut(prefix) else {
3394 return;
3395 };
3396 if peers.len() <= 1 {
3397 peers.retain(|peer| self.routes.at(prefix).any(|entry| entry.hops.contains(peer)));
3399 } else {
3400 let mut unreferenced = peers.clone();
3401 for entry in self.routes.at(prefix) {
3402 if unreferenced.is_empty() {
3403 break;
3404 }
3405 for hop in entry.hops.iter() {
3406 unreferenced.remove(hop);
3407 }
3408 }
3409 peers.retain(|peer| !unreferenced.contains(peer));
3410 }
3411 if peers.is_empty() {
3412 self.withdrawn.remove(prefix);
3413 }
3414 }
3415
3416 fn sync_route(&mut self, prefix: &Path, claim: &Pattern) {
3421 let routes = &self.routes;
3423 let mut candidates = None;
3424 for id in routes.cursors_touching(prefix) {
3425 let Some(cursor) = self.cursors.get_mut(&id) else {
3426 continue;
3427 };
3428 if let Some(presented) = cursor.presented(prefix, claim) {
3429 if presented.is_empty() {
3430 Self::sync_cursor(routes, cursor, &presented);
3432 } else {
3433 let candidates = candidates.get_or_insert_with(|| {
3436 let mut entries: Vec<_> = routes
3437 .at(prefix)
3438 .map(|entry| (route_order(prefix, entry), entry))
3439 .collect();
3440 entries.sort_unstable_by_key(|(order, _)| *order);
3441 entries
3442 });
3443 let best = candidates
3444 .iter()
3445 .map(|(_, entry)| *entry)
3446 .find(|entry| cursor.visible(entry));
3447 cursor.update(&presented, best);
3448 }
3449 }
3450 }
3451 routes.poke_below(prefix);
3453 }
3454
3455 fn watch(&mut self, shared: &kio::Shared<OriginState>, path: &Path) -> Watch {
3457 let id = self.next_watch;
3458 self.next_watch += 1;
3459 let signal = self.routes.add_watch(path, id);
3460 Watch {
3461 shared: shared.clone(),
3462 path: path.to_owned(),
3463 id,
3464 signal,
3465 }
3466 }
3467
3468 fn sync_cursor(routes: &RouteTable, cursor: &mut TableCursor, presented: &PathOwned) {
3471 let candidates: Vec<&RouteEntry> = match presented.is_empty() {
3477 true => routes
3478 .covering(&cursor.root)
3479 .filter(|entry| cursor.visible(entry))
3480 .collect(),
3481 false => {
3482 let absolute = cursor.root.join(presented);
3483 routes.at(&absolute).filter(|entry| cursor.visible(entry)).collect()
3484 }
3485 };
3486 let most = candidates.iter().map(|entry| entry.prefix.len()).max();
3487 let best = most.and_then(|most| {
3488 candidates
3489 .into_iter()
3490 .filter(|entry| entry.prefix.len() == most)
3491 .min_by_key(|entry| route_order(&entry.prefix, entry))
3492 });
3493
3494 cursor.update(presented, best);
3495 }
3496
3497 fn register_cursor(&mut self, id: ConsumerId, mut cursor: TableCursor) {
3499 let mut presented: BTreeSet<PathOwned> = BTreeSet::new();
3502 for head in &cursor.heads {
3503 let (above, at) = self.routes.split(head);
3504 let mut nodes = above;
3505 if let Some(node) = at {
3506 node.walk(&mut |node| nodes.push(node));
3507 }
3508 for entry in nodes.into_iter().flat_map(|node| node.entries.iter()) {
3509 if let Some(p) = cursor.presented(&entry.prefix, &entry.claim) {
3510 presented.insert(p);
3511 }
3512 }
3513 }
3514 for p in &presented {
3515 Self::sync_cursor(&self.routes, &mut cursor, p);
3516 }
3517 for head in &cursor.heads {
3518 self.routes.add_cursor(head, id);
3519 }
3520 self.cursors.insert(id, cursor);
3521 }
3522
3523 fn best_route(&self, path: &Path, horizon: Horizon, pin: Pin, refused: &HashSet<u64>) -> Option<&RouteEntry> {
3538 let (above, at) = self.routes.split(path);
3542 let mut best = None;
3543 for node in above.into_iter().chain(at) {
3544 let mut candidates = node
3545 .entries
3546 .iter()
3547 .filter(|entry| entry.live())
3548 .filter(|entry| entry.scope.matches(path.as_str()))
3549 .filter(|entry| horizon.admits(entry))
3550 .filter(|entry| entry.qualifies(pin))
3551 .filter(|entry| !refused.contains(&entry.id))
3552 .peekable();
3553 if candidates.peek().is_some() {
3554 best = candidates
3555 .filter(|entry| entry.serves(path))
3556 .min_by_key(|entry| route_order(&entry.prefix, entry));
3557 }
3558 }
3559 best
3560 }
3561}
3562
3563#[derive(Default)]
3570struct PendingBroadcast {
3571 resolved: Option<Result<broadcast::Consumer, Error>>,
3572}
3573
3574#[must_use = "dropping an origin::Dynamic retracts the route"]
3589pub struct Dynamic {
3590 announcement: AnnounceProducer,
3592 state: kio::Shared<ServeState>,
3593 _keepalive: Keepalive,
3597}
3598
3599impl Dynamic {
3600 pub(crate) fn withdrawn(self) {
3603 self.announcement.retract(true);
3604 }
3605
3606 pub fn update(&self, route: Route) -> Result<(), Error> {
3616 self.announcement.update(route)
3617 }
3618
3619 pub fn poll_requested_broadcast(&self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
3624 let mut state = ready!(self.state.poll(waiter, |state| {
3625 if state.closed || state.requests.has_queued() {
3626 Poll::Ready(())
3627 } else {
3628 Poll::Pending
3629 }
3630 }));
3631
3632 if state.closed {
3634 return Poll::Ready(Err(Error::Closed));
3635 }
3636
3637 let path = state.requests.pop().expect("predicate guaranteed a request");
3638 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
3644 Poll::Ready(Ok(Request {
3645 path,
3646 producer,
3647 home: self.state.clone(),
3648 }))
3649 }
3650
3651 pub async fn requested_broadcast(&self) -> Result<Request, Error> {
3657 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
3658 }
3659}
3660
3661impl ServeState {
3662 fn resolve(
3673 shared: &kio::Shared<Self>,
3674 path: &PathOwned,
3675 producer: &kio::Producer<PendingBroadcast>,
3676 result: Result<broadcast::Consumer, Error>,
3677 ) {
3678 let mut state = shared.lock();
3679 if state.closed {
3680 return;
3681 }
3682 let resolved = match result {
3683 Ok(broadcast) => {
3684 let existing = state.served.insert(path.clone(), broadcast.weak());
3688 Ok(existing.map(|weak| weak.consume()).unwrap_or(broadcast))
3689 }
3690 Err(err) => Err(err),
3691 };
3692 state.requests.remove_if(path, |p| p.same_channel(producer));
3693 if let Ok(mut pending) = producer.write() {
3694 pending.resolved.get_or_insert(resolved);
3695 drop(state);
3696 }
3697 }
3698
3699 fn forget(shared: &kio::Shared<Self>, path: &PathOwned, producer: &kio::Producer<PendingBroadcast>) {
3701 shared.lock().requests.remove_if(path, |p| p.same_channel(producer));
3702 }
3703}
3704
3705pub struct Request {
3712 path: PathOwned,
3714
3715 producer: kio::Producer<PendingBroadcast>,
3718
3719 home: kio::Shared<ServeState>,
3722}
3723
3724impl Request {
3725 pub fn path(&self) -> &Path<'_> {
3727 &self.path
3728 }
3729
3730 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
3736 let broadcast = broadcast.consume();
3737 ServeState::resolve(&self.home, &self.path, &self.producer, Ok(broadcast));
3738 }
3740
3741 pub fn reject(self, err: Error) {
3743 ServeState::resolve(&self.home, &self.path, &self.producer, Err(err));
3744 }
3745}
3746
3747impl Drop for Request {
3748 fn drop(&mut self) {
3749 ServeState::forget(&self.home, &self.path, &self.producer);
3757 }
3758}
3759
3760pub struct Requesting {
3767 inner: RequestState,
3768 path: PathOwned,
3772 stats: stats::Scope,
3775}
3776
3777enum RequestState {
3778 Failed(Error),
3781 Pending(kio::Consumer<PendingBroadcast>),
3783}
3784
3785impl Requesting {
3786 fn failed(error: Error) -> Self {
3787 Self::new(RequestState::Failed(error))
3788 }
3789
3790 fn queued(consumer: kio::Consumer<PendingBroadcast>) -> Self {
3791 Self::new(RequestState::Pending(consumer))
3792 }
3793
3794 pub fn is_queued(&self) -> bool {
3803 matches!(self.inner, RequestState::Pending(_))
3804 }
3805
3806 fn new(inner: RequestState) -> Self {
3807 Self {
3808 inner,
3809 path: PathOwned::default(),
3810 stats: stats::Scope::default(),
3811 }
3812 }
3813
3814 fn with_path(mut self, path: PathOwned) -> Self {
3815 self.path = path;
3816 self
3817 }
3818
3819 fn with_stats(mut self, scope: stats::Scope) -> Self {
3821 self.stats = scope;
3822 self
3823 }
3824
3825 fn hand_out(&self, broadcast: broadcast::Consumer) -> broadcast::Consumer {
3827 broadcast.with_path(self.path.clone()).with_stats(self.stats.clone())
3828 }
3829
3830 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
3832 match &self.inner {
3833 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
3834 RequestState::Pending(consumer) => Poll::Ready(
3835 match ready!(consumer.poll(waiter, |state| match &state.resolved {
3836 Some(result) => Poll::Ready(result.clone()),
3837 None => Poll::Pending,
3838 })) {
3839 Ok(result) => result.map(|broadcast| self.hand_out(broadcast)),
3840 Err(_closed) => Err(Error::Unroutable),
3842 },
3843 ),
3844 }
3845 }
3846}
3847
3848impl kio::Pollable for Requesting {
3849 type Output = Result<broadcast::Consumer, Error>;
3850
3851 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
3852 self.poll_ok(waiter)
3853 }
3854}
3855
3856pub trait Consume<T> {
3864 fn consume(&self) -> T;
3866}
3867
3868impl<T, U: Consume<T>> Consume<T> for &U {
3869 fn consume(&self) -> T {
3870 (**self).consume()
3871 }
3872}
3873
3874impl Consume<Consumer> for Producer {
3875 fn consume(&self) -> Consumer {
3876 Consumer::from_producer(self, stats::Session::default())
3880 }
3881}
3882
3883impl Consume<Consumer> for Consumer {
3884 fn consume(&self) -> Consumer {
3885 self.clone()
3886 }
3887}
3888
3889impl Consume<broadcast::Consumer> for broadcast::Producer {
3890 fn consume(&self) -> broadcast::Consumer {
3891 self.consume()
3893 }
3894}
3895
3896impl Consume<broadcast::Consumer> for broadcast::Consumer {
3897 fn consume(&self) -> broadcast::Consumer {
3898 self.clone()
3899 }
3900}
3901
3902impl Consume<track::Consumer> for track::Producer {
3903 fn consume(&self) -> track::Consumer {
3904 self.consume()
3905 }
3906}
3907
3908impl Consume<track::Consumer> for track::Consumer {
3909 fn consume(&self) -> track::Consumer {
3910 self.clone()
3911 }
3912}
3913
3914#[derive(Clone)]
3920pub struct Consumer {
3921 hop: Hop,
3923 scope: OriginScope,
3924
3925 root: PathOwned,
3927
3928 shared: kio::Shared<OriginState>,
3931
3932 stats: stats::Session,
3936
3937 horizon: Horizon,
3942
3943 hidden: Hidden,
3945
3946 pool: cache::Pool,
3949 cache_duration: Duration,
3950
3951 tasks: TasksWeak,
3955
3956 timers: Clock,
3959}
3960
3961impl Consumer {
3962 fn from_producer(producer: &Producer, stats: stats::Session) -> Self {
3963 Self {
3964 hop: producer.hop,
3965 scope: producer.scope.clone(),
3966 root: producer.root.clone(),
3967 shared: producer.shared.clone(),
3968 stats,
3969 horizon: Horizon::default(),
3970 hidden: Hidden::default(),
3971 pool: producer.pool.clone(),
3972 cache_duration: producer.cache_duration,
3973 tasks: producer.tasks.downgrade(),
3974 timers: producer.timers.clone(),
3975 }
3976 }
3977
3978 pub fn hop(&self) -> Hop {
3980 self.hop
3981 }
3982
3983 pub(crate) fn excluding(mut self, peer: Hop) -> Self {
3991 self.horizon.exclude = Some(peer);
3992 self
3993 }
3994
3995 pub fn local(mut self) -> Self {
4002 self.horizon.local = true;
4003 self
4004 }
4005
4006 pub fn with_hidden(mut self, hidden: bool) -> Self {
4011 self.hidden.include = hidden;
4012 self
4013 }
4014
4015 pub(crate) fn beyond(mut self, outer: &Consumer) -> Self {
4018 self.hidden.beyond = Some(
4019 outer
4020 .hidden
4021 .from
4022 .clone()
4023 .map(|head| vec![head])
4024 .unwrap_or_else(|| interest_prefixes(&outer.scope.allowed)),
4025 );
4026 self
4027 }
4028
4029 pub(crate) fn discovery(mut self, hidden: bool) -> Self {
4031 self.hidden.include = hidden;
4032 self.hidden.from = Some(self.root.clone());
4033 self
4034 }
4035
4036 pub(crate) fn includes_hidden(&self) -> bool {
4038 self.hidden.include
4039 }
4040
4041 pub fn with_stats(mut self, session: stats::Session) -> Self {
4045 self.stats = session;
4046 self
4047 }
4048
4049 fn untagged(&self) -> Self {
4053 Self {
4054 stats: stats::Session::default(),
4055 ..self.clone()
4056 }
4057 }
4058
4059 pub(crate) fn empty(&self) -> Self {
4064 Self {
4065 scope: OriginScope::empty(),
4066 ..self.clone()
4067 }
4068 }
4069
4070 pub fn announced(&self) -> AnnounceConsumer {
4079 let state = kio::Producer::<OriginConsumerState>::default();
4080 let cursor = |root: PathOwned,
4081 allowed: Patterns,
4082 hidden: Hidden,
4083 mount: Option<&Mount>,
4084 under: PathOwned,
4085 holes: Vec<PathOwned>| TableCursor {
4086 root,
4087 heads: interest_prefixes(&allowed),
4088 allowed,
4089 horizon: self.horizon,
4090 hidden,
4091 state: state.clone(),
4092 under,
4093 holes,
4094 named: mount.map(|mount| (mount.clone(), self.scope.allowed.clone())),
4095 current: HashMap::new(),
4096 };
4097
4098 if let Some(mount) = self.scope.mount(&self.root) {
4101 let Some(root) = mount.resolve(&self.root) else {
4103 return AnnounceConsumer::new(
4104 self.root.clone(),
4105 Vec::new(),
4106 state,
4107 self.stats.clone(),
4108 &self.shared,
4109 self.timers.clone(),
4110 );
4111 };
4112 let cursors = vec![cursor(
4113 root,
4114 mount.translate(&self.scope.allowed),
4115 self.hidden.translate(mount),
4116 Some(mount),
4117 PathOwned::default(),
4118 Vec::new(),
4119 )];
4120 return AnnounceConsumer::new(
4121 self.root.clone(),
4122 cursors,
4123 state,
4124 self.stats.clone(),
4125 &self.shared,
4126 self.timers.clone(),
4127 );
4128 }
4129
4130 let heads = interest_prefixes(&self.scope.allowed);
4133 let mut holes = Vec::new();
4134 let mut cursors = Vec::new();
4135 for mount in self.scope.mounts.iter() {
4136 let Some(under) = mount.at.strip_prefix(&self.root) else {
4137 continue;
4138 };
4139 holes.push(mount.at.clone());
4140 let allowed = mount.translate(&self.scope.allowed);
4141 if allowed.is_empty()
4144 || !(self.hidden.include
4145 || !hides(
4146 self.hidden.from.as_ref().map(std::slice::from_ref).unwrap_or(&heads),
4147 &mount.at,
4148 )) {
4149 continue;
4150 }
4151 cursors.push(cursor(
4152 mount.target.clone(),
4153 allowed,
4154 self.hidden.translate(mount),
4155 Some(mount),
4156 under.to_owned(),
4157 Vec::new(),
4158 ));
4159 }
4160 cursors.insert(
4161 0,
4162 cursor(
4163 self.root.clone(),
4164 self.scope.allowed.clone(),
4165 self.hidden.clone(),
4166 None,
4167 PathOwned::default(),
4168 holes,
4169 ),
4170 );
4171 AnnounceConsumer::new(
4172 self.root.clone(),
4173 cursors,
4174 state,
4175 self.stats.clone(),
4176 &self.shared,
4177 self.timers.clone(),
4178 )
4179 }
4180
4181 pub fn consume(&self) -> Self {
4183 self.clone()
4184 }
4185
4186 #[cfg(test)]
4189 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
4190 let full = self.root.join(path).to_owned();
4191 if !self.scope.permits(&full) {
4192 return None;
4193 }
4194 let full = self.scope.resolve(&full)?;
4195 let table = self.shared.lock();
4196 table
4197 .routes
4198 .at(&full)
4199 .filter(|entry| entry.local)
4200 .min_by_key(|entry| route_order(&entry.prefix, entry))
4201 .and_then(|entry| entry.source.clone())
4202 }
4203
4204 pub async fn routed(&self, path: impl AsPath) -> Option<Route> {
4214 let path = path.as_path();
4215
4216 let consumer = match Pattern::subtree(path.as_str()) {
4221 Ok(subtree) => self.scope("", &Patterns::from(subtree)).ok()?,
4222 Err(InvalidPattern::TooManySegments) => self.clone(),
4223 Err(_) => return None,
4224 };
4225
4226 if !consumer.allowed().matches(path.as_str()) {
4230 return None;
4231 }
4232
4233 let mut announced = consumer.untagged().with_hidden(true).announced();
4237 loop {
4238 let update = announced.next().await?;
4239 if update.kind.is_active() && path.has_prefix(&update.prefix) {
4240 return Some(update.route);
4241 }
4242 }
4243 }
4244
4245 pub async fn routed_broadcast(&self, path: impl AsPath) -> Result<broadcast::Consumer, Error> {
4259 let path = path.as_path();
4260
4261 if !self.allowed().matches(path.as_str()) {
4264 return Err(Error::Unauthorized);
4265 }
4266 loop {
4267 let (watch, seen) = {
4276 let mut table = self.shared.lock();
4277 if table.closed {
4278 return Err(Error::Closed);
4279 }
4280 let named = self.root.join(&path);
4281 let resolved = self.scope.resolve(&named).ok_or(BoundsExceeded)?;
4282 let watch = table.watch(&self.shared, &resolved);
4283 let seen = watch.seen();
4284 (watch, seen)
4285 };
4286 match self.request_broadcast(&path).await {
4287 Ok(broadcast) => return Ok(broadcast),
4288 Err(Error::Unroutable) => {
4289 kio::wait(|waiter| watch.poll_changed(waiter, seen)).await;
4290 }
4291 Err(Error::Dropped) if self.shared.lock().closed => return Err(Error::Closed),
4294 Err(err) => return Err(err),
4295 }
4296 }
4297 }
4298
4299 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Consumer, Error> {
4306 let root = self.root.join(root).to_owned();
4307 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
4308 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
4309 Ok(Consumer {
4310 scope,
4311 root,
4312 ..self.clone()
4313 })
4314 }
4315
4316 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
4339 let path = path.as_path();
4340
4341 let named = self.root.join(&path).to_owned();
4345 let scope = self.stats.egress(&named);
4346 let requested = path.to_owned();
4350
4351 if !self.scope.permits(&named) {
4353 return kio::Pending::new(Requesting::failed(Error::Unauthorized));
4354 }
4355
4356 let Some(absolute) = self.scope.resolve(&named).map(|path| path.to_owned()) else {
4359 return kio::Pending::new(Requesting::failed(BoundsExceeded.into()));
4360 };
4361
4362 let mut state = self.shared.lock();
4363
4364 if state.closed {
4366 return kio::Pending::new(Requesting::failed(Error::Closed));
4367 }
4368
4369 if state
4373 .best_route(&absolute.as_path(), self.horizon, Pin::Any, &HashSet::new())
4374 .is_none()
4375 {
4376 return kio::Pending::new(Requesting::failed(Error::Unroutable));
4377 }
4378
4379 let key = (absolute.clone(), self.horizon);
4387 if let Some(front) = state.fronts.get(&key) {
4388 let pin = *front.pin.lock();
4389 let current = state
4390 .best_route(&absolute.as_path(), self.horizon, Pin::Any, &HashSet::new())
4391 .is_some_and(|entry| entry.qualifies(pin));
4392 if current {
4393 let pending = Requesting::queued(front.request.consume())
4394 .with_path(requested)
4395 .with_stats(scope);
4396 return kio::Pending::new(pending);
4397 }
4398 state.fronts.remove(&key);
4399 }
4400
4401 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
4406 pool: self.pool.clone(),
4407 cache_duration: self.cache_duration,
4408 path: absolute.clone(),
4409 });
4410 let request = kio::Producer::<PendingBroadcast>::default();
4411 let consumer = request.consume();
4412 let watch = state.watch(&self.shared, &absolute);
4413 let pin = kio::Lock::new(Pin::Any);
4414 state.fronts.insert(
4415 key,
4416 RemoteFront {
4417 request: request.clone(),
4418 broadcast: broadcast.consume().weak(),
4419 pin: pin.clone(),
4420 },
4421 );
4422 drop(state);
4425 self.tasks.push(run_front(FrontTask {
4426 shared: self.shared.clone(),
4427 broadcast,
4428 path: absolute,
4429 horizon: self.horizon,
4430 watch,
4431 request,
4432 pin,
4433 timers: self.timers.clone(),
4434 }));
4435 kio::Pending::new(Requesting::queued(consumer).with_path(requested).with_stats(scope))
4436 }
4437
4438 pub fn root(&self) -> &Path<'_> {
4440 &self.root
4441 }
4442
4443 pub fn allowed(&self) -> Patterns {
4445 self.scope.relative(&self.root)
4446 }
4447
4448 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
4450 self.root.join(path)
4451 }
4452}
4453
4454pub struct AnnounceConsumer {
4459 ids: Vec<ConsumerId>,
4462 shared: kio::Shared<OriginState>,
4463 root: PathOwned,
4464
4465 state: kio::Producer<OriginConsumerState>,
4468
4469 stats: stats::Session,
4472
4473 guards: HashMap<PathOwned, stats::Announce>,
4477
4478 park: kio::Park,
4481
4482 update_hold: Duration,
4485 held: HashMap<PathOwned, Held>,
4486 timers: Clock,
4487 hold: crate::runtime::Deadline<Clock>,
4488}
4489
4490impl AnnounceConsumer {
4491 fn new(
4492 root: PathOwned,
4493 cursors: Vec<TableCursor>,
4494 state: kio::Producer<OriginConsumerState>,
4495 stats: stats::Session,
4496 shared: &kio::Shared<OriginState>,
4497 timers: Clock,
4498 ) -> Self {
4499 let mut ids = Vec::with_capacity(cursors.len());
4500 let update_hold;
4501 {
4502 let mut table = shared.lock();
4503 update_hold = table.update_hold;
4504 if table.closed {
4505 if let Ok(mut state) = state.write() {
4507 state.ended = true;
4508 }
4509 } else {
4510 for cursor in cursors {
4511 let id = ConsumerId::new();
4512 table.register_cursor(id, cursor);
4513 ids.push(id);
4514 }
4515 }
4516 }
4517
4518 Self {
4519 ids,
4520 shared: shared.clone(),
4521 root,
4522 state,
4523 stats,
4524 guards: HashMap::new(),
4525 park: kio::Park::default(),
4526 update_hold,
4527 held: HashMap::new(),
4528 hold: crate::runtime::Deadline::new(&timers),
4529 timers,
4530 }
4531 }
4532
4533 fn hand_out(&mut self, update: AnnounceUpdate) -> AnnounceUpdate {
4535 let absolute = self.root.join(&update.prefix).to_owned();
4536 if update.kind.is_active() {
4537 let scope = self.stats.egress(&absolute);
4538 self.guards
4539 .entry(update.prefix.clone())
4540 .or_insert_with(|| scope.announce());
4541 } else {
4542 self.guards.remove(&update.prefix);
4543 }
4544 update
4545 }
4546
4547 pub async fn next(&mut self) -> Option<AnnounceUpdate> {
4555 kio::wait(|waiter| self.poll_next(waiter)).await
4556 }
4557
4558 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<AnnounceUpdate>> {
4564 loop {
4565 let now = self.timers.try_now();
4566 let hold = self.update_hold;
4567 let held = &mut self.held;
4568 let mut ready = None;
4569 let mut wake = None;
4570 let update = match self.state.poll(waiter, |state| {
4571 if state.pending.is_empty() {
4572 return match state.ended {
4573 true => Poll::Ready(()),
4574 false => Poll::Pending,
4575 };
4576 }
4577 match state.scan(held, now, hold) {
4578 Ok(prefix) => {
4579 ready = Some(prefix);
4580 Poll::Ready(())
4581 }
4582 Err(at) => {
4583 wake = at;
4584 Poll::Pending
4585 }
4586 }
4587 }) {
4588 Poll::Ready(Ok(mut state)) => match ready {
4589 Some(prefix) => state.take_prefix(prefix),
4590 None => {
4591 state.close();
4594 None
4595 }
4596 },
4597 Poll::Ready(Err(_)) => None,
4599 Poll::Pending => {
4600 self.hold.set(wake);
4601 ready!(self.hold.poll(waiter));
4602 continue;
4603 }
4604 };
4605 let Some(update) = update else {
4606 return Poll::Ready(None);
4607 };
4608 self.held.remove(&update.prefix);
4609 return Poll::Ready(Some(self.hand_out(update)));
4610 }
4611 }
4612
4613 pub fn try_next(&mut self) -> Option<AnnounceUpdate> {
4618 let now = self.timers.try_now();
4619 let mut state = self.state.write().ok()?;
4620 let prefix = match state.scan(&mut self.held, now, self.update_hold) {
4621 Ok(prefix) => prefix,
4622 Err(wake) => {
4623 drop(state);
4626 self.hold.set(wake);
4627 return None;
4628 }
4629 };
4630 self.held.remove(&prefix);
4631 let update = state.take_prefix(prefix)?;
4632 drop(state);
4633 Some(self.hand_out(update))
4634 }
4635
4636 pub fn is_closed(&self) -> bool {
4638 let state = self.state.read();
4639 state.is_closed() || state.ended
4640 }
4641
4642 pub fn root(&self) -> &Path<'_> {
4644 &self.root
4645 }
4646
4647 pub fn absolute(&self, prefix: impl AsPath) -> Path<'_> {
4649 self.root.join(prefix)
4650 }
4651}
4652
4653impl futures::Stream for AnnounceConsumer {
4654 type Item = AnnounceUpdate;
4655
4656 fn poll_next(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> {
4657 let this = self.get_mut();
4658 let waiter = this.park.hold(cx).clone();
4659 this.poll_next(&waiter)
4660 }
4661}
4662
4663impl Drop for AnnounceConsumer {
4664 fn drop(&mut self) {
4665 let mut shared = self.shared.lock();
4666 for id in &self.ids {
4667 if let Some(cursor) = shared.cursors.remove(id) {
4668 for head in &cursor.heads {
4669 shared.routes.remove_cursor(head, *id);
4670 }
4671 }
4672 }
4673 }
4674}
4675
4676#[cfg(test)]
4677use futures::FutureExt;
4678
4679#[cfg(test)]
4680#[allow(missing_docs)] impl AnnounceConsumer {
4682 pub fn assert_next_active(&mut self, expected: impl AsPath) -> Route {
4684 let expected = expected.as_path();
4685 let update = self.next().now_or_never().expect("next blocked").expect("no next");
4686 assert_eq!(update.prefix, expected, "wrong prefix");
4687 assert!(update.kind.is_active(), "should be an active route");
4688 update.route
4689 }
4690
4691 pub fn assert_try_next_active(&mut self, expected: impl AsPath) -> Route {
4693 let expected = expected.as_path();
4694 let update = self.try_next().expect("no next");
4695 assert_eq!(update.prefix, expected, "wrong prefix");
4696 assert!(update.kind.is_active(), "should be an active route");
4697 update.route
4698 }
4699
4700 pub fn assert_next_ended(&mut self, expected: impl AsPath) {
4702 let expected = expected.as_path();
4703 let update = self.next().now_or_never().expect("next blocked").expect("no next");
4704 assert_eq!(update.prefix, expected, "wrong prefix");
4705 assert_eq!(update.kind, AnnounceKind::Retracted, "should be a retraction");
4706 }
4707
4708 pub fn assert_next_wait(&mut self) {
4709 if let Some(res) = self.next().now_or_never() {
4710 panic!("next should block: got {:?}", res.map(|u| u.prefix));
4711 }
4712 }
4713}
4714
4715#[cfg(test)]
4719pub(crate) trait ProduceTest {
4720 fn produce(self) -> Producer;
4721}
4722
4723#[cfg(test)]
4724impl ProduceTest for Config {
4725 fn produce(self) -> Producer {
4726 let (producer, driver) = Producer::new(self);
4727 if tokio::runtime::Handle::try_current().is_ok() {
4728 tokio::spawn(crate::time::run(driver));
4729 } else {
4730 std::mem::forget(driver);
4733 }
4734 producer
4735 }
4736}
4737
4738#[cfg(test)]
4739impl ProduceTest for Hop {
4740 fn produce(self) -> Producer {
4741 Config::new(self).produce()
4742 }
4743}
4744
4745#[cfg(test)]
4746mod tests {
4747 use super::*;
4748 use futures::FutureExt;
4749
4750 fn origin(id: u64) -> Hop {
4751 Hop::new(id).unwrap()
4752 }
4753
4754 fn hops(ids: &[u64]) -> Hops {
4755 let mut list = Hops::new();
4756 for &id in ids {
4757 list.push(if id == 0 { Hop::UNKNOWN } else { origin(id) }).unwrap();
4758 }
4759 list
4760 }
4761
4762 fn scopes(prefixes: &[&str]) -> Patterns {
4764 prefixes
4765 .iter()
4766 .map(|prefix| Pattern::subtree(prefix).unwrap())
4767 .collect()
4768 }
4769
4770 #[test]
4771 fn default_config_mints_a_real_hop() {
4772 let config = Config::default();
4773 assert_ne!(config.hop, Hop::UNKNOWN);
4774 let (producer, _driver) = Producer::new(config.clone());
4775 assert_eq!(producer.hop(), config.hop);
4776 assert_eq!(producer.consume().hop(), config.hop);
4777 }
4778
4779 #[test]
4780 fn random_hops_fit_legacy_lite_clients() {
4781 for _ in 0..32 {
4782 assert!(Hop::random().id() < 1u64 << 53);
4783 }
4784 }
4785
4786 async fn settle(mut check: impl FnMut() -> bool) {
4789 for _ in 0..100 {
4790 if check() {
4791 return;
4792 }
4793 tokio::task::yield_now().await;
4794 }
4795 panic!("condition never settled");
4796 }
4797
4798 async fn queued(server: &Dynamic) -> Request {
4800 let mut request = None;
4801 settle(|| match server.poll_requested_broadcast(&kio::Waiter::noop()) {
4802 Poll::Ready(Ok(popped)) => {
4803 request = Some(popped);
4804 true
4805 }
4806 _ => false,
4807 })
4808 .await;
4809 request.unwrap()
4810 }
4811
4812 async fn next_group(subscription: &mut crate::track::Subscriber) -> Result<Option<crate::group::Consumer>, Error> {
4814 let mut next = None;
4815 settle(|| match subscription.poll_recv_group(&kio::Waiter::noop()) {
4816 Poll::Ready(result) => {
4817 next = Some(result);
4818 true
4819 }
4820 Poll::Pending => false,
4821 })
4822 .await;
4823 next.unwrap()
4824 }
4825
4826 #[tokio::test]
4827 async fn announce_and_retract() {
4828 let producer = origin(1).produce();
4829 let consumer = producer.consume();
4830 let mut announced = consumer.announced();
4831 announced.assert_next_wait();
4832
4833 let announcement = producer.announce("room/alice", Route::default()).unwrap();
4834 let route = announced.assert_next_active("room/alice");
4835 assert!(route.hops.is_empty());
4836 assert_eq!(route.cost, Cost::default());
4837 announced.assert_next_wait();
4838
4839 drop(announcement);
4840 announced.assert_next_ended("room/alice");
4841 announced.assert_next_wait();
4842 }
4843
4844 #[tokio::test]
4847 async fn hidden_routes_need_an_opt_in() {
4848 let producer = origin(1).produce();
4849 let consumer = producer.consume();
4850 let _visible = producer.announce("room/alice", Route::default()).unwrap();
4851 let _stats = producer.announce(".stats/node", Route::default()).unwrap();
4852 let _nested = producer.announce("room/.internal", Route::default()).unwrap();
4853 let _suffix = producer.announce("room/catalog.pro", Route::default()).unwrap();
4855
4856 let mut announced = consumer.announced();
4857 announced.assert_next_active("room/alice");
4858 announced.assert_next_active("room/catalog.pro");
4859 announced.assert_next_wait();
4860
4861 let mut announced = consumer.clone().with_hidden(true).announced();
4862 announced.assert_next_active(".stats/node");
4863 announced.assert_next_active("room/.internal");
4864 announced.assert_next_active("room/alice");
4865 announced.assert_next_active("room/catalog.pro");
4866 announced.assert_next_wait();
4867
4868 let mut announced = consumer
4870 .scope(".stats", &Patterns::from(Pattern::all()))
4871 .unwrap()
4872 .announced();
4873 announced.assert_next_active("node");
4874 announced.assert_next_wait();
4875 let mut announced = consumer.scope("", &scopes(&["room/.internal"])).unwrap().announced();
4876 announced.assert_next_active("room/.internal");
4877 announced.assert_next_wait();
4878
4879 let mut announced = consumer.clone().with_hidden(true).beyond(&consumer).announced();
4882 announced.assert_next_active(".stats/node");
4883 announced.assert_next_active("room/.internal");
4884 announced.assert_next_wait();
4885 let room = consumer.scope("", &scopes(&["room"])).unwrap().beyond(&consumer);
4886 let mut announced = room.announced();
4887 announced.assert_next_wait();
4888 let stats = consumer.scope("", &scopes(&[".stats"])).unwrap().beyond(&consumer);
4889 let mut announced = stats.announced();
4890 announced.assert_next_active(".stats/node");
4891 announced.assert_next_wait();
4892 }
4893
4894 #[tokio::test]
4896 async fn hidden_broadcast_resolves_by_path() {
4897 let producer = origin(1).produce();
4898 let consumer = producer.consume();
4899 let broadcast = producer.create_broadcast(".stats/node").unwrap();
4900 broadcast.announce(Route::default()).unwrap();
4901
4902 consumer.announced().assert_next_wait();
4903 let resolved = consumer.request_broadcast(".stats/node").await.expect("resolves");
4904 assert_eq!(resolved.info().path.as_str(), ".stats/node");
4905 }
4906
4907 fn mounted(producer: &Producer, pid: &str, patterns: &[&str]) -> Consumer {
4909 let patterns: Patterns = patterns.iter().map(|pattern| pattern.parse().unwrap()).collect();
4910 producer
4911 .mount(format!("{pid}/.svc"), format!(".svc/{pid}"))
4912 .unwrap()
4913 .scope(pid, &patterns)
4914 .unwrap()
4915 .consume()
4916 }
4917
4918 #[tokio::test]
4921 async fn mount_resolves_on_the_target_front() {
4922 let producer = origin(1).produce();
4923 let server = producer.dynamic(".svc", Route::default()).unwrap();
4924 let project = mounted(&producer, "p1", &["**"]);
4925
4926 let through = project.request_broadcast(".svc/foo");
4927 let direct = producer.consume().request_broadcast(".svc/p1/foo");
4928
4929 let request = queued(&server).await;
4930 assert_eq!(request.path().as_str(), ".svc/p1/foo");
4931 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
4932 let source = broadcast::Info::new().produce();
4933 request.accept(&source);
4934
4935 let through = through.await.expect("resolves");
4936 let direct = direct.await.expect("resolves");
4937 assert!(through.is_clone(&direct));
4938 assert_eq!(through.info().path.as_str(), ".svc/foo");
4940 }
4941
4942 #[tokio::test]
4945 async fn mount_presents_target_routes_under_the_mount() {
4946 let producer = origin(1).produce();
4947 let _claim = producer.announce(".svc", Route::default()).unwrap();
4948 let foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4949 let _other = producer.publish(".svc/p2/bar", Route::default()).unwrap();
4950 let _cam = producer.publish("p1/cam", Route::default()).unwrap();
4951 let _shadowed = producer.publish("p1/.svc/forged", Route::default()).unwrap();
4952 let project = mounted(&producer, "p1", &["**"]);
4953
4954 let mut announced = project.announced();
4956 announced.assert_next_active("cam");
4957 announced.assert_next_wait();
4958
4959 let mut announced = project.clone().with_hidden(true).announced();
4960 announced.assert_next_active(".svc");
4961 announced.assert_next_active(".svc/foo");
4962 announced.assert_next_active("cam");
4963 announced.assert_next_wait();
4964
4965 let mut inside = project
4967 .scope(".svc", &Patterns::from(Pattern::all()))
4968 .unwrap()
4969 .announced();
4970 inside.assert_next_active("");
4971 inside.assert_next_active("foo");
4972 inside.assert_next_wait();
4973
4974 drop(foo);
4975 announced.assert_next_ended(".svc/foo");
4976 inside.assert_next_ended("foo");
4977 announced.assert_next_wait();
4978
4979 let err = project.request_broadcast(".svc/forged").await.err().unwrap();
4981 assert!(matches!(err, Error::NotFound | Error::Unroutable), "{err:?}");
4982 }
4983
4984 #[tokio::test]
4987 async fn mount_authorizes_the_named_path() {
4988 let producer = origin(1).produce();
4989 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4990 let _bar = producer.publish(".svc/p1/bar", Route::default()).unwrap();
4991
4992 let granted = mounted(&producer, "p1", &["foo", ".svc/foo"]);
4993 granted.request_broadcast(".svc/foo").await.expect("granted");
4994 let refused = granted
4995 .request_broadcast(".svc/bar")
4996 .now_or_never()
4997 .expect("refused at once");
4998 assert!(matches!(refused, Err(Error::Unauthorized)));
4999 let mut announced = granted.clone().with_hidden(true).announced();
5000 announced.assert_next_active(".svc/foo");
5001 announced.assert_next_wait();
5002
5003 let narrower = mounted(&producer, "p1", &["foo"]);
5004 let refused = narrower
5005 .request_broadcast(".svc/foo")
5006 .now_or_never()
5007 .expect("refused at once");
5008 assert!(matches!(refused, Err(Error::Unauthorized)));
5009 narrower.clone().with_hidden(true).announced().assert_next_wait();
5010
5011 let other = mounted(&producer, "p2", &["**"]);
5013 other.clone().with_hidden(true).announced().assert_next_wait();
5014 let err = other.request_broadcast(".svc/foo").await.err().unwrap();
5015 assert!(matches!(err, Error::Unroutable));
5016 }
5017
5018 #[tokio::test]
5021 async fn mount_captures_the_named_path() {
5022 let producer = origin(1).produce();
5023 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
5024 let project = mounted(&producer, "p1", &["**"]).with_hidden(true);
5025
5026 let update = project.announced().try_next().expect("foo");
5027 assert_eq!(update.prefix.as_str(), ".svc/foo");
5028 assert_eq!(update.captures, Some(vec![".svc/foo".parse::<Pattern>().unwrap()]));
5029
5030 let inside = project.scope(".svc", &Patterns::from(Pattern::all())).unwrap();
5031 let update = inside.announced().try_next().expect("foo");
5032 assert_eq!(update.prefix.as_str(), "foo");
5033 assert_eq!(update.captures, Some(vec!["foo".parse::<Pattern>().unwrap()]));
5034 }
5035
5036 #[tokio::test]
5039 async fn mount_keeps_a_max_depth_target() {
5040 let producer = origin(1).produce();
5041 let deep = vec!["d"; Path::MAX_PARTS].join("/");
5042 let _leaf = producer.publish(deep.as_str(), Route::default()).unwrap();
5043 let project = producer
5044 .mount("p1/.svc", deep.as_str())
5045 .unwrap()
5046 .scope("p1", &Patterns::from(Pattern::all()))
5047 .unwrap()
5048 .consume()
5049 .with_hidden(true);
5050
5051 let mut announced = project.announced();
5052 announced.assert_next_active(".svc");
5053 announced.assert_next_wait();
5054
5055 let err = project.request_broadcast(".svc/x").await.err().unwrap();
5057 assert!(matches!(err, Error::BoundsExceeded(_)), "{err:?}");
5058 let inside = project.scope(".svc/x", &Patterns::from(Pattern::all())).unwrap();
5059 inside.announced().assert_next_wait();
5060 }
5061
5062 #[tokio::test]
5065 async fn mount_bounds_the_named_path() {
5066 let producer = origin(1).produce();
5067 let _near = producer.publish("t/x", Route::default()).unwrap();
5068 let deep = format!("t/{}", vec!["d"; Path::MAX_PARTS - 1].join("/"));
5069 let _deep = producer.publish(deep.as_str(), Route::default()).unwrap();
5070 let project = producer
5071 .mount("p1/a/b", "t")
5072 .unwrap()
5073 .scope("p1", &Patterns::from(Pattern::all()))
5074 .unwrap()
5075 .consume();
5076
5077 let mut announced = project.announced();
5078 announced.assert_next_active("a/b/x");
5079 announced.assert_next_wait();
5080
5081 let over = vec!["d"; Path::MAX_PARTS + 1].join("/");
5082 assert!(matches!(
5083 producer.mount(over.as_str(), "t"),
5084 Err(Error::BoundsExceeded(_))
5085 ));
5086 }
5087
5088 #[tokio::test]
5090 async fn mount_is_read_only() {
5091 let producer = origin(1).produce();
5092 let project = producer
5093 .mount("p1/.svc", ".svc/p1")
5094 .unwrap()
5095 .scope("p1", &Patterns::from(Pattern::all()))
5096 .unwrap();
5097 assert!(matches!(project.create_broadcast(".svc/foo"), Err(Error::Unauthorized)));
5098 assert!(matches!(
5099 project.dynamic(".svc", Route::default()),
5100 Err(Error::Unauthorized)
5101 ));
5102 assert!(matches!(
5103 project.dynamic(".svc/foo", Route::default()),
5104 Err(Error::Unauthorized)
5105 ));
5106 project.publish("cam", Route::default()).unwrap();
5107 }
5108
5109 #[test]
5111 fn mount_never_widens_a_scope() {
5112 let (producer, _driver) = Producer::new(Config::new(origin(1)));
5113 let project = producer.scope("", &scopes(&["p1"])).unwrap();
5114 assert!(matches!(project.mount("p1/.svc", ".svc/p1"), Err(Error::Unauthorized)));
5115
5116 let mounted = producer.mount("p1/.svc", ".svc/p1").unwrap();
5117 assert!(matches!(mounted.mount("p1/.svc/x", ".other"), Err(Error::Duplicate)));
5118 assert!(matches!(mounted.mount("p1", ".other"), Err(Error::Duplicate)));
5119 assert!(matches!(mounted.mount("p2", "p1/.svc/x"), Err(Error::Duplicate)));
5121 assert!(matches!(mounted.mount("p2", "p1"), Err(Error::Duplicate)));
5122 assert!(matches!(mounted.mount(".svc/p1/x", ".other"), Err(Error::Duplicate)));
5124 assert!(matches!(mounted.mount(".svc", ".other"), Err(Error::Duplicate)));
5125 mounted.mount("p1/.other", ".other/p1").unwrap();
5126 mounted.mount("p2/.svc", ".svc/p1").unwrap();
5128 for (at, target) in [("a", "a/b"), ("a/b", "a"), ("a", "a")] {
5130 assert!(
5131 matches!(producer.mount(at, target), Err(Error::Duplicate)),
5132 "{at} -> {target}"
5133 );
5134 }
5135 }
5136
5137 #[test]
5139 fn mount_refuses_a_wildcard_mount_point() {
5140 let (producer, _driver) = Producer::new(Config::new(origin(1)));
5141 for at in ["*", "p1/*", "p1/**", "p1/a*"] {
5142 assert!(
5143 matches!(producer.mount(at, ".svc/p1"), Err(Error::InvalidPath(_))),
5144 "{at}"
5145 );
5146 }
5147 }
5148
5149 #[tokio::test]
5152 async fn mount_egress_counts_under_the_named_path() {
5153 let registry = stats::Registry::new(stats::Config::new());
5154 let producer = origin(1).produce();
5155 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
5156 let project = mounted(&producer, "p1", &["**"])
5157 .with_stats(registry.tier(stats::Tier::default()).session("p1"))
5158 .with_hidden(true);
5159
5160 let mut announced = project.announced();
5161 announced.assert_next_active(".svc/foo");
5162 project.request_broadcast(".svc/foo").await.expect("resolves");
5163
5164 let mut report = stats::Report::default();
5165 registry.report(&mut report);
5166 let paths: Vec<_> = report
5167 .traffic
5168 .iter()
5169 .map(|entry| entry.path.as_str().to_string())
5170 .collect();
5171 assert_eq!(paths, ["p1/.svc/foo"]);
5172 }
5173
5174 #[tokio::test]
5176 async fn hidden_route_announced_later_stays_hidden() {
5177 let producer = origin(1).produce();
5178 let consumer = producer.consume();
5179 let mut announced = consumer.announced();
5180 let mut opted = consumer.clone().with_hidden(true).announced();
5181
5182 let hidden = producer.announce(".stats/node", Route::default()).unwrap();
5183 announced.assert_next_wait();
5184 opted.assert_next_active(".stats/node");
5185
5186 drop(hidden);
5187 announced.assert_next_wait();
5188 opted.assert_next_ended(".stats/node");
5189 }
5190
5191 #[tokio::test]
5192 async fn broadcast_announces_its_own_path() {
5193 let producer = origin(1).produce();
5194 let consumer = producer.consume();
5195 let mut announced = consumer.announced();
5196 let mut peer = consumer.clone().excluding(Hop::UNKNOWN).announced();
5197
5198 let broadcast = producer.create_broadcast("room/alice").unwrap();
5200 announced.assert_next_wait();
5201 peer.assert_next_wait();
5202
5203 broadcast.announce(Route::default().with_cost(3)).unwrap();
5204 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(3));
5205 assert_eq!(peer.assert_next_active("room/alice").cost, Cost::new(3));
5206
5207 broadcast.announce(Route::default().with_cost(1)).unwrap();
5209 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(1));
5210 assert_eq!(peer.assert_next_active("room/alice").cost, Cost::new(1));
5211
5212 broadcast.unannounce();
5214 announced.assert_next_ended("room/alice");
5215 peer.assert_next_ended("room/alice");
5216 broadcast.unannounce();
5217 announced.assert_next_wait();
5218 let err = consumer.request_broadcast("room/alice").await.err().unwrap();
5219 assert!(matches!(err, Error::Unroutable));
5220
5221 broadcast.announce(Route::default()).unwrap();
5223 announced.assert_next_active("room/alice");
5224 peer.assert_next_active("room/alice");
5225 broadcast.close();
5226 announced.assert_next_ended("room/alice");
5227 peer.assert_next_ended("room/alice");
5228 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
5229 announced.assert_next_wait();
5230 }
5231
5232 #[tokio::test]
5233 async fn broadcast_announcement_retracts_with_the_last_producer() {
5234 let producer = origin(1).produce();
5235 let consumer = producer.consume();
5236 let mut announced = consumer.announced();
5237
5238 let broadcast = producer.create_broadcast("room/alice").unwrap();
5239 let clone = broadcast.clone();
5240 broadcast.announce(Route::default()).unwrap();
5241 announced.assert_next_active("room/alice");
5242
5243 drop(broadcast);
5245 announced.assert_next_wait();
5246 drop(clone);
5247 announced.assert_next_ended("room/alice");
5248 }
5249
5250 #[tokio::test]
5251 async fn publish_creates_and_announces_together() {
5252 let producer = origin(1).produce();
5253 let mut announced = producer.consume().announced();
5254 let _broadcast = producer.publish("room/alice", Route::default()).unwrap();
5255 announced.assert_next_active("room/alice");
5256 }
5257
5258 #[tokio::test]
5259 async fn standalone_broadcast_cannot_announce() {
5260 let broadcast = broadcast::Info::new().produce();
5261 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
5262 broadcast.unannounce();
5264 }
5265
5266 #[tokio::test]
5267 async fn announce_replays_to_late_cursor() {
5268 let producer = origin(1).produce();
5269 let _a = producer.announce("room/alice", Route::default()).unwrap();
5270 let _b = producer.announce("room/bob", Route::default()).unwrap();
5271
5272 let mut announced = producer.consume().announced();
5273 announced.assert_next_active("room/alice");
5275 announced.assert_next_active("room/bob");
5276 announced.assert_next_wait();
5277 }
5278
5279 #[tokio::test]
5280 async fn announce_keeps_its_prefix_under_a_producer_scope() {
5281 let producer = origin(1).produce();
5282 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
5283
5284 let _a = scoped.announce("", Route::default()).unwrap();
5286 let mut announced = producer.consume().announced();
5287 announced.assert_next_active("");
5288
5289 assert!(matches!(
5291 scoped.announce("other", Route::default()),
5292 Err(Error::Unauthorized)
5293 ));
5294 }
5295
5296 #[tokio::test]
5297 async fn cursor_keeps_an_overlapping_prefix_above_its_scope() {
5298 let producer = origin(1).produce();
5299 let _a = producer.announce("", Route::default()).unwrap();
5300
5301 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
5302 let mut announced = consumer.announced();
5303 announced.assert_next_active("");
5304 }
5305
5306 #[tokio::test]
5307 async fn cursor_root_strips_prefix() {
5308 let producer = origin(1).produce();
5309 let _a = producer.announce("room/alice", Route::default()).unwrap();
5310
5311 let consumer = producer
5312 .consume()
5313 .scope("room", &Patterns::from(Pattern::all()))
5314 .unwrap();
5315 let mut announced = consumer.announced();
5316 announced.assert_next_active("alice");
5317 }
5318
5319 #[tokio::test]
5320 async fn best_route_wins_and_fails_over() {
5321 let producer = origin(1).produce();
5322 let mut announced = producer.consume().announced();
5323
5324 let expensive = producer
5325 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5326 .unwrap();
5327 let route = announced.assert_next_active("room");
5328 assert_eq!(route.cost, Cost::new(5));
5329
5330 let cheap = producer
5332 .announce("room", Route::default().with_hops(hops(&[20])).with_cost(1))
5333 .unwrap();
5334 let route = announced.assert_next_active("room");
5335 assert_eq!(route.cost, Cost::new(1));
5336
5337 drop(cheap);
5339 let route = announced.assert_next_active("room");
5340 assert_eq!(route.cost, Cost::new(5));
5341
5342 drop(expensive);
5344 announced.assert_next_ended("room");
5345 }
5346
5347 #[tokio::test]
5348 async fn identical_reannounce_is_invisible() {
5349 let producer = origin(1).produce();
5350 let mut announced = producer.consume().announced();
5351
5352 let old = producer
5353 .announce("room", Route::default().with_hops(hops(&[10])))
5354 .unwrap();
5355 let first = announced.assert_next_active("room");
5356 assert_eq!(first.hops.as_slice(), hops(&[10]).as_slice());
5357
5358 let _new = producer
5362 .announce("room", Route::default().with_hops(hops(&[10])))
5363 .unwrap();
5364 announced.assert_next_wait();
5365
5366 drop(old);
5368 announced.assert_next_wait();
5369 }
5370
5371 #[tokio::test]
5372 async fn exclude_hides_routes_through_the_peer() {
5373 let producer = origin(1).produce();
5374 let _a = producer
5375 .announce("room", Route::default().with_hops(hops(&[7])))
5376 .unwrap();
5377
5378 let mut hidden = producer.consume().excluding(origin(7)).announced();
5379 hidden.assert_next_wait();
5380
5381 let mut visible = producer.consume().excluding(origin(8)).announced();
5382 visible.assert_next_active("room");
5383 }
5384
5385 #[tokio::test]
5386 async fn exclude_matches_via_when_the_chain_is_anonymous() {
5387 let producer = origin(1).produce();
5388 let assigned = origin(777);
5389 let _echoed = producer
5390 .announce("echoed", Route::default().with_hops(hops(&[0])).with_via(assigned))
5391 .unwrap();
5392 let _local = producer
5393 .announce("local", Route::default().with_hops(hops(&[10])))
5394 .unwrap();
5395
5396 let mut hidden = producer.consume().excluding(assigned).announced();
5397 hidden.assert_next_active("local");
5398 hidden.assert_next_wait();
5399 }
5400
5401 #[tokio::test]
5402 async fn anonymous_route_loses_to_identified_at_any_cost() {
5403 let producer = origin(1).produce();
5404 let mut announced = producer.consume().announced();
5405
5406 let _anonymous = producer
5407 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(1))
5408 .unwrap();
5409 let route = announced.assert_next_active("room");
5410 assert!(route.is_anonymous());
5411 assert_eq!(route.cost, Cost::new(1));
5412
5413 let _identified = producer
5414 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5415 .unwrap();
5416 let route = announced.assert_next_active("room");
5417 assert!(!route.is_anonymous());
5418 assert_eq!(route.cost, Cost::new(5));
5419 }
5420
5421 #[tokio::test]
5424 async fn a_stamped_route_ranks_below_an_identified_one() {
5425 let producer = origin(1).produce();
5426 let mut announced = producer.consume().announced();
5427
5428 let mut stamped = Hops::new();
5429 stamped.stamp(origin(5)).unwrap();
5430 assert_eq!(stamped.as_slice(), &[origin(5), Hop::UNKNOWN]);
5431
5432 let _legacy = producer
5433 .announce("room", Route::default().with_hops(stamped).with_cost(0))
5434 .unwrap();
5435 let route = announced.assert_next_active("room");
5436 assert!(route.is_anonymous());
5437
5438 let _identified = producer
5439 .announce("room", Route::default().with_hops(hops(&[10, 11])).with_cost(5))
5440 .unwrap();
5441 let route = announced.assert_next_active("room");
5442 assert_eq!(route.hops.as_slice(), hops(&[10, 11]).as_slice());
5443 }
5444
5445 #[test]
5446 fn stamping_keeps_the_leading_zero_and_names_nothing_else() {
5447 let mut chain = hops(&[0, 7]);
5448 chain.stamp(origin(5)).unwrap();
5449 assert_eq!(chain.as_slice(), hops(&[5, 0, 7]).as_slice());
5450
5451 let mut named = hops(&[7, 0]);
5453 named.stamp(origin(5)).unwrap();
5454 assert_eq!(named.as_slice(), hops(&[7, 0]).as_slice());
5455
5456 let mut full = Hops::try_from(vec![Hop::UNKNOWN; MAX_HOPS]).unwrap();
5458 assert_eq!(full.stamp(origin(5)), Err(InvalidHop::TooMany));
5459 }
5460
5461 #[tokio::test]
5462 async fn anonymous_routes_order_by_cost() {
5463 let producer = origin(1).produce();
5464 let mut announced = producer.consume().announced();
5465
5466 let expensive = producer
5467 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(5))
5468 .unwrap();
5469 let route = announced.assert_next_active("room");
5470 assert_eq!(route.cost, Cost::new(5));
5471
5472 let _cheap = producer
5473 .announce("room", Route::default().with_hops(hops(&[0, 7])).with_cost(1))
5474 .unwrap();
5475 let route = announced.assert_next_active("room");
5476 assert!(route.is_anonymous());
5477 assert_eq!(route.cost, Cost::new(1));
5478
5479 drop(expensive);
5480 announced.assert_next_wait();
5481 }
5482
5483 #[tokio::test]
5484 async fn anonymous_chain_from_identified_peer_still_ranks_last() {
5485 let producer = origin(1).produce();
5486 let mut announced = producer.consume().announced();
5487
5488 let _anonymous = producer
5489 .announce(
5490 "room",
5491 Route::default()
5492 .with_hops(hops(&[0, 7]))
5493 .with_cost(1)
5494 .with_via(origin(7)),
5495 )
5496 .unwrap();
5497 announced.assert_next_active("room");
5498
5499 let _identified = producer
5500 .announce("room", Route::default().with_hops(hops(&[10, 20])).with_cost(5))
5501 .unwrap();
5502 let route = announced.assert_next_active("room");
5503 assert!(!route.is_anonymous());
5504 assert_eq!(route.cost, Cost::new(5));
5505 }
5506
5507 #[tokio::test]
5508 async fn request_prefers_identified_over_cheaper_anonymous() {
5509 let producer = origin(1).produce();
5510 let consumer = producer.consume();
5511
5512 let anonymous = producer
5513 .dynamic("room", Route::default().with_hops(hops(&[0])).with_cost(1))
5514 .unwrap();
5515 let identified = producer
5516 .dynamic("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5517 .unwrap();
5518
5519 let _pending = consumer.request_broadcast("room/alice");
5520 let request = queued(&identified).await;
5521 assert_eq!(request.path().as_str(), "room/alice");
5522 assert!(
5523 anonymous.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
5524 "the cheaper anonymous route must not serve"
5525 );
5526 }
5527
5528 #[tokio::test]
5529 async fn update_reprices_in_place() {
5530 let producer = origin(1).produce();
5531 let mut announced = producer.consume().announced();
5532
5533 let announcement = producer.announce("room", Route::default()).unwrap();
5534 announced.assert_next_active("room");
5535
5536 announcement.update(Route::default().with_cost(9)).unwrap();
5537 let route = announced.assert_next_active("room");
5538 assert_eq!(route.cost, Cost::new(9));
5539 }
5540
5541 #[tokio::test]
5542 async fn retract_after_undelivered_reprice_still_delivered() {
5543 let producer = origin(1).produce();
5544 let mut announced = producer.consume().announced();
5545
5546 let announcement = producer.announce("room", Route::default()).unwrap();
5547 announced.assert_next_active("room");
5548
5549 announcement.update(Route::default().with_cost(9)).unwrap();
5553 drop(announcement);
5554 announced.assert_next_ended("room");
5555 announced.assert_next_wait();
5556 }
5557
5558 #[tokio::test]
5559 async fn scoped_cursor_advertises_most_specific_covering_route() {
5560 let producer = origin(1).produce();
5561 let _broad = producer.announce("room", Route::default().with_cost(1)).unwrap();
5564 let _narrow = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
5565
5566 let consumer = producer
5567 .consume()
5568 .scope("room/alice", &Patterns::from(Pattern::all()))
5569 .unwrap();
5570 let mut announced = consumer.announced();
5571 let route = announced.assert_next_active("");
5572 assert_eq!(route.cost, Cost::new(9));
5573 announced.assert_next_wait();
5574 }
5575
5576 #[tokio::test]
5577 async fn capture_change_retracts_before_reannouncing_a_presented_prefix() {
5578 let producer = origin(1).produce();
5579 let _broad = producer.announce("room", Route::default()).unwrap();
5580 let exact = producer.announce("room/alice", Route::default()).unwrap();
5581 let consumer = producer
5582 .consume()
5583 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
5584 .unwrap()
5585 .scope("room/alice", &Patterns::from(Pattern::all()))
5586 .unwrap();
5587 let mut announced = consumer.announced();
5588
5589 let first = announced.next().now_or_never().expect("next").expect("announce");
5590 assert_eq!(first.prefix.as_str(), "");
5591 assert_eq!(first.kind, AnnounceKind::Announced);
5592 assert_eq!(first.captures, Some(Vec::new()));
5593
5594 drop(exact);
5595 let retracted = announced.next().now_or_never().expect("next").expect("retract");
5596 assert_eq!(retracted.prefix.as_str(), "");
5597 assert_eq!(retracted.kind, AnnounceKind::Retracted);
5598 assert_eq!(retracted.captures, Some(Vec::new()));
5599 let replacement = announced.next().now_or_never().expect("next").expect("announce");
5600 assert_eq!(replacement.prefix.as_str(), "");
5601 assert_eq!(replacement.kind, AnnounceKind::Announced);
5602 assert_eq!(replacement.captures, None);
5603 }
5604
5605 #[tokio::test]
5606 async fn routed_broadcast_resolves_once_announced() {
5607 let producer = origin(1).produce();
5608 let consumer = producer.consume();
5609
5610 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
5612 assert!((&mut resolving).now_or_never().is_none());
5613
5614 let broadcast = producer.create_broadcast("room/alice").unwrap();
5616 for _ in 0..20 {
5617 tokio::task::yield_now().await;
5618 }
5619 assert!((&mut resolving).now_or_never().is_none());
5620
5621 broadcast.announce(Route::default()).unwrap();
5622 let resolved = resolving.await.expect("resolves once announced");
5623 assert_eq!(resolved.info().path.as_str(), "room/alice");
5624 drop(broadcast);
5625 }
5626
5627 #[tokio::test]
5630 async fn cheaper_remote_route_beats_a_local_broadcast() {
5631 let producer = origin(1).produce();
5632 let consumer = producer.consume();
5633 let mut announced = consumer.announced();
5634
5635 let _local = producer.publish("room/alice", Route::default().with_cost(5)).unwrap();
5636 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(5));
5637
5638 let server = producer
5639 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(1))
5640 .unwrap();
5641 let route = announced.assert_next_active("room/alice");
5642 assert_eq!(route.cost, Cost::new(1));
5643 assert_eq!(route.hops, hops(&[10]));
5644
5645 let pending = consumer.request_broadcast("room/alice");
5647 let request = queued(&server).await;
5648 let upstream = broadcast::Info::new().produce();
5649 request.accept(&upstream);
5650 pending.await.expect("resolves through the cheaper route");
5651 }
5652
5653 #[tokio::test]
5657 async fn cheaper_route_after_a_front_wins_new_requests() {
5658 let producer = origin(1).produce();
5659 let consumer = producer.consume();
5660
5661 let _local = producer.publish("room/alice", Route::default().with_cost(5)).unwrap();
5662 let first = consumer
5663 .request_broadcast("room/alice")
5664 .await
5665 .expect("resolves locally");
5666
5667 let server = producer
5668 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(1))
5669 .unwrap();
5670
5671 let pending = consumer.request_broadcast("room/alice");
5672 let request = queued(&server).await;
5673 let upstream = broadcast::Info::new().produce();
5674 request.accept(&upstream);
5675 let second = pending.await.expect("resolves through the cheaper route");
5676
5677 assert!(!first.is_closed(), "the old front must keep serving its readers");
5678 assert!(!first.is_clone(&second), "the newcomer must not join the old front");
5679 }
5680
5681 #[tokio::test]
5685 async fn local_broadcast_wins_a_tie_with_a_hopless_route() {
5686 let producer = origin(1).produce();
5687 let consumer = producer.consume();
5688
5689 let _local = producer.publish("room/alice", Route::default()).unwrap();
5690 let server = producer.dynamic("room/alice", Route::default()).unwrap();
5691
5692 let resolved = tokio::time::timeout(Duration::from_secs(1), consumer.request_broadcast("room/alice"))
5694 .await
5695 .expect("the newer hopless route won the tie")
5696 .expect("resolves");
5697 assert_eq!(resolved.info().path.as_str(), "room/alice");
5698 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
5699 }
5700
5701 #[tokio::test]
5704 async fn announce_stats_follow_the_advertisement() {
5705 let registry = stats::Registry::new(stats::Config::new());
5706 let producer = origin(1)
5707 .produce()
5708 .with_stats(registry.tier(stats::Tier::default()).session("root"));
5709 let announces = || {
5710 registry
5711 .snapshot()
5712 .traffic()
5713 .into_iter()
5714 .find(|(_, role, _)| *role == stats::Role::Subscriber)
5715 .map(|(_, _, traffic)| (traffic.announces_started, traffic.announces_ended))
5716 .unwrap_or_default()
5717 };
5718
5719 let broadcast = producer.create_broadcast("room/alice").unwrap();
5720 assert_eq!(announces(), (0, 0), "a hidden broadcast is not announced");
5721 broadcast.announce(Route::default()).unwrap();
5722 broadcast
5723 .announce(Route {
5724 cost: Cost::new(3),
5725 ..Route::default()
5726 })
5727 .unwrap();
5728 assert_eq!(announces(), (1, 0), "a re-price is not another announce");
5729 broadcast.unannounce();
5730 assert_eq!(announces(), (1, 1));
5731 broadcast.announce(Route::default()).unwrap();
5732 drop(broadcast);
5733 assert_eq!(announces(), (2, 2));
5734 }
5735
5736 #[tokio::test]
5738 async fn local_broadcast_wins_a_cost_tie() {
5739 let producer = origin(1).produce();
5740 let consumer = producer.consume();
5741 let mut announced = consumer.announced();
5742
5743 let server = producer
5744 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(2))
5745 .unwrap();
5746 announced.assert_next_active("room/alice");
5747 let _local = producer.publish("room/alice", Route::default().with_cost(2)).unwrap();
5748 assert!(announced.assert_next_active("room/alice").hops.is_empty());
5749
5750 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5751 assert_eq!(resolved.info().path.as_str(), "room/alice");
5752 for _ in 0..20 {
5753 tokio::task::yield_now().await;
5754 }
5755 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
5756 }
5757
5758 #[tokio::test]
5761 async fn unannounce_ends_the_front() {
5762 let producer = origin(1).produce();
5763 let consumer = producer.consume();
5764
5765 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5766 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5767
5768 broadcast.unannounce();
5769 let err = consumer.request_broadcast("room/alice").await.err().unwrap();
5770 assert!(matches!(err, Error::Unroutable), "joined a retracted front: {err}");
5771 settle(|| resolved.is_closed()).await;
5772 assert!(!broadcast.consume().is_closed(), "the broadcast itself lives on");
5773
5774 broadcast.announce(Route::default()).unwrap();
5776 let again = consumer.request_broadcast("room/alice").await.expect("resolves again");
5777 assert!(!again.is_clone(&resolved));
5778 }
5779
5780 #[tokio::test]
5784 async fn unannounce_keeps_a_track_awaiting_its_info() {
5785 let producer = origin(1).produce();
5786 let consumer = producer.consume();
5787
5788 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5789 let mut dynamic = broadcast.dynamic();
5790 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5791 let track = resolved.track("video").unwrap();
5792 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5793 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5794 .await
5795 .expect("the front asked the source")
5796 .expect("request");
5797
5798 broadcast.unannounce();
5799 settle(|| resolved.is_closed()).await;
5800
5801 let source = request.accept(None);
5802 let mut group = source.append_group().unwrap();
5803 group.write_frame(crate::Timestamp::ZERO, b"late".as_ref()).unwrap();
5804 group.finish().unwrap();
5805 source.finish().unwrap();
5806
5807 let mut subscription = subscribing.await.unwrap().expect("subscribe survives the retraction");
5808 let mut group = subscription.recv_group().await.unwrap().expect("the source's group");
5809 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"late");
5810 assert!(matches!(subscription.recv_group().await, Ok(None)), "ends cleanly");
5811 }
5812
5813 async fn served_front() -> (Dynamic, broadcast::Producer, broadcast::Dynamic, broadcast::Consumer) {
5816 let producer = origin(1).produce();
5817 let consumer = producer.consume();
5818 let server = producer
5819 .dynamic("room/alice", Route::default().with_hops(hops(&[10])))
5820 .unwrap();
5821 let pending = consumer.request_broadcast("room/alice");
5822 let upstream = broadcast::Info::new().produce();
5823 let dynamic = upstream.dynamic();
5824 queued(&server).await.accept(&upstream);
5825 let resolved = pending.await.expect("resolves");
5826 (server, upstream, dynamic, resolved)
5827 }
5828
5829 #[tokio::test]
5834 async fn returning_reader_skips_a_warm_cache_the_copy_resolved_past() {
5835 let ms = |v: u64| crate::Timestamp::from_millis(v).unwrap();
5836 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5837 let budget = track::Subscription::default().with_max_age(Duration::from_millis(100));
5838
5839 let track = resolved.track("audio").unwrap();
5840 let b = budget.clone();
5841 let subscribing = tokio::spawn(async move { track.subscribe(b).await });
5842 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5843 .await
5844 .expect("the front asked the source")
5845 .expect("request");
5846 let source = request.resolving_start().accept(None);
5847 for seq in 0..4u64 {
5848 let mut group = source.create_group(seq.into()).unwrap();
5849 group.write_frame(ms(seq * 20), b"old".as_ref()).unwrap();
5850 group.finish().unwrap();
5851 }
5852 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5853 subscription.recv_group().await.unwrap().expect("the live group");
5854 drop(subscription);
5855 tokio::time::timeout(Duration::from_secs(1), source.unused())
5856 .await
5857 .expect("parked")
5858 .expect("source open");
5859 drop(source);
5860
5861 let track = resolved.track("audio").unwrap();
5862 let subscribing = tokio::spawn(async move { track.subscribe(budget).await });
5863 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5864 .await
5865 .expect("the front asked the source again")
5866 .expect("request");
5867 let mut source = request.resolving_start().accept(None);
5868 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
5869
5870 assert!(
5872 tokio::time::timeout(Duration::from_millis(50), subscription.recv_group())
5873 .await
5874 .is_err(),
5875 "the warm cache was served before the copy resolved its start"
5876 );
5877
5878 source.start_at(20).unwrap();
5880 let mut group = source.create_group(20u64.into()).unwrap();
5881 group.write_frame(ms(2000), b"new".as_ref()).unwrap();
5882 group.finish().unwrap();
5883 let group = subscription.recv_group().await.unwrap().expect("the live group");
5884 assert_eq!(group.sequence, 20, "a stale warm group was served");
5885 }
5886
5887 #[tokio::test]
5892 async fn returning_reader_replays_a_current_warm_cache() {
5893 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5894
5895 let track = resolved.track("catalog").unwrap();
5896 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5897 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5898 .await
5899 .expect("the front asked the source")
5900 .expect("request");
5901 let source = request.resolving_start().accept(None);
5902 let mut group = source.create_group(0u64.into()).unwrap();
5903 group.write_frame(crate::Timestamp::ZERO, b"snapshot".as_ref()).unwrap();
5904 group.finish().unwrap();
5905 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5906 subscription.recv_group().await.unwrap().expect("the catalog");
5907 drop(subscription);
5908 tokio::time::timeout(Duration::from_secs(1), source.unused())
5909 .await
5910 .expect("parked")
5911 .expect("source open");
5912 drop(source);
5913
5914 let track = resolved.track("catalog").unwrap();
5915 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5916 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5917 .await
5918 .expect("the front asked the source again")
5919 .expect("request");
5920 let mut source = request.resolving_start().accept(None);
5921 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
5922
5923 let reading = tokio::spawn(async move {
5926 let mut group = subscription.recv_group().await.unwrap().expect("the catalog");
5927 assert_eq!(group.sequence, 0);
5928 group.read_frame().await.unwrap().expect("the snapshot").payload
5929 });
5930 tokio::task::yield_now().await;
5931 assert_eq!(
5932 source.subscription().and_then(|sub| sub.start),
5933 Some(track::Position { group: 0, frame: 1 }),
5934 "the re-splice asked past the cached catalog"
5935 );
5936 source.start_at(0).unwrap();
5937 let mut tail = source.create_group(0u64.into()).unwrap();
5938 tail.start_at(1).unwrap();
5939 tail.finish().unwrap();
5940 let payload = tokio::time::timeout(Duration::from_secs(1), reading)
5941 .await
5942 .expect("the returning reader never got the catalog")
5943 .unwrap();
5944 assert_eq!(&payload[..], b"snapshot");
5945 }
5946
5947 #[tokio::test(start_paused = true)]
5951 async fn finished_track_is_forgotten_after_the_linger() {
5952 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5953
5954 let track = resolved.track("catalog").unwrap();
5955 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5956 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5957 .await
5958 .expect("the front asked the source")
5959 .expect("request");
5960 let source = request.resolving_start().accept(None);
5961 let mut group = source.create_group(0u64.into()).unwrap();
5962 group.write_frame(crate::Timestamp::ZERO, b"snapshot".as_ref()).unwrap();
5963 group.finish().unwrap();
5964 source.finish().unwrap();
5965 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5966 assert_eq!(
5967 next_group(&mut subscription)
5968 .await
5969 .unwrap()
5970 .expect("the catalog")
5971 .sequence,
5972 0
5973 );
5974 assert!(next_group(&mut subscription).await.unwrap().is_none());
5975 drop(subscription);
5976 drop(source);
5977
5978 let mut subscription = resolved.track("catalog").unwrap().subscribe(None).await.unwrap();
5980 assert_eq!(
5981 next_group(&mut subscription)
5982 .await
5983 .unwrap()
5984 .expect("the catalog")
5985 .sequence,
5986 0
5987 );
5988 assert!(next_group(&mut subscription).await.unwrap().is_none());
5989 drop(subscription);
5990 assert!(
5991 tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5992 .await
5993 .is_err(),
5994 "a finished track within the linger asked the source again"
5995 );
5996
5997 tokio::time::sleep(TRACK_IDLE_LINGER).await;
5999
6000 let track = resolved.track("catalog").unwrap();
6001 let _subscribing = tokio::spawn(async move { track.subscribe(None).await });
6002 tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6003 .await
6004 .expect("the finished track outlived the linger")
6005 .expect("request");
6006 }
6007
6008 #[tokio::test]
6013 async fn returning_reader_continues_an_open_warm_group() {
6014 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
6015
6016 async fn read(group: &mut group::Consumer) -> Vec<u8> {
6017 let frame = tokio::time::timeout(Duration::from_secs(1), group.read_frame())
6018 .await
6019 .expect("frame")
6020 .unwrap()
6021 .expect("group ended");
6022 frame.payload.to_vec()
6023 }
6024
6025 let mut expect: Vec<&[u8]> = Vec::new();
6026 let mut floor: Option<track::Position> = None;
6027 for (round, payload) in [b"a".as_ref(), b"b", b"c"].into_iter().enumerate() {
6028 let track = resolved.track("log").unwrap();
6029 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6030 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6031 .await
6032 .expect("the front asked the source")
6033 .expect("request");
6034 let mut source = request.resolving_start().accept(None);
6035 let mut subscription = subscribing.await.unwrap().expect("subscribe");
6036
6037 source.start_at(0).unwrap();
6039 let mut group = source.create_group(0u64.into()).unwrap();
6040 if let Some(floor) = floor {
6041 group.start_at(floor.frame).unwrap();
6042 }
6043 group.write_frame(crate::Timestamp::ZERO, payload).unwrap();
6044 expect.push(payload);
6045 source
6046 .insert_datagram(10, crate::Timestamp::ZERO, b"datagram".as_ref())
6047 .unwrap();
6048
6049 let mut reading = tokio::time::timeout(Duration::from_secs(1), subscription.recv_group())
6050 .await
6051 .expect("group 0")
6052 .unwrap()
6053 .expect("track ended");
6054 assert_eq!(reading.sequence, 0);
6055 for frame in &expect {
6056 assert_eq!(read(&mut reading).await, *frame, "round {round}");
6057 }
6058 assert_eq!(
6059 source.subscription().and_then(|sub| sub.start),
6060 floor,
6061 "round {round} asked for the wrong continuation"
6062 );
6063
6064 drop(reading);
6065 drop(subscription);
6066 tokio::time::timeout(Duration::from_secs(1), source.unused())
6067 .await
6068 .expect("parked")
6069 .expect("source open");
6070 drop(group);
6071 drop(source);
6072 floor = Some(track::Position {
6073 group: 0,
6074 frame: expect.len() as u64,
6075 });
6076 }
6077 }
6078
6079 #[tokio::test]
6086 async fn a_takeover_keeps_the_open_group_head_for_later_readers() {
6087 tokio::time::pause();
6088 for dies_first in [true, false] {
6089 takeover_keeps_the_open_group_head(dies_first).await;
6090 }
6091 }
6092
6093 async fn takeover_keeps_the_open_group_head(dies_first: bool) {
6094 let producer = origin(1).produce();
6095 let first_server = producer
6096 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(5))
6097 .unwrap();
6098 let pending = producer.consume().request_broadcast("room/alice");
6099 let upstream = broadcast::Info::new().produce();
6100 let mut dynamic = upstream.dynamic();
6101 queued(&first_server).await.accept(&upstream);
6102 let resolved = pending.await.unwrap();
6103
6104 async fn read(group: &mut group::Consumer) -> Vec<u8> {
6105 let frame = tokio::time::timeout(Duration::from_secs(1), group.read_frame())
6106 .await
6107 .expect("frame")
6108 .unwrap()
6109 .expect("group ended");
6110 frame.payload.to_vec()
6111 }
6112
6113 let track = resolved.track("log").unwrap();
6115 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6116 let source = dynamic.requested_track().await.unwrap().accept(None);
6117 let mut group = source.create_group(0u64.into()).unwrap();
6118 group.write_frame(crate::Timestamp::ZERO, b"a".as_ref()).unwrap();
6119 group.write_frame(crate::Timestamp::ZERO, b"b".as_ref()).unwrap();
6120 let mut subscription = subscribing.await.unwrap().unwrap();
6121 let mut reading = subscription.recv_group().await.unwrap().unwrap();
6122 assert_eq!(read(&mut reading).await, b"a");
6123 assert_eq!(read(&mut reading).await, b"b");
6124
6125 let mut first = Some((group, source, upstream, dynamic));
6128 let mut joining = None;
6129 if dies_first {
6130 drop(first.take());
6131 assert!(futures::FutureExt::now_or_never(subscription.recv_group()).is_none());
6134 let track = resolved.track("log").unwrap();
6135 joining = Some(tokio::spawn(async move {
6136 let mut joined = track.subscribe(None).await.unwrap();
6137 let group = joined.recv_group().await.unwrap().expect("track ended");
6138 (joined, group)
6139 }));
6140 }
6141 let cost = if dies_first { 5 } else { 0 };
6142 let replacement_server = producer
6143 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(cost))
6144 .unwrap();
6145 let replacement = broadcast::Info::new().produce();
6146 let mut replacement_dynamic = replacement.dynamic();
6147 queued(&replacement_server).await.accept(&replacement);
6148 let mut next = tokio::time::timeout(Duration::from_secs(1), replacement_dynamic.requested_track())
6149 .await
6150 .expect("the front asked the replacement")
6151 .unwrap()
6152 .resolving_start()
6153 .accept(None);
6154 next.start_at(0).unwrap();
6155 let mut next_group = next.create_group(0u64.into()).unwrap();
6156 next_group.start_at(2).unwrap();
6157 next_group.write_frame(crate::Timestamp::ZERO, b"c".as_ref()).unwrap();
6158 assert_eq!(read(&mut reading).await, b"c", "dies_first={dies_first}");
6159
6160 if let Some((group, ..)) = &mut first {
6162 let mut independent = group.consume();
6163 group.write_frame(crate::Timestamp::ZERO, b"x".as_ref()).unwrap();
6164 for expect in [b"a", b"b", b"x"] {
6165 assert_eq!(read(&mut independent).await, expect);
6166 }
6167 group.finish().unwrap();
6168 }
6169 drop(first);
6170
6171 let track = resolved.track("log").unwrap();
6172 let mut later = tokio::time::timeout(Duration::from_secs(1), track.subscribe(None))
6173 .await
6174 .expect("subscribe")
6175 .unwrap();
6176 let mut late = tokio::time::timeout(Duration::from_secs(1), later.recv_group())
6177 .await
6178 .unwrap_or_else(|_| panic!("dies_first={dies_first}: the later reader never got the group"))
6179 .unwrap()
6180 .expect("track ended");
6181 assert_eq!(late.sequence, 0);
6182 for expect in [b"a", b"b", b"c"] {
6183 assert_eq!(read(&mut late).await, expect, "dies_first={dies_first}");
6184 }
6185 let mut joined = None;
6186 if let Some(joining) = joining {
6187 let (subscription, mut group) = tokio::time::timeout(Duration::from_secs(1), joining)
6188 .await
6189 .expect("the reader joining mid-outage never got the group")
6190 .unwrap();
6191 assert_eq!(group.sequence, 0);
6192 for expect in [b"a", b"b", b"c"] {
6193 assert_eq!(read(&mut group).await, expect);
6194 }
6195 joined = Some((subscription, group));
6196 }
6197
6198 assert!(
6200 tokio::time::timeout(Duration::from_secs(1), subscription.recv_group())
6201 .await
6202 .is_err(),
6203 "dies_first={dies_first}: the group was handed out twice"
6204 );
6205
6206 next_group.write_frame(crate::Timestamp::ZERO, b"d".as_ref()).unwrap();
6208 assert_eq!(read(&mut reading).await, b"d");
6209 assert_eq!(read(&mut late).await, b"d");
6210 if let Some((_, group)) = &mut joined {
6211 assert_eq!(read(group).await, b"d");
6212 }
6213 next_group.finish().unwrap();
6214 drop((
6215 joined,
6216 reading,
6217 late,
6218 subscription,
6219 later,
6220 next,
6221 replacement,
6222 replacement_server,
6223 first_server,
6224 ));
6225 }
6226
6227 #[tokio::test]
6232 async fn returning_reader_gets_a_live_open_warm_group() {
6233 use futures::FutureExt;
6234
6235 tokio::time::pause();
6236 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
6237
6238 let track = resolved.track("video").unwrap();
6240 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6241 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6242 .await
6243 .expect("the front asked the source")
6244 .expect("request");
6245 let mut source = request.resolving_start().accept(None);
6246 let mut subscription = subscribing.await.unwrap().expect("subscribe");
6247 source.start_at(0).unwrap();
6248 let mut group = source.create_group(0u64.into()).unwrap();
6249 group.write_frame(crate::Timestamp::ZERO, b"a".as_ref()).unwrap();
6250 let mut reading = subscription.recv_group().await.unwrap().expect("group 0");
6251 assert_eq!(reading.read_frame().await.unwrap().unwrap().payload.as_ref(), b"a");
6252 drop(reading);
6253 drop(subscription);
6254 tokio::time::timeout(Duration::from_secs(1), source.unused())
6255 .await
6256 .expect("parked")
6257 .expect("source open");
6258 drop(group);
6259 drop(source);
6260
6261 let track = resolved.track("video").unwrap();
6263 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6264 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6265 .await
6266 .expect("the front asked the source")
6267 .expect("request");
6268 let mut source = request.resolving_start().accept(None);
6269 let mut subscription = subscribing.await.unwrap().expect("subscribe");
6270 source.start_at(0).unwrap();
6271 let mut group = source.create_group(0u64.into()).unwrap();
6272 group.start_at(1).unwrap();
6273 group.write_frame(crate::Timestamp::ZERO, b"b".as_ref()).unwrap();
6274
6275 let mut reading = tokio::time::timeout(Duration::from_secs(1), subscription.recv_group())
6276 .await
6277 .expect("group 0")
6278 .unwrap()
6279 .expect("track ended");
6280 assert_eq!(reading.sequence, 0);
6281 let early = reading.finished().now_or_never();
6282 assert!(
6283 early.is_none(),
6284 "the live group already resolved its end before a frame was read: {early:?}"
6285 );
6286
6287 group.finish().unwrap();
6289 assert_eq!(reading.finished().await.unwrap(), 0);
6290 assert_eq!(reading.read_frame().await.unwrap().unwrap().payload.as_ref(), b"a");
6291 assert_eq!(reading.finished().await.unwrap(), 1);
6292 assert_eq!(reading.read_frame().await.unwrap().unwrap().payload.as_ref(), b"b");
6293 assert_eq!(reading.finished().await.unwrap(), 2);
6294 }
6295
6296 #[tokio::test]
6297 async fn warm_head_survives_another_takeover_before_park() {
6298 tokio::time::pause();
6299 let producer = origin(1).produce();
6300 let _server = producer
6301 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(5))
6302 .unwrap();
6303 let pending = producer.consume().request_broadcast("room/alice");
6304 let upstream = broadcast::Info::new().produce();
6305 let mut dynamic = upstream.dynamic();
6306 queued(&_server).await.accept(&upstream);
6307 let resolved = pending.await.unwrap();
6308
6309 let track = resolved.track("log").unwrap();
6310 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6311 let source = dynamic.requested_track().await.unwrap().accept(None);
6312 let mut group = source.create_group(0u64.into()).unwrap();
6313 group.write_frame(crate::Timestamp::ZERO, b"a".as_ref()).unwrap();
6314 group.write_frame(crate::Timestamp::ZERO, b"b".as_ref()).unwrap();
6315 let mut subscription = subscribing.await.unwrap().unwrap();
6316 let mut reading = subscription.recv_group().await.unwrap().unwrap();
6317 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"a");
6318 drop(reading);
6319 drop(subscription);
6320 source.unused().await.unwrap();
6321 drop(group);
6322 drop(source);
6323
6324 let track = resolved.track("log").unwrap();
6325 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6326 let mut resumed = dynamic.requested_track().await.unwrap().resolving_start().accept(None);
6327 resumed.start_at(0).unwrap();
6328 let mut subscription = subscribing.await.unwrap().unwrap();
6329 let mut reading = subscription.recv_group().await.unwrap().unwrap();
6330 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"a");
6331
6332 let replacement_server = producer
6333 .dynamic("room/alice", Route::default().with_hops(hops(&[10])))
6334 .unwrap();
6335 let replacement = broadcast::Info::new().produce();
6336 let mut replacement_dynamic = replacement.dynamic();
6337 queued(&replacement_server).await.accept(&replacement);
6338 let mut source = replacement_dynamic
6339 .requested_track()
6340 .await
6341 .unwrap()
6342 .resolving_start()
6343 .accept(None);
6344 source.start_at(0).unwrap();
6345 let mut group = source.create_group(0u64.into()).unwrap();
6346 group.start_at(2).unwrap();
6347 group.write_frame(crate::Timestamp::ZERO, b"c".as_ref()).unwrap();
6348 let mut continuation = reading.clone();
6349 continuation.start_at(2);
6350 let frame = tokio::time::timeout(Duration::from_secs(1), continuation.read_frame())
6351 .await
6352 .expect("replacement frame")
6353 .unwrap()
6354 .unwrap();
6355 assert_eq!(&frame.payload[..], b"c");
6356 drop(continuation);
6357 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"b");
6358 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"c");
6359 drop(reading);
6360 drop(subscription);
6361 source.unused().await.unwrap();
6362 drop(group);
6363 drop(source);
6364
6365 let track = resolved.track("log").unwrap();
6366 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6367 let mut next = replacement_dynamic
6368 .requested_track()
6369 .await
6370 .unwrap()
6371 .resolving_start()
6372 .accept(None);
6373 next.start_at(0).unwrap();
6374 let mut subscription = subscribing.await.unwrap().unwrap();
6375 let mut reading = subscription.recv_group().await.unwrap().unwrap();
6376 for expected in [b"a", b"b", b"c"] {
6377 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], expected);
6378 }
6379 }
6380
6381 #[tokio::test]
6384 async fn unannounce_keeps_a_returning_reader_awaiting_its_info() {
6385 let producer = origin(1).produce();
6386 let consumer = producer.consume();
6387
6388 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
6389 let mut dynamic = broadcast.dynamic();
6390 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
6391 let track = resolved.track("video").unwrap();
6392 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6393 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6394 .await
6395 .expect("the front asked the source")
6396 .expect("request");
6397 let source = request.accept(None);
6398 let mut group = source.append_group().unwrap();
6399 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
6400 group.finish().unwrap();
6401 let mut subscription = subscribing.await.unwrap().expect("subscribe");
6402 subscription.recv_group().await.unwrap().expect("the cached group");
6403 drop(subscription);
6404
6405 tokio::time::timeout(Duration::from_secs(1), source.unused())
6408 .await
6409 .expect("parked")
6410 .expect("source open");
6411 drop(source);
6412 let track = resolved.track("video").unwrap();
6413 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
6414 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
6415 .await
6416 .expect("the front asked the source again")
6417 .expect("request");
6418
6419 broadcast.unannounce();
6420 settle(|| resolved.is_closed()).await;
6421
6422 let source = request.accept(None);
6424 let mut group = source.create_group(1u64.into()).unwrap();
6425 group.write_frame(crate::Timestamp::ZERO, b"late".as_ref()).unwrap();
6426 group.finish().unwrap();
6427 source.finish().unwrap();
6428
6429 let mut subscription = subscribing.await.unwrap().expect("subscribe survives the retraction");
6430 let mut payloads = Vec::new();
6431 while let Some(mut group) = subscription.recv_group().await.expect("ends cleanly") {
6432 payloads.push(group.read_frame().await.unwrap().unwrap().payload);
6433 }
6434 assert_eq!(payloads.last().map(|p| &p[..]), Some(&b"late"[..]));
6435 }
6436
6437 #[tokio::test]
6440 async fn reannounce_before_the_front_acts_keeps_it() {
6441 let producer = origin(1).produce();
6442 let consumer = producer.consume();
6443
6444 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
6445 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
6446
6447 broadcast.unannounce();
6448 broadcast.announce(Route::default()).unwrap();
6449 for _ in 0..20 {
6450 tokio::task::yield_now().await;
6451 }
6452 assert!(!resolved.is_closed(), "the front ended across a reannouncement");
6453 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
6454 assert!(again.is_clone(&resolved));
6455 }
6456
6457 #[tokio::test]
6458 async fn local_broadcast_resolves_once_announced() {
6459 let producer = origin(1).produce();
6460 let consumer = producer.consume();
6461
6462 let broadcast = producer.create_broadcast("room/alice").unwrap();
6464 let err = consumer
6465 .request_broadcast("room/alice")
6466 .now_or_never()
6467 .expect("unroutable is synchronous")
6468 .err()
6469 .unwrap();
6470 assert!(matches!(err, Error::Unroutable));
6471
6472 broadcast.announce(Route::default()).unwrap();
6473 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
6474 assert_eq!(resolved.info().path.as_str(), "room/alice");
6475 drop(broadcast);
6476
6477 let err = consumer
6479 .request_broadcast("room/bob")
6480 .now_or_never()
6481 .expect("unroutable is synchronous")
6482 .err()
6483 .unwrap();
6484 assert!(matches!(err, Error::Unroutable));
6485 }
6486
6487 #[test]
6488 fn create_broadcast_accepts_a_max_depth_path() {
6489 let producer = origin(1).produce();
6490 let path = vec!["a"; Path::MAX_PARTS].join("/");
6491 let _broadcast = producer.create_broadcast(path.as_str()).expect("max depth is allowed");
6492 let deeper = vec!["a"; Path::MAX_PARTS + 1].join("/");
6493 assert!(matches!(
6494 producer.create_broadcast(deeper.as_str()),
6495 Err(Error::BoundsExceeded(_))
6496 ));
6497 }
6498
6499 #[tokio::test]
6500 async fn duplicate_routes_aggregate_until_the_last_leaves() {
6501 let producer = origin(1).produce();
6502 let first = producer.dynamic("live", Route::default().with_cost(3)).unwrap();
6503 let second = producer.dynamic("live", Route::default().with_cost(1)).unwrap();
6504
6505 let mut announced = producer.consume().announced();
6506 let update = announced.next().now_or_never().expect("next").expect("no next");
6507 assert_eq!(update.prefix.as_str(), "live");
6508 assert_eq!(update.kind, AnnounceKind::Announced);
6509 assert_eq!(update.route.cost, Cost::new(1));
6510 announced.assert_next_wait();
6511
6512 drop(second);
6513 let update = announced.next().now_or_never().expect("next").expect("no next");
6514 assert_eq!(update.prefix.as_str(), "live");
6515 assert_eq!(update.kind, AnnounceKind::Updated);
6516 assert_eq!(update.route.cost, Cost::new(3));
6517
6518 drop(first);
6519 announced.assert_next_ended("live");
6520 announced.assert_next_wait();
6521 }
6522
6523 #[test]
6524 fn dynamic_may_cover_a_scope_but_disjoint_prefixes_are_refused() {
6525 let producer = origin(1).produce();
6526 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
6527 let _broad = scoped
6528 .dynamic("", Route::default())
6529 .expect("an overlapping prefix is accepted");
6530
6531 let _ok = scoped
6532 .dynamic("room/alice", Route::default())
6533 .expect("a contained prefix is accepted");
6534 assert!(matches!(
6535 scoped.dynamic("other", Route::default()),
6536 Err(Error::Unauthorized)
6537 ));
6538 }
6539
6540 #[tokio::test]
6541 async fn dynamic_route_keeps_its_producer_scope() {
6542 let producer = origin(1).produce();
6543 let scope = Patterns::from("*/chat".parse::<Pattern>().unwrap());
6544 let scoped = producer.scope("", &scope).unwrap();
6545 let dynamic = scoped.dynamic("", Route::default()).unwrap();
6546
6547 let mut matching = producer
6548 .consume()
6549 .scope("", &scopes(&["room/chat"]))
6550 .unwrap()
6551 .announced();
6552 matching.assert_next_active("");
6553 let mut outside = producer
6554 .consume()
6555 .scope("", &scopes(&["room/video"]))
6556 .unwrap()
6557 .announced();
6558 outside.assert_next_wait();
6559
6560 let refused = producer
6561 .consume()
6562 .request_broadcast("room/video")
6563 .now_or_never()
6564 .expect("an out-of-scope request must be refused synchronously");
6565 assert!(matches!(refused, Err(Error::Unroutable)));
6566 assert!(dynamic.requested_broadcast().now_or_never().is_none());
6567
6568 let _pending = producer.consume().request_broadcast("room/chat");
6569 let request = queued(&dynamic).await;
6570 assert_eq!(request.path().as_str(), "room/chat");
6571 }
6572
6573 #[tokio::test]
6574 async fn scoped_cursor_selects_among_the_routes_it_can_see() {
6575 let producer = origin(1).produce();
6576 let scoped = |pattern: &str| {
6577 producer
6578 .scope("", &Patterns::from(pattern.parse::<Pattern>().unwrap()))
6579 .unwrap()
6580 };
6581 let _chat = scoped("*/chat").dynamic("", Route::default().with_cost(1)).unwrap();
6582 let _video = scoped("*/video").dynamic("", Route::default().with_cost(5)).unwrap();
6583
6584 let mut video = producer
6585 .consume()
6586 .scope("", &scopes(&["room/video"]))
6587 .unwrap()
6588 .announced();
6589 assert_eq!(video.assert_next_active("").cost, Cost::new(5));
6590 }
6591
6592 #[tokio::test]
6593 async fn route_changes_preserve_each_cursors_visible_winner() {
6594 let producer = origin(1).produce();
6595 let consumer = producer.consume();
6596 let mut all = consumer.clone().announced();
6597 let mut excluded = consumer.clone().excluding(origin(7)).announced();
6598 let mut video = consumer
6599 .clone()
6600 .scope("", &scopes(&["room/live/video"]))
6601 .unwrap()
6602 .announced();
6603 let mut rooted = consumer
6604 .scope("room/live", &Patterns::from(Pattern::all()))
6605 .unwrap()
6606 .announced();
6607
6608 let chat = producer
6609 .scope("", &scopes(&["room/live/chat"]))
6610 .unwrap()
6611 .dynamic("room/live", Route::default().with_hops(hops(&[7])).with_cost(1))
6612 .unwrap();
6613 assert_eq!(all.assert_next_active("room/live").cost, Cost::new(1));
6614 assert_eq!(rooted.assert_next_active("").cost, Cost::new(1));
6615 excluded.assert_next_wait();
6616 video.assert_next_wait();
6617
6618 let video_route = producer
6619 .scope("", &scopes(&["room/live/video"]))
6620 .unwrap()
6621 .dynamic("room/live", Route::default().with_hops(hops(&[8])).with_cost(5))
6622 .unwrap();
6623 all.assert_next_wait();
6624 rooted.assert_next_wait();
6625 assert_eq!(excluded.assert_next_active("room/live").cost, Cost::new(5));
6626 assert_eq!(video.assert_next_active("room/live").cost, Cost::new(5));
6627
6628 chat.update(Route::default().with_hops(hops(&[7])).with_cost(9))
6629 .unwrap();
6630 assert_eq!(all.assert_next_active("room/live").cost, Cost::new(5));
6631 assert_eq!(rooted.assert_next_active("").cost, Cost::new(5));
6632 excluded.assert_next_wait();
6633 video.assert_next_wait();
6634
6635 drop(video_route);
6636 assert_eq!(all.assert_next_active("room/live").cost, Cost::new(9));
6637 assert_eq!(rooted.assert_next_active("").cost, Cost::new(9));
6638 excluded.assert_next_ended("room/live");
6639 video.assert_next_ended("room/live");
6640 }
6641
6642 #[tokio::test]
6643 async fn root_cursor_keeps_the_more_specific_covering_route() {
6644 let producer = origin(1).produce();
6645 let broad = producer.dynamic("room", Route::default().with_cost(1)).unwrap();
6646 let narrow = producer.dynamic("room/live", Route::default().with_cost(9)).unwrap();
6647 let mut announced = producer
6648 .consume()
6649 .scope("room/live/video", &Patterns::from(Pattern::all()))
6650 .unwrap()
6651 .announced();
6652 assert_eq!(announced.assert_next_active("").cost, Cost::new(9));
6653
6654 broad.update(Route::default().with_cost(0)).unwrap();
6655 announced.assert_next_wait();
6656 narrow.update(Route::default().with_cost(8)).unwrap();
6657 assert_eq!(announced.assert_next_active("").cost, Cost::new(8));
6658 drop(narrow);
6659 assert_eq!(announced.assert_next_active("").cost, Cost::new(0));
6660 }
6661
6662 #[tokio::test]
6663 async fn dynamic_accepts_a_max_depth_prefix() {
6664 let producer = origin(1).produce();
6665 let path = (0..Path::MAX_PARTS)
6666 .map(|i| format!("s{i}"))
6667 .collect::<Vec<_>>()
6668 .join("/");
6669 let mut announced = producer.consume().announced();
6670
6671 let dynamic = producer.dynamic(&path, Route::default()).expect("max depth is allowed");
6672 announced.assert_next_active(&path);
6673
6674 let _pending = producer.consume().request_broadcast(&path);
6675 let request = queued(&dynamic).await;
6676 assert_eq!(request.path().as_str(), path);
6677 }
6678
6679 #[tokio::test]
6680 async fn dynamic_exclusion_skips_routes_through_the_subscriber() {
6681 let producer = origin(1).produce();
6682 let _server = producer
6683 .dynamic("live", Route::default().with_hops(hops(&[7])))
6684 .unwrap();
6685
6686 let mut excluded = producer.consume().excluding(origin(7)).announced();
6687 excluded.assert_next_wait();
6688
6689 let mut clean = producer.consume().excluding(origin(8)).announced();
6690 clean.assert_next_active("live");
6691 }
6692
6693 #[tokio::test]
6695 async fn announce_consumer_is_a_stream() {
6696 use futures::StreamExt;
6697 let producer = origin(1).produce();
6698 let server = producer.dynamic("live", Route::default()).unwrap();
6699 let mut announced = producer.consume().announced();
6700 let update = StreamExt::next(&mut announced)
6701 .now_or_never()
6702 .expect("next")
6703 .expect("no next");
6704 assert_eq!(update.prefix.as_str(), "live");
6705 assert_eq!(update.kind, AnnounceKind::Announced);
6706 assert!(StreamExt::next(&mut announced).now_or_never().is_none());
6707 drop(server);
6708 let update = StreamExt::next(&mut announced)
6709 .now_or_never()
6710 .expect("next")
6711 .expect("no next");
6712 assert_eq!(update.kind, AnnounceKind::Retracted);
6713 }
6714
6715 #[tokio::test]
6716 async fn dynamic_retracts() {
6717 let producer = origin(1).produce();
6718 let server = producer.dynamic("live", Route::default()).unwrap();
6719 let mut announced = producer.consume().announced();
6720 announced.assert_next_active("live");
6721
6722 drop(server);
6723 announced.assert_next_ended("live");
6724 }
6725
6726 #[test]
6727 fn charged_wildcard_cost_accumulates_across_hops() {
6728 let first = Cost::new(4).charged(1);
6729 let second = first.charged(2);
6730 assert_eq!(second, Cost { warm: 7, cold: 7 });
6731 }
6732
6733 #[tokio::test]
6734 async fn local_broadcast_is_invisible_until_announced() {
6735 let producer = origin(1).produce();
6736 let mut local = producer.consume().announced();
6737 let mut peer = producer.consume().excluding(Hop::UNKNOWN).announced();
6738 let broadcast = producer.create_broadcast("room/alice").unwrap();
6739 local.assert_next_wait();
6740 peer.assert_next_wait();
6741
6742 broadcast.announce(Route::default()).unwrap();
6743 local.assert_next_active("room/alice");
6744 peer.assert_next_active("room/alice");
6745
6746 drop(broadcast);
6747 local.assert_next_ended("room/alice");
6748 peer.assert_next_ended("room/alice");
6749 }
6750
6751 #[tokio::test]
6752 async fn served_route_materializes_on_demand() {
6753 let producer = origin(1).produce();
6754 let consumer = producer.consume();
6755
6756 let server = producer.dynamic("room", Route::default()).unwrap();
6757
6758 let pending = consumer.request_broadcast("room/alice");
6759 let request = queued(&server).await;
6760 assert_eq!(request.path().as_str(), "room/alice");
6761
6762 let source = broadcast::Info::new().produce();
6763 request.accept(&source);
6764
6765 let resolved = pending.await.expect("resolves");
6766 assert_eq!(resolved.info().path.as_str(), "room/alice");
6768
6769 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
6771 assert!(again.is_clone(&resolved));
6772 }
6773
6774 #[tokio::test]
6775 async fn served_requests_coalesce() {
6776 let producer = origin(1).produce();
6777 let consumer = producer.consume();
6778 let server = producer.dynamic("room", Route::default()).unwrap();
6779
6780 let first = consumer.request_broadcast("room/alice");
6781 let second = consumer.request_broadcast("room/alice");
6782
6783 let request = queued(&server).await;
6784 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
6786
6787 let source = broadcast::Info::new().produce();
6788 request.accept(&source);
6789
6790 let first = first.await.expect("resolves");
6791 let second = second.await.expect("resolves");
6792 assert!(first.is_clone(&second));
6793 }
6794
6795 #[tokio::test]
6796 async fn retract_rejects_pending_requests() {
6797 let producer = origin(1).produce();
6798 let consumer = producer.consume();
6799 let server = producer.dynamic("room", Route::default()).unwrap();
6800
6801 let pending = consumer.request_broadcast("room/alice");
6802 drop(server);
6803
6804 let err = pending.await.err().unwrap();
6805 assert!(matches!(err, Error::Unroutable));
6806
6807 let err = consumer
6809 .request_broadcast("room/alice")
6810 .now_or_never()
6811 .expect("unroutable")
6812 .err()
6813 .unwrap();
6814 assert!(matches!(err, Error::Unroutable));
6815 }
6816
6817 #[tokio::test]
6818 async fn routed_broadcast_survives_serving_route_retraction() {
6819 let producer = origin(1).produce();
6820 let consumer = producer.consume();
6821
6822 let standby_server = producer.dynamic("room", Route::default()).unwrap();
6825 let second_server = producer.dynamic("room", Route::default()).unwrap();
6826 let incumbent_server = producer.dynamic("room", Route::default()).unwrap();
6827
6828 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
6829 assert!((&mut resolving).now_or_never().is_none());
6830
6831 drop(incumbent_server);
6837 assert!((&mut resolving).now_or_never().is_none());
6838 drop(second_server);
6839 assert!((&mut resolving).now_or_never().is_none());
6840
6841 let request = queued(&standby_server).await;
6842 let source = broadcast::Info::new().produce();
6843 request.accept(&source);
6844
6845 let resolved = resolving.await.expect("resolves via the standby");
6846 assert_eq!(resolved.info().path.as_str(), "room/alice");
6847 }
6848
6849 #[tokio::test]
6850 async fn split_horizon_skips_routes_through_the_requester() {
6851 let producer = origin(1).produce();
6852 let _server = producer
6853 .dynamic("room", Route::default().with_hops(hops(&[7])))
6854 .unwrap();
6855
6856 let excluded = producer.consume().excluding(origin(7));
6858 let err = excluded
6859 .request_broadcast("room/alice")
6860 .now_or_never()
6861 .expect("unroutable")
6862 .err()
6863 .unwrap();
6864 assert!(matches!(err, Error::Unroutable));
6865
6866 let clean = producer.consume().excluding(origin(8));
6868 let pending = clean.request_broadcast("room/alice");
6869 assert!(pending.now_or_never().is_none());
6870 }
6871
6872 #[tokio::test]
6873 async fn routes_report_where_they_entered() {
6874 let producer = origin(1).produce();
6875 let peer = producer.clone().peer();
6876 let mut announced = producer.consume().announced();
6877
6878 let _ingest = producer
6879 .dynamic("client", Route::default().with_hops(hops(&[5])).with_via(origin(5)))
6880 .unwrap();
6881 let _gateway = producer.publish("gateway", Route::default()).unwrap();
6882 let _forwarded = peer
6883 .dynamic(
6884 "forwarded",
6885 Route::default().with_hops(hops(&[5, 7])).with_via(origin(7)),
6886 )
6887 .unwrap();
6888
6889 assert_eq!(announced.assert_next_active("client").source(), Source::Local);
6890 assert_eq!(
6891 announced.assert_next_active("forwarded").source(),
6892 Source::Peer(origin(7))
6893 );
6894 assert_eq!(announced.assert_next_active("gateway").source(), Source::Local);
6895
6896 let scoped = peer.scope("room", &Patterns::from(Pattern::all())).unwrap();
6898 let _nested = scoped.dynamic("x", Route::default().with_via(origin(8))).unwrap();
6899 assert_eq!(announced.assert_next_active("room/x").source(), Source::Peer(origin(8)));
6900 }
6901
6902 #[tokio::test(start_paused = true)]
6906 async fn try_next_delivers_a_held_update() {
6907 async fn step(by: Duration) {
6909 tokio::time::advance(by).await;
6910 for _ in 0..10 {
6911 tokio::task::yield_now().await;
6912 }
6913 }
6914
6915 let producer = origin(1).produce();
6916 let peer = producer.clone().peer();
6917 let mut announced = producer.consume().announced();
6918 step(Duration::from_millis(1)).await;
6919
6920 let route = |chain: &[u64], via| Route::default().with_hops(hops(chain)).with_via(origin(via));
6921 let near = peer.dynamic("room", route(&[9, 2], 2)).unwrap();
6922 assert!(announced.try_next().unwrap().kind.is_active());
6923
6924 let _far = peer.dynamic("room", route(&[9, 3, 4], 4)).unwrap();
6925 drop(near);
6926 assert!(announced.try_next().is_none(), "the update is held");
6927 step(Duration::from_millis(1)).await;
6928 assert!(announced.try_next().is_none(), "the update is still held");
6929 step(DEFAULT_UPDATE_HOLD).await;
6930 let update = announced.try_next().expect("the hold passed");
6931 assert_eq!(update.kind, AnnounceKind::Updated);
6932 assert_eq!(update.route.hops, hops(&[9, 3, 4]));
6933 }
6934
6935 #[tokio::test(start_paused = true)]
6938 async fn a_changed_route_waits_out_the_hold() {
6939 let producer = origin(1).produce();
6940 let peer = producer.clone().peer();
6941 let mut announced = producer.consume().announced();
6942 tokio::time::sleep(Duration::from_millis(1)).await;
6944
6945 let route = |chain: &[u64], via| Route::default().with_hops(hops(chain)).with_via(origin(via));
6946 let near = peer.dynamic("room", route(&[9, 2], 2)).unwrap();
6947 announced.assert_next_active("room");
6948
6949 let far = peer.dynamic("room", route(&[9, 3, 4], 4)).unwrap();
6951 let farther = peer.dynamic("room", route(&[9, 5, 6, 7], 7)).unwrap();
6952 let start = tokio::time::Instant::now();
6953 drop(near);
6954 drop(far);
6955 let update = announced.next().await.unwrap();
6956 assert_eq!(update.kind, AnnounceKind::Updated);
6957 assert_eq!(update.route.hops, hops(&[9, 5, 6, 7]));
6958 assert_eq!(start.elapsed(), DEFAULT_UPDATE_HOLD);
6959 announced.assert_next_wait();
6960
6961 drop(farther);
6963 announced.assert_next_ended("room");
6964 }
6965
6966 #[tokio::test]
6970 async fn withdrawn_peer_hides_routes_through_it() {
6971 let producer = origin(1).produce();
6972 let peer = producer.clone().peer();
6973 let mut announced = producer.consume().announced();
6974
6975 let direct = || Route::default().with_hops(hops(&[9, 2])).with_via(origin(2));
6977 let first = peer.dynamic("room", direct()).unwrap();
6978 let relayed = peer
6979 .dynamic("room", Route::default().with_hops(hops(&[9, 2, 3])).with_via(origin(3)))
6980 .unwrap();
6981 announced.assert_next_active("room");
6982 announced.assert_next_wait();
6983
6984 first.withdrawn();
6986 announced.assert_next_ended("room");
6987 announced.assert_next_wait();
6988
6989 let second = peer.dynamic("room", direct()).unwrap();
6991 announced.assert_next_active("room");
6992 assert!(producer.shared.lock().withdrawn.is_empty());
6993
6994 second.withdrawn();
6996 drop(relayed);
6997 announced.assert_next_ended("room");
6998 assert!(producer.shared.lock().withdrawn.is_empty());
6999 }
7000
7001 #[tokio::test]
7003 async fn old_session_withdrawal_keeps_newer_route() {
7004 for restart in [false, true] {
7005 let producer = origin(1).produce();
7006 let peer = producer.clone().peer();
7007 let mut announced = producer.consume().announced();
7008 let route = Route::default().with_hops(hops(&[9, 2]));
7009 let first = peer.dynamic("room", route.clone()).unwrap();
7010 let second = peer.dynamic("room", route.clone()).unwrap();
7011 announced.assert_next_active("room");
7012 announced.assert_next_wait();
7013 let remaining = if restart {
7014 first.update(route).unwrap();
7016 second.withdrawn();
7017 first
7018 } else {
7019 first.withdrawn();
7020 second
7021 };
7022 announced.assert_next_wait();
7023 assert!(producer.shared.lock().withdrawn.is_empty());
7024 remaining.withdrawn();
7025 announced.assert_next_ended("room");
7026 assert!(producer.shared.lock().withdrawn.is_empty());
7027 }
7028 }
7029
7030 #[tokio::test]
7031 async fn withdrawn_route_no_longer_covers_requests() {
7032 let producer = origin(1).produce();
7033 let peer = producer.clone().peer();
7034 let direct = peer.dynamic("room", Route::default().with_hops(hops(&[9, 2]))).unwrap();
7035 let _relayed = peer
7036 .dynamic("room", Route::default().with_hops(hops(&[9, 2, 3])))
7037 .unwrap();
7038 let path = Path::new("room/video");
7039 let id = producer
7040 .shared
7041 .read()
7042 .routes
7043 .covering(&path)
7044 .find(|entry| entry.hops.iter().last() == Some(&origin(3)))
7045 .unwrap()
7046 .id;
7047 assert!(producer.shared.read().routes.covers(&path, id));
7048 direct.withdrawn();
7049 assert!(!producer.shared.read().routes.covers(&path, id));
7050 }
7051
7052 #[tokio::test]
7055 async fn source_change_is_an_update() {
7056 let producer = origin(1).produce();
7057 let peer = producer.clone().peer();
7058 let mut announced = producer.consume().announced();
7059
7060 let route = Route::default().with_hops(hops(&[7])).with_via(origin(7));
7061 let _forwarded = peer.dynamic("room", route.clone()).unwrap();
7062 assert_eq!(announced.assert_next_active("room").source(), Source::Peer(origin(7)));
7063
7064 let local = producer.dynamic("room", route).unwrap();
7066 let update = announced.next().now_or_never().expect("next blocked").expect("no next");
7067 assert_eq!(update.kind, AnnounceKind::Updated);
7068 assert_eq!(update.route.source(), Source::Local);
7069
7070 drop(local);
7071 assert_eq!(announced.assert_next_active("room").source(), Source::Peer(origin(7)));
7072 }
7073
7074 #[tokio::test]
7075 async fn local_view_hides_peer_routes() {
7076 let producer = origin(1).produce();
7077 let peer = producer.clone().peer();
7078 let mut local = producer.consume().local().announced();
7079
7080 let _forwarded = peer
7081 .dynamic("remote", Route::default().with_hops(hops(&[7])).with_via(origin(7)))
7082 .unwrap();
7083 local.assert_next_wait();
7084
7085 let _shadow = peer
7089 .dynamic("both", Route::default().with_hops(hops(&[7])).with_via(origin(7)))
7090 .unwrap();
7091 let ingest = producer
7092 .dynamic(
7093 "both",
7094 Route::default().with_hops(hops(&[5])).with_via(origin(5)).with_cost(9),
7095 )
7096 .unwrap();
7097 assert_eq!(local.assert_next_active("both").source(), Source::Local);
7098 drop(ingest);
7099 local.assert_next_ended("both");
7100
7101 let err = producer
7104 .consume()
7105 .local()
7106 .request_broadcast("remote/alice")
7107 .now_or_never()
7108 .expect("unroutable")
7109 .err()
7110 .unwrap();
7111 assert!(matches!(err, Error::Unroutable));
7112 assert!(
7113 producer
7114 .consume()
7115 .request_broadcast("remote/alice")
7116 .now_or_never()
7117 .is_none()
7118 );
7119 }
7120
7121 #[tokio::test]
7125 async fn handler_rejection_is_final() {
7126 let producer = origin(1).produce();
7127 let consumer = producer.consume();
7128 let server = producer.dynamic("room", Route::default()).unwrap();
7129
7130 let pending = consumer.request_broadcast("room/alice");
7131 let request = queued(&server).await;
7132 request.reject(Error::Unroutable);
7133 let err = tokio::time::timeout(Duration::from_secs(5), pending)
7134 .await
7135 .expect("the front must give up, not spin")
7136 .err()
7137 .unwrap();
7138 assert!(matches!(err, Error::Unroutable));
7139
7140 let pending = consumer.request_broadcast("room/bob");
7142 let request = queued(&server).await;
7143 assert_eq!(request.path().as_str(), "room/bob");
7144 let served = broadcast::Info::new().produce();
7145 request.accept(&served);
7146 pending.await.expect("resolves");
7147 }
7148
7149 #[tokio::test]
7152 async fn routed_broadcast_waits_out_a_rejection() {
7153 let producer = origin(1).produce();
7154 let consumer = producer.consume();
7155 let server = producer.dynamic("room", Route::default()).unwrap();
7156
7157 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
7158 assert!((&mut resolving).now_or_never().is_none());
7159 let request = queued(&server).await;
7160 request.reject(Error::Unroutable);
7161
7162 for _ in 0..20 {
7164 tokio::task::yield_now().await;
7165 }
7166 assert!((&mut resolving).now_or_never().is_none());
7167 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
7168
7169 server.update(Route::default().with_cost(2)).unwrap();
7171 assert!((&mut resolving).now_or_never().is_none());
7172 let request = queued(&server).await;
7173 let served = broadcast::Info::new().produce();
7174 request.accept(&served);
7175 resolving.await.expect("resolves");
7176 }
7177
7178 #[tokio::test]
7181 async fn routed_broadcast_reports_teardown_as_closed() {
7182 let (producer, driver) = Producer::new(Config::new(origin(1)));
7183 let consumer = producer.consume();
7184 let _server = producer.dynamic("room", Route::default()).unwrap();
7185
7186 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
7188 assert!((&mut resolving).now_or_never().is_none());
7189
7190 drop(driver);
7191
7192 let err = tokio::time::timeout(Duration::from_secs(5), resolving)
7193 .await
7194 .expect("teardown resolves the wait")
7195 .err()
7196 .unwrap();
7197 assert!(matches!(err, Error::Closed), "unexpected end: {err}");
7198 }
7199
7200 #[tokio::test]
7203 async fn routed_broadcast_wakes_for_a_local_broadcast() {
7204 let producer = origin(1).produce();
7205 let consumer = producer.consume();
7206 let server = producer.dynamic("room", Route::default()).unwrap();
7207
7208 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
7209 assert!((&mut resolving).now_or_never().is_none());
7210 queued(&server).await.reject(Error::Unroutable);
7211 for _ in 0..20 {
7212 tokio::task::yield_now().await;
7213 }
7214 assert!((&mut resolving).now_or_never().is_none());
7215
7216 let _local = producer.publish("room/alice", Route::default()).unwrap();
7218 let resolved = resolving.await.expect("resolves locally");
7219 assert_eq!(resolved.info().path.as_str(), "room/alice");
7220 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
7221 }
7222
7223 #[tokio::test]
7226 async fn late_track_on_a_served_front_replays() {
7227 let producer = origin(1).produce();
7228 let consumer = producer.consume();
7229 let server = producer.dynamic("room", Route::default()).unwrap();
7230
7231 let source = broadcast::Info::new().produce();
7232 for name in ["a", "b"] {
7233 let track = source.create_track(name, None).unwrap();
7234 let mut group = track.append_group().unwrap();
7235 group.write_frame(crate::Timestamp::ZERO, name.as_bytes()).unwrap();
7236 group.finish().unwrap();
7237 std::mem::forget(track);
7239 }
7240
7241 let pending = consumer.request_broadcast("room/alice");
7242 queued(&server).await.accept(&source);
7243 let resolved = pending.await.expect("resolves");
7244
7245 let budget = track::Subscription::default().with_max_age(Duration::from_secs(3600));
7246 for name in ["a", "b"] {
7247 let mut subscription = resolved
7248 .track(name)
7249 .unwrap()
7250 .subscribe(budget.clone())
7251 .await
7252 .expect("subscribe");
7253 let mut group = tokio::time::timeout(Duration::from_secs(5), subscription.recv_group())
7254 .await
7255 .expect("the late track must replay, not park")
7256 .expect("recv group")
7257 .expect("track ended early");
7258 let frame = group.read_frame().await.expect("read frame").expect("frame");
7259 assert_eq!(&frame.payload[..], name.as_bytes());
7260 }
7261 }
7262
7263 #[tokio::test]
7264 async fn most_specific_prefix_shadows() {
7265 let producer = origin(1).produce();
7266 let consumer = producer.consume();
7267
7268 let broad_server = producer.dynamic("", Route::default()).unwrap();
7269 let _narrow = producer.announce(".dash", Route::default()).unwrap();
7272
7273 let err = consumer
7274 .request_broadcast(".dash/pid")
7275 .now_or_never()
7276 .expect("unroutable")
7277 .err()
7278 .unwrap();
7279 assert!(matches!(err, Error::Unroutable));
7280
7281 let _pending = consumer.request_broadcast("room/alice");
7283 let request = queued(&broad_server).await;
7284 assert_eq!(request.path().as_str(), "room/alice");
7285 }
7286
7287 #[tokio::test]
7288 async fn root_dynamic_serves_any_path() {
7289 let producer = origin(1).produce();
7290 let consumer = producer.consume();
7291 let mut announced = consumer.announced();
7292 let dynamic = producer.dynamic("", Route::default()).unwrap();
7293 announced.assert_next_active("");
7295
7296 let pending = consumer.request_broadcast("anything/at/all");
7297 let request = queued(&dynamic).await;
7298 assert_eq!(request.path().as_str(), "anything/at/all");
7299
7300 let source = broadcast::Info::new().produce();
7301 request.accept(&source);
7302 let resolved = pending.await.expect("resolves");
7303 assert_eq!(resolved.info().path.as_str(), "anything/at/all");
7304
7305 drop(dynamic);
7307 announced.assert_next_ended("");
7308 let err = consumer
7309 .request_broadcast("something/else")
7310 .now_or_never()
7311 .expect("unroutable")
7312 .err()
7313 .unwrap();
7314 assert!(matches!(err, Error::Unroutable));
7315 }
7316
7317 #[tokio::test]
7323 async fn out_of_scope_request_never_reaches_the_dynamic_handler() {
7324 let producer = origin(1).produce();
7325 let dynamic = producer.dynamic("", Route::default()).unwrap();
7326 let scoped = producer.consume().scope("", &scopes(&["tenant-a"])).unwrap();
7327
7328 for path in ["tenant-b/live", "tenant-a-other/live"] {
7331 let refused = scoped
7332 .request_broadcast(path)
7333 .now_or_never()
7334 .expect("an out-of-scope request must be refused synchronously, not queued");
7335 assert!(matches!(refused, Err(Error::Unauthorized)));
7336 assert!(
7337 dynamic.requested_broadcast().now_or_never().is_none(),
7338 "the dynamic handler was asked to create a broadcast the requester may not read"
7339 );
7340 }
7341 }
7342
7343 #[tokio::test]
7344 async fn routed_waits_for_coverage() {
7345 let producer = origin(1).produce();
7346 let consumer = producer.consume();
7347
7348 let mut fut = consumer.routed("room/alice").boxed();
7349 assert!((&mut fut).now_or_never().is_none());
7350
7351 let _a = producer.announce("room", Route::default().with_cost(3)).unwrap();
7353 let route = fut.now_or_never().expect("covered").expect("routed");
7354 assert_eq!(route.cost, Cost::new(3));
7355
7356 consumer
7358 .routed("room/alice/cam")
7359 .now_or_never()
7360 .expect("covered")
7361 .expect("routed");
7362 }
7363
7364 #[tokio::test]
7365 async fn routed_ignores_deeper_routes() {
7366 let producer = origin(1).produce();
7367 let consumer = producer.consume();
7368
7369 let _deep = producer.announce("room/alice/cam", Route::default()).unwrap();
7371 let mut fut = consumer.routed("room/alice").boxed();
7372 assert!((&mut fut).now_or_never().is_none());
7373
7374 let _exact = producer.announce("room/alice", Route::default()).unwrap();
7375 fut.now_or_never().expect("covered").expect("routed");
7376 }
7377
7378 #[tokio::test]
7379 async fn routed_accepts_a_max_depth_path() {
7380 let producer = origin(1).produce();
7381 let consumer = producer.consume();
7382 let path = (0..Path::MAX_PARTS)
7383 .map(|i| format!("s{i}"))
7384 .collect::<Vec<_>>()
7385 .join("/");
7386 assert_eq!(Path::new(&path).parts().count(), Path::MAX_PARTS);
7387
7388 assert!(consumer.allowed().matches(&path));
7389
7390 let mut fut = consumer.routed(&path).boxed();
7391 assert!((&mut fut).now_or_never().is_none());
7392
7393 let _a = producer.announce("", Route::default()).unwrap();
7395 fut.now_or_never().expect("covered").expect("routed");
7396 }
7397
7398 #[tokio::test]
7399 async fn teardown_ends_everything() {
7400 let (producer, driver) = Producer::new(Config::new(origin(1)));
7401 let consumer = producer.consume();
7402 let _announcement = producer.announce("room", Route::default()).unwrap();
7403 let mut announced = consumer.announced();
7404 announced.assert_next_active("room");
7405
7406 let _server = producer.dynamic("served", Route::default()).unwrap();
7407 let pending = consumer.request_broadcast("served/path");
7408
7409 drop(driver);
7410
7411 announced.assert_next_active("served");
7413 assert!(announced.next().now_or_never().expect("ended").is_none());
7414
7415 assert!(pending.now_or_never().expect("rejected").is_err());
7417 assert!(matches!(producer.announce("x", Route::default()), Err(Error::Closed)));
7418 assert!(matches!(producer.create_broadcast("x"), Err(Error::Closed)));
7419 let err = consumer
7420 .request_broadcast("y")
7421 .now_or_never()
7422 .expect("closed")
7423 .err()
7424 .unwrap();
7425 assert!(matches!(err, Error::Closed));
7426
7427 let mut late = consumer.announced();
7429 assert!(late.next().now_or_never().expect("ended").is_none());
7430 }
7431
7432 struct ResumeRig {
7435 producer: Producer,
7436 resolved: broadcast::Consumer,
7437 subscription: track::Subscriber,
7438 incumbent_track: track::Producer,
7441 }
7442
7443 impl ResumeRig {
7444 async fn new(first: &[u64]) -> (Self, Dynamic, broadcast::Producer) {
7447 let producer = origin(1).produce();
7448 let consumer = producer.consume();
7449
7450 let server = producer
7451 .dynamic("room", Route::default().with_hops(hops(first)))
7452 .unwrap();
7453
7454 let pending = consumer.request_broadcast("room/alice");
7455 let request = queued(&server).await;
7456 let source = broadcast::Info::new().produce();
7457 let track = source.create_track("video", None).unwrap();
7458 let mut group = track.append_group().unwrap();
7459 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
7460 group.finish().unwrap();
7461 request.accept(&source);
7462
7463 let resolved = pending.await.expect("resolves");
7464 let mut subscription = resolved
7465 .track("video")
7466 .unwrap()
7467 .subscribe(None)
7468 .await
7469 .expect("subscribe");
7470 let mut group = subscription
7471 .recv_group()
7472 .await
7473 .expect("recv group")
7474 .expect("track ended early");
7475 let frame = group.read_frame().await.expect("read frame").expect("frame");
7476 assert_eq!(&frame.payload[..], b"before");
7477
7478 (
7479 Self {
7480 producer,
7481 resolved,
7482 subscription,
7483 incumbent_track: track,
7484 },
7485 server,
7486 source,
7487 )
7488 }
7489
7490 fn standby(&self, first: &[u64]) -> Dynamic {
7493 self.producer
7494 .dynamic("room", Route::default().with_hops(hops(first)))
7495 .unwrap()
7496 }
7497 }
7498
7499 async fn assert_resumes(rig: &mut ResumeRig, server: &Dynamic) {
7504 let request = queued(server).await;
7505 let replacement = broadcast::Info::new().produce();
7506 let track = replacement.create_track("video", None).unwrap();
7507 let mut group = track.append_group().unwrap();
7510 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
7511 group.finish().unwrap();
7512 request.accept(&replacement);
7513
7514 let mut group = track.append_group().unwrap();
7515 group.write_frame(crate::Timestamp::ZERO, b"resumed".as_ref()).unwrap();
7516 group.finish().unwrap();
7517
7518 let mut group = rig
7519 .subscription
7520 .recv_group()
7521 .await
7522 .expect("subscription survives the failover")
7523 .expect("track ended early");
7524 let frame = group.read_frame().await.expect("read frame").expect("frame");
7525 assert_eq!(&frame.payload[..], b"resumed");
7526 }
7527
7528 #[tokio::test]
7531 async fn driver_resolves_with_live_consumers() {
7532 let (producer, driver) = Producer::new(Config::new(origin(1)));
7533 let consumer = producer.consume();
7534 let run = crate::time::run(driver);
7535 drop(producer);
7536 tokio::time::timeout(Duration::from_secs(5), run)
7537 .await
7538 .expect("driver must finish once the producers are gone");
7539 drop(consumer);
7540 }
7541
7542 #[tokio::test]
7543 async fn remote_source_resumes_through_same_first_hop() {
7544 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
7545 let standby_server = rig.standby(&[10, 20]);
7546
7547 drop(incumbent);
7549 drop(source);
7550
7551 assert_resumes(&mut rig, &standby_server).await;
7553 }
7554
7555 #[tokio::test]
7559 async fn incompatible_successor_is_refused() {
7560 for replacement in [
7561 track::Info::default().with_timescale(crate::Timescale::MICRO),
7562 track::Info::default().with_priority(7),
7563 track::Info::default().with_max_age(Duration::from_secs(7)),
7564 ] {
7565 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
7566 let standby_server = rig.standby(&[10, 20]);
7567 drop(incumbent);
7568 drop(source);
7569
7570 let request = queued(&standby_server).await;
7573 let successor = broadcast::Info::new().produce();
7574 let track = successor.create_track("video", replacement).unwrap();
7575 let mut group = track.append_group().unwrap();
7576 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
7577 group.finish().unwrap();
7578 request.accept(&successor);
7579
7580 assert!(
7581 matches!(rig.subscription.recv_group().await, Err(Error::Unsupported)),
7582 "the subscription must abort rather than resume onto incompatible metadata"
7583 );
7584
7585 let reopened = rig.resolved.track("video").unwrap();
7587 assert!(matches!(reopened.query().await, Err(Error::Unsupported)));
7588 assert!(matches!(reopened.subscribe(None).await, Err(Error::Unsupported)));
7589 }
7590 }
7591
7592 #[tokio::test]
7593 async fn different_first_hop_ends_the_subscription() {
7594 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
7595 let rival_server = rig.standby(&[11]);
7597
7598 drop(incumbent);
7601 drop(source);
7602 rig.incumbent_track.abort(Error::Dropped).unwrap();
7603
7604 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
7606 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
7607
7608 let consumer = rig.producer.consume();
7610 let pending = consumer.request_broadcast("room/alice");
7611 let request = queued(&rival_server).await;
7612 let replacement = broadcast::Info::new().produce();
7613 request.accept(&replacement);
7614 pending.await.expect("re-request resolves through the rival");
7615 }
7616
7617 #[tokio::test]
7622 async fn a_first_hop_update_drains_the_old_publisher() {
7623 for first in [&[][..], &[0][..], &[10][..]] {
7624 let (mut rig, server, _source) = ResumeRig::new(first).await;
7625 server.update(Route::default().with_hops(hops(&[11]))).unwrap();
7626
7627 let mut group = rig.incumbent_track.append_group().unwrap();
7629 group.write_frame(crate::Timestamp::ZERO, b"draining".as_ref()).unwrap();
7630 group.finish().unwrap();
7631 let mut group = next_group(&mut rig.subscription)
7632 .await
7633 .expect("the in-flight subscription survives the update")
7634 .expect("track ended early");
7635 let frame = group.read_frame().await.expect("read frame").expect("frame");
7636 assert_eq!(&frame.payload[..], b"draining", "first hop {first:?}");
7637
7638 let consumer = rig.producer.consume();
7641 let pending = consumer.request_broadcast("room/alice");
7642 let request = queued(&server).await;
7643 let replacement = broadcast::Info::new().produce();
7644 let track = replacement
7645 .create_track("video", track::Info::default().with_priority(7))
7646 .unwrap();
7647 let mut group = track.append_group().unwrap();
7648 group.write_frame(crate::Timestamp::ZERO, b"new".as_ref()).unwrap();
7649 group.finish().unwrap();
7650 request.accept(&replacement);
7651
7652 let resolved = pending.await.expect("resolves through the new publisher");
7653 assert!(
7654 !resolved.is_clone(&rig.resolved),
7655 "first hop {first:?} joined the old front"
7656 );
7657 let mut fresh = resolved
7658 .track("video")
7659 .unwrap()
7660 .subscribe(None)
7661 .await
7662 .expect("the old publisher's track info does not apply");
7663 let mut group = next_group(&mut fresh)
7664 .await
7665 .expect("recv group")
7666 .expect("track ended early");
7667 let frame = group.read_frame().await.expect("read frame").expect("frame");
7668 assert_eq!(&frame.payload[..], b"new");
7669
7670 rig.incumbent_track.finish().unwrap();
7673 let end = next_group(&mut rig.subscription).await;
7674 assert!(
7675 !matches!(end, Ok(Some(_))),
7676 "first hop {first:?} spliced the new publisher into a live subscription"
7677 );
7678 }
7679 }
7680
7681 #[tokio::test]
7684 async fn a_first_hop_update_resumes_through_the_same_publisher() {
7685 let (mut rig, incumbent, _source) = ResumeRig::new(&[10]).await;
7686 let standby_server = rig.standby(&[10, 20]);
7687
7688 incumbent.update(Route::default().with_hops(hops(&[11]))).unwrap();
7689
7690 assert_resumes(&mut rig, &standby_server).await;
7691 assert!(
7692 incumbent.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
7693 "the front never asks the new publisher"
7694 );
7695 }
7696
7697 #[tokio::test]
7698 async fn anonymous_routes_never_resume() {
7699 let (mut rig, incumbent, source) = ResumeRig::new(&[]).await;
7702 let _twin_server = rig.standby(&[]);
7703
7704 drop(incumbent);
7705 drop(source);
7706 rig.incumbent_track.abort(Error::Dropped).unwrap();
7707
7708 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
7709 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
7710 }
7711
7712 #[tokio::test]
7719 async fn anonymous_handoff_serves_the_newcomer_immediately() {
7720 let producer = origin(1).produce();
7721
7722 let server_a = producer
7724 .dynamic("room", Route::default().with_hops(hops(&[10])))
7725 .unwrap();
7726
7727 let consumer = producer.consume().excluding(origin(30));
7730 let pending = consumer.request_broadcast("room/alice");
7731 let request = queued(&server_a).await;
7732 let source_a = broadcast::Info::new().produce();
7733 let track_a = source_a.create_track("video", None).unwrap();
7734 let mut group = track_a.append_group().unwrap();
7735 group.write_frame(crate::Timestamp::ZERO, b"from-a".as_ref()).unwrap();
7736 group.finish().unwrap();
7737 request.accept(&source_a);
7738
7739 let resolved_a = pending.await.expect("resolves");
7740 let mut sub_a = resolved_a
7741 .track("video")
7742 .unwrap()
7743 .subscribe(None)
7744 .await
7745 .expect("subscribe");
7746 let mut group = sub_a
7747 .recv_group()
7748 .await
7749 .expect("recv group")
7750 .expect("track ended early");
7751 assert_eq!(
7752 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
7753 b"from-a"
7754 );
7755
7756 drop(track_a);
7759 drop(source_a);
7760 drop(server_a);
7761
7762 let err = sub_a.recv_group().await.err().expect("front closed");
7764 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
7765
7766 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
7770 settle(|| {
7771 matches!(
7772 consumer.request_broadcast("room/alice").now_or_never(),
7773 Some(Err(Error::Unroutable))
7774 )
7775 })
7776 .await;
7777
7778 let server_b = producer
7780 .dynamic("room", Route::default().with_hops(hops(&[20])))
7781 .unwrap();
7782 let pending = consumer.request_broadcast("room/alice");
7783 let request = queued(&server_b).await;
7784 let source_b = broadcast::Info::new().produce();
7785 let track_b = source_b.create_track("video", None).unwrap();
7786 let mut group = track_b.append_group().unwrap();
7787 group.write_frame(crate::Timestamp::ZERO, b"from-b".as_ref()).unwrap();
7788 group.finish().unwrap();
7789 request.accept(&source_b);
7790
7791 let resolved_b = pending.await.expect("B's front is served immediately");
7792 assert!(
7793 !resolved_b.is_clone(&resolved_a),
7794 "B must not splice into A's closed front"
7795 );
7796
7797 let mut sub_b = resolved_b
7798 .track("video")
7799 .unwrap()
7800 .subscribe(None)
7801 .await
7802 .expect("subscribe");
7803 let mut group = sub_b
7804 .recv_group()
7805 .await
7806 .expect("recv group")
7807 .expect("track ended early");
7808 assert_eq!(
7809 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
7810 b"from-b"
7811 );
7812 }
7813
7814 #[tokio::test]
7815 async fn reprice_is_invisible_to_the_subscription() {
7816 let (rig, incumbent, source) = ResumeRig::new(&[10]).await;
7817
7818 incumbent
7821 .update(Route::default().with_hops(hops(&[10])).with_cost(9))
7822 .unwrap();
7823
7824 let track = source.create_track("audio", None).unwrap();
7825 let mut group = track.append_group().unwrap();
7826 group.write_frame(crate::Timestamp::ZERO, b"steady".as_ref()).unwrap();
7827 group.finish().unwrap();
7828
7829 let mut audio = rig
7830 .resolved
7831 .track("audio")
7832 .unwrap()
7833 .subscribe(None)
7834 .await
7835 .expect("subscribe survives the reprice");
7836 let mut group = audio
7837 .recv_group()
7838 .await
7839 .expect("recv group")
7840 .expect("track ended early");
7841 let frame = group.read_frame().await.expect("read frame").expect("frame");
7842 assert_eq!(&frame.payload[..], b"steady");
7843 }
7844
7845 #[tokio::test]
7846 async fn drain_reprice_migrates_before_the_session_dies() {
7847 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
7848 let standby_server = rig.standby(&[10, 20]);
7849
7850 incumbent
7854 .update(Route::default().with_hops(hops(&[10])).with_cost(Cost::DRAIN))
7855 .unwrap();
7856
7857 assert_resumes(&mut rig, &standby_server).await;
7858
7859 drop(incumbent);
7861 drop(source);
7862 }
7863
7864 #[tokio::test]
7865 async fn local_sources_splice_newest_first() {
7866 let producer = origin(1).produce();
7867 let consumer = producer.consume();
7868
7869 let first = producer.publish("room/alice", Route::default()).unwrap();
7870 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7871
7872 let second = producer.publish("room/alice", Route::default()).unwrap();
7874 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
7875 assert!(again.is_clone(&resolved));
7876
7877 first.close();
7879 settle(|| consumer.get_broadcast("room/alice").is_some()).await;
7880 second.close();
7881 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
7882
7883 let _third = producer.publish("room/alice", Route::default()).unwrap();
7885 assert!(consumer.get_broadcast("room/alice").is_some());
7886 }
7887
7888 #[tokio::test]
7896 async fn a_finished_broadcast_concludes_in_flight_subscriptions() {
7897 let producer = origin(1).produce();
7898 let consumer = producer.consume();
7899
7900 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
7901 let track = broadcast.create_track("video", None).unwrap();
7902
7903 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7904 let mut subscription = resolved
7905 .track("video")
7906 .unwrap()
7907 .subscribe(None)
7908 .await
7909 .expect("subscribe");
7910 let mut group = track.append_group().unwrap();
7912 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
7913 group.finish().unwrap();
7914 track.finish().unwrap();
7915 drop(track);
7916 broadcast.close();
7917
7918 let mut group = next_group(&mut subscription)
7919 .await
7920 .expect("a cleanly finished track was served as an error")
7921 .expect("the track ended before its last group");
7922 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
7923 drop(group);
7924
7925 let end = next_group(&mut subscription)
7926 .await
7927 .expect("a cleanly finished track ended as an error");
7928 assert!(end.is_none(), "a group followed the final one");
7929 }
7930
7931 #[tokio::test]
7938 async fn a_retracted_route_concludes_in_flight_subscriptions() {
7939 let producer = origin(1).produce();
7940 let consumer = producer.consume();
7941 let server = producer
7942 .dynamic("room", Route::default().with_hops(hops(&[10])))
7943 .unwrap();
7944
7945 let pending = consumer.request_broadcast("room/alice");
7946 let request = queued(&server).await;
7947 let source = broadcast::Info::new().produce();
7948 let track = source.create_track("video", None).unwrap();
7949 request.accept(&source);
7950
7951 let resolved = pending.await.expect("resolves");
7952 let mut subscription = resolved
7953 .track("video")
7954 .unwrap()
7955 .subscribe(None)
7956 .await
7957 .expect("subscribe");
7958
7959 source.close();
7962 drop(server);
7963 settle(|| resolved.is_closed()).await;
7964 let mut group = track.append_group().unwrap();
7965 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
7966 group.finish().unwrap();
7967 track.finish().unwrap();
7968 drop(track);
7969
7970 let mut group = next_group(&mut subscription)
7971 .await
7972 .expect("a retracted route's track was served as an error")
7973 .expect("the track ended before its last group");
7974 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
7975 drop(group);
7976
7977 let end = next_group(&mut subscription)
7978 .await
7979 .expect("a cleanly finished track ended as an error");
7980 assert!(end.is_none(), "a group followed the final one");
7981 }
7982
7983 #[tokio::test]
7986 async fn a_closed_source_is_not_requested_again_from_its_standing_route() {
7987 let producer = origin(1).produce();
7988 let server = producer
7989 .dynamic("room", Route::default().with_hops(hops(&[10])))
7990 .unwrap();
7991 let pending = producer.consume().request_broadcast("room/alice");
7992 let source = broadcast::Info::new().produce();
7993 queued(&server).await.accept(&source);
7994 let resolved = pending.await.unwrap();
7995
7996 source.close();
7997 settle(|| resolved.is_closed()).await;
7998 assert!(
7999 server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
8000 "the closed source was requested again"
8001 );
8002 }
8003
8004 #[tokio::test]
8009 async fn origin_front_drops_the_source_when_unused() {
8010 let producer = origin(1).produce();
8011 let consumer = producer.consume();
8012
8013 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
8014 let track = broadcast.create_track("video", None).unwrap();
8015 let mut group = track.append_group().unwrap();
8016 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
8017 group.finish().unwrap();
8018
8019 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
8020 let mut subscription = resolved
8021 .track("video")
8022 .unwrap()
8023 .subscribe(None)
8024 .await
8025 .expect("subscribe");
8026 let mut group = subscription.recv_group().await.unwrap().unwrap();
8027 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
8028 drop(group);
8029 drop(subscription);
8030
8031 tokio::time::timeout(Duration::from_secs(1), track.unused())
8032 .await
8033 .expect("source unused should resolve far below TRACK_IDLE_LINGER")
8034 .expect("source closed");
8035
8036 let mut again = resolved
8039 .track("video")
8040 .unwrap()
8041 .subscribe(track::Subscription::default().with_max_age(Duration::from_secs(3600)))
8042 .await
8043 .expect("resubscribe");
8044 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
8045 .await
8046 .expect("cached group is still on the front")
8047 .expect("recv group")
8048 .expect("track ended early");
8049 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
8050
8051 tokio::time::timeout(Duration::from_secs(1), track.used())
8052 .await
8053 .expect("returning reader re-splices the source")
8054 .expect("source closed");
8055
8056 let mut group = track.append_group().unwrap();
8057 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
8058 group.finish().unwrap();
8059 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
8060 .await
8061 .expect("groups past the cached edge come from the re-splice")
8062 .expect("recv group")
8063 .expect("track ended early");
8064 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
8065 }
8066
8067 #[tokio::test]
8077 async fn resumed_reader_skips_warm_groups_behind_the_new_edge() {
8078 let ms = |v: u64| crate::Timestamp::from_millis(v).unwrap();
8079 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
8080 let budget = track::Subscription::default().with_max_age(Duration::from_millis(100));
8081 let write = |source: &track::Producer, sequence: u64, millis: u64| {
8082 let mut group = source.create_group(sequence.into()).unwrap();
8083 group.write_frame(ms(millis), b"x".as_ref()).unwrap();
8084 group.finish().unwrap();
8085 };
8086 let drain = |subscription: &mut track::Subscriber| {
8087 let mut sequences = Vec::new();
8088 while let Poll::Ready(group) = subscription.poll_recv_group(&kio::Waiter::noop()) {
8089 sequences.push(group.unwrap().expect("track ended").sequence);
8090 }
8091 sequences
8092 };
8093
8094 let track = resolved.track("video").unwrap();
8095 let first = budget.clone();
8096 let subscribing = tokio::spawn(async move { track.subscribe(first).await });
8097 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
8098 .await
8099 .expect("the front asked the source")
8100 .expect("request");
8101 let source = request.resolving_start().accept(None);
8102 for (sequence, millis) in [(0, 0), (1, 20), (2, 40), (3, 60)] {
8103 write(&source, sequence, millis);
8104 }
8105 let mut subscription = subscribing.await.unwrap().expect("subscribe");
8106 next_group(&mut subscription).await.unwrap().expect("a cached group");
8107 drain(&mut subscription);
8108 drop(subscription);
8109
8110 tokio::time::timeout(Duration::from_secs(1), source.unused())
8111 .await
8112 .expect("parked")
8113 .expect("source open");
8114 drop(source);
8115
8116 let track = resolved.track("video").unwrap();
8117 let second = budget.clone();
8118 let subscribing = tokio::spawn(async move { track.subscribe(second).await });
8119 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
8120 .await
8121 .expect("the front asked the source again")
8122 .expect("request");
8123 let mut source = request.resolving_start().accept(None);
8124 for (sequence, millis) in [(20, 400), (21, 420), (22, 440)] {
8125 write(&source, sequence, millis);
8126 }
8127 source.start_at(3).unwrap();
8130 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
8131 settle(|| subscription.latest() == Some(22)).await;
8132
8133 assert_eq!(drain(&mut subscription), [3, 20, 21, 22]);
8136 }
8137
8138 #[tokio::test]
8143 async fn chained_front_drops_the_source_when_unused() {
8144 let leaf = origin(1).produce();
8145 let leaf_consumer = leaf.consume();
8146
8147 let broadcast = leaf.publish("room/alice", Route::default()).unwrap();
8148 let track = broadcast.create_track("video", None).unwrap();
8149 let mut group = track.append_group().unwrap();
8150 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
8151 group.finish().unwrap();
8152
8153 let leaf_front = leaf_consumer.request_broadcast("room/alice").await.expect("resolves");
8156
8157 let mid = origin(2).produce();
8158 let mid_server = mid.dynamic("room", Route::default().with_hops(hops(&[10]))).unwrap();
8159 let mid_pending = mid.consume().request_broadcast("room/alice");
8160 queued(&mid_server).await.accept(&leaf_front);
8161 let mid_resolved = mid_pending.await.expect("mid resolves");
8162
8163 let edge = origin(3).produce();
8164 let edge_server = edge.dynamic("room", Route::default().with_hops(hops(&[20]))).unwrap();
8165 let edge_pending = edge.consume().request_broadcast("room/alice");
8166 queued(&edge_server).await.accept(&mid_resolved);
8167 let edge_resolved = edge_pending.await.expect("edge resolves");
8168
8169 let mut subscription = edge_resolved
8170 .track("video")
8171 .unwrap()
8172 .subscribe(None)
8173 .await
8174 .expect("subscribe");
8175 let mut group = subscription.recv_group().await.unwrap().unwrap();
8176 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
8177 drop(group);
8178 drop(subscription);
8179
8180 tokio::time::timeout(Duration::from_secs(5), track.unused())
8181 .await
8182 .expect("chained unused should resolve far below TRACK_IDLE_LINGER")
8183 .expect("source closed");
8184
8185 let cached = edge_resolved.track("video").unwrap().cached_groups();
8186 assert_eq!(
8187 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
8188 vec![0],
8189 "every front keeps the delivered groups after releasing its source"
8190 );
8191
8192 let mut subscription = edge_resolved
8194 .track("video")
8195 .unwrap()
8196 .subscribe(track::Subscription::default().with_max_age(Duration::from_secs(3600)))
8197 .await
8198 .expect("resubscribe");
8199 tokio::time::timeout(Duration::from_secs(5), track.used())
8200 .await
8201 .expect("resubscribe should reach the leaf")
8202 .expect("source open");
8203 let mut group = track.append_group().unwrap();
8204 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
8205 group.finish().unwrap();
8206 let mut group = subscription.recv_group().await.unwrap().unwrap();
8207 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
8208 drop(group);
8209 let mut group = subscription.recv_group().await.unwrap().unwrap();
8210 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
8211 drop(group);
8212 drop(subscription);
8213
8214 tokio::time::timeout(Duration::from_secs(5), track.unused())
8215 .await
8216 .expect("second chained unused should resolve far below TRACK_IDLE_LINGER")
8217 .expect("source closed");
8218
8219 let cached = edge_resolved.track("video").unwrap().cached_groups();
8220 assert_eq!(
8221 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
8222 vec![0, 1],
8223 "repeated demand keeps every complete group while releasing its source"
8224 );
8225
8226 let fetch = edge_resolved.track("video").unwrap().fetch_group(2, None);
8227 let mut fetch = std::pin::pin!(fetch);
8228 assert!(futures::poll!(fetch.as_mut()).is_pending(), "fetch should re-splice");
8229 tokio::time::timeout(Duration::from_secs(5), track.used())
8230 .await
8231 .expect("fetch should reach the leaf")
8232 .expect("source open");
8233 let mut group = track.append_group().unwrap();
8234 group.write_frame(crate::Timestamp::ZERO, b"fetched".as_ref()).unwrap();
8235 group.finish().unwrap();
8236 let mut group = tokio::time::timeout(Duration::from_secs(5), fetch)
8237 .await
8238 .expect("re-spliced source should answer the fetch")
8239 .expect("fetch succeeds");
8240 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"fetched");
8241 }
8242
8243 #[tokio::test]
8247 async fn incompatible_local_source_keeps_the_incumbent() {
8248 let producer = origin(1).produce();
8249 let consumer = producer.consume();
8250
8251 let first = producer.publish("room/alice", Route::default()).unwrap();
8252 let track = first.create_track("video", None).unwrap();
8253 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
8254 let mut subscription = resolved
8255 .track("video")
8256 .unwrap()
8257 .subscribe(None)
8258 .await
8259 .expect("subscribe");
8260 let mut group = track.append_group().unwrap();
8261 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
8262 group.finish().unwrap();
8263 let mut group = subscription.recv_group().await.unwrap().unwrap();
8264 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"before");
8265
8266 let second = producer.publish("room/alice", Route::default()).unwrap();
8268 let _incompatible = second
8269 .create_track("video", track::Info::default().with_timescale(crate::Timescale::MICRO))
8270 .unwrap();
8271 for _ in 0..10 {
8272 tokio::task::yield_now().await;
8273 }
8274
8275 let mut group = track.append_group().unwrap();
8277 group.write_frame(crate::Timestamp::ZERO, b"still".as_ref()).unwrap();
8278 group.finish().unwrap();
8279 let mut group = subscription.recv_group().await.unwrap().unwrap();
8280 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"still");
8281
8282 drop(track);
8284 first.close();
8285 assert!(matches!(subscription.recv_group().await, Err(Error::Unsupported)));
8286 }
8287
8288 #[tokio::test]
8289 async fn multiple_scopes_present_one_broad_prefix() {
8290 let producer = origin(1).produce();
8291 let _a = producer.announce("", Route::default()).unwrap();
8292
8293 let consumer = producer.consume().scope("", &scopes(&["alpha", "beta"])).unwrap();
8294 let mut announced = consumer.announced();
8295 announced.assert_next_active("");
8296 announced.assert_next_wait();
8297 }
8298
8299 #[test]
8300 fn scope_accepts_every_pattern_union() {
8301 let producer = origin(1).produce();
8302
8303 let root = producer.scope("", &Patterns::from(Pattern::all())).unwrap();
8305 assert_eq!(root.allowed(), Patterns::from(Pattern::all()));
8306
8307 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
8309 assert_eq!(scoped.allowed(), scopes(&["room"]));
8310
8311 let multi = producer.scope("", &scopes(&["room", "room/chat", "anon"])).unwrap();
8313 assert_eq!(multi.allowed(), scopes(&["room", "anon"]));
8314
8315 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
8317 assert_eq!(consumer.allowed(), scopes(&["room"]));
8318
8319 for text in ["room", "", "*room", "room/*", "*", "**/room", "room/**/chat", "*.hang"] {
8320 let union = Patterns::from(text.parse::<Pattern>().unwrap());
8321 assert_eq!(producer.scope("", &union).expect(text).allowed(), union, "{text}");
8322 assert_eq!(
8323 producer.consume().scope("", &union).expect(text).allowed(),
8324 union,
8325 "{text}"
8326 );
8327 }
8328
8329 let mixed: Patterns = ["room/**".parse().unwrap(), "other".parse().unwrap()]
8330 .into_iter()
8331 .collect();
8332 assert_eq!(producer.scope("", &mixed).unwrap().allowed(), mixed);
8333 }
8334
8335 #[test]
8336 fn route_table_prunes_to_empty() {
8337 let producer = origin(1).produce();
8338 let consumer = producer.consume();
8339
8340 let cursor = consumer
8343 .scope("", &scopes(&["room/a", "other/deep/head"]))
8344 .unwrap()
8345 .announced();
8346 let route = producer.announce("room/a/b/c", Route::default()).unwrap();
8347 {
8348 let table = producer.shared.lock();
8349 assert!(table.routes.root.find(Path::new("room/a/b/c").parts()).is_some());
8350 assert!(table.routes.root.find(Path::new("other/deep/head").parts()).is_some());
8351 assert_eq!(table.routes.root.cursors_below, 2);
8352 }
8353
8354 drop(route);
8355 drop(cursor);
8356 let table = producer.shared.lock();
8357 assert!(table.routes.root.is_empty());
8358 assert_eq!(table.routes.root.cursors_below, 0);
8359 }
8360
8361 #[test]
8365 fn a_published_broadcast_keeps_the_driver_running() {
8366 let (producer, mut driver) = Producer::new(Config::new(origin(1)));
8367 let waiter = kio::Waiter::noop();
8368 let broadcast = producer.create_broadcast("room/a").unwrap();
8369 drop(producer);
8370 assert!(
8371 driver.poll(Instant::now(), &waiter).is_ok(),
8372 "the broadcast is lifecycle work"
8373 );
8374 drop(broadcast);
8375 assert!(matches!(driver.poll(Instant::now(), &waiter), Err(Error::Closed)));
8376 }
8377
8378 #[tokio::test]
8381 async fn a_dynamic_keeps_the_driver_running() {
8382 let (producer, driver) = Producer::new(Config::new(origin(1)));
8383 let run = tokio::spawn(crate::time::run(driver));
8384 let consumer = producer.consume();
8385 let server = producer.dynamic("room", Route::default()).unwrap();
8386 drop(producer);
8387
8388 let pending = consumer.request_broadcast("room/alice");
8389 let request = queued(&server).await;
8390 let source = broadcast::Info::new().produce();
8391 request.accept(&source);
8392 let resolved = pending.await.expect("the dynamic still serves");
8393
8394 drop(resolved);
8395 drop(source);
8396 drop(server);
8397 tokio::time::timeout(Duration::from_secs(5), run)
8398 .await
8399 .expect("driver must finish once the dynamic is gone")
8400 .unwrap();
8401 drop(consumer);
8402 }
8403
8404 #[test]
8405 fn watch_wakes_only_for_covering_changes() {
8406 let producer = origin(1).produce();
8407 let waiter = kio::Waiter::noop();
8408 let watch = producer.shared.lock().watch(&producer.shared, &Path::new("room/a"));
8409 let seen = watch.seen();
8410
8411 let _other = producer.announce("other", Route::default()).unwrap();
8413 let _below = producer.announce("room/a/b", Route::default()).unwrap();
8414 assert!(watch.poll_changed(&waiter, seen).is_pending());
8415
8416 let above = producer.announce("room", Route::default()).unwrap();
8418 assert!(watch.poll_changed(&waiter, seen).is_ready());
8419 let seen = watch.seen();
8420 drop(above);
8421 assert!(watch.poll_changed(&waiter, seen).is_ready());
8422 let seen = watch.seen();
8423
8424 let _beside = producer.create_broadcast("room/b").unwrap();
8426 assert!(watch.poll_changed(&waiter, seen).is_pending());
8427 let _here = producer.create_broadcast("room/a").unwrap();
8428 assert!(watch.poll_changed(&waiter, seen).is_ready());
8429
8430 drop(watch);
8432 let table = producer.shared.lock();
8433 let node = table
8434 .routes
8435 .root
8436 .find(Path::new("room/a").parts())
8437 .expect("route below keeps the node");
8438 assert!(node.watches.is_empty());
8439 assert_eq!(table.routes.root.watches_below, 0);
8440 }
8441
8442 #[test]
8443 fn a_discarded_front_task_unregisters_its_watch() {
8444 let (producer, _driver) = Producer::new(Config {
8445 hop: origin(1),
8446 ..Default::default()
8447 });
8448 let consumer = producer.consume();
8449 let _served = producer.dynamic("room", Route::default()).unwrap();
8450 drop(producer);
8455 let _pending = consumer.request_broadcast("room/a");
8456 }
8457
8458 #[test]
8459 fn create_broadcast_refuses_a_path_no_pattern_can_spell() {
8460 let producer = origin(1).produce();
8461
8462 assert!(matches!(
8466 producer.create_broadcast("room/*"),
8467 Err(Error::InvalidPath(_))
8468 ));
8469 assert!(matches!(
8470 producer.announce("room/**", Route::default()),
8471 Err(Error::InvalidPath(_))
8472 ));
8473 }
8474
8475 #[test]
8476 fn scope_empty_union_grants_nothing() {
8477 let producer = origin(1).produce();
8478
8479 assert!(matches!(producer.scope("", &Patterns::new()), Err(Error::Unauthorized)));
8481 assert!(matches!(
8482 producer.consume().scope("", &Patterns::new()),
8483 Err(Error::Unauthorized)
8484 ));
8485 }
8486
8487 #[test]
8488 fn scope_nests_and_rebases_roots() {
8489 let producer = origin(1).produce();
8490
8491 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
8493 let nested = scoped.scope("", &scopes(&["room/chat"])).unwrap();
8494 assert_eq!(nested.allowed(), scopes(&["room/chat"]));
8495
8496 assert!(matches!(
8498 scoped.scope("", &scopes(&["other"])),
8499 Err(Error::Unauthorized)
8500 ));
8501
8502 let rooted = nested.scope("room/chat", &Patterns::from(Pattern::all())).unwrap();
8504 assert_eq!(rooted.allowed(), scopes(&[""]));
8505
8506 let broadcast = nested.create_broadcast("room/chat/live").unwrap();
8508 assert!(producer.consume().get_broadcast("room/chat/live").is_some());
8509 broadcast.close();
8510 }
8511
8512 #[test]
8513 fn scope_intersects_and_rebases_arbitrary_grants() {
8514 let producer = origin(1).produce();
8515 let rooms = producer
8516 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
8517 .unwrap();
8518 let chats = rooms
8519 .scope("", &Patterns::from("*/chat".parse::<Pattern>().unwrap()))
8520 .unwrap();
8521 assert_eq!(chats.allowed(), Patterns::from("room/chat".parse::<Pattern>().unwrap()));
8522
8523 let exact = producer
8524 .scope("", &Patterns::from("room/alice".parse::<Pattern>().unwrap()))
8525 .unwrap();
8526 let rooted = exact.scope("room", &Patterns::from(Pattern::all())).unwrap();
8527 assert_eq!(rooted.allowed(), Patterns::from("alice".parse::<Pattern>().unwrap()));
8528 assert!(matches!(
8529 exact.scope("room/bob", &Patterns::from(Pattern::all())),
8530 Err(Error::Unauthorized)
8531 ));
8532
8533 let broadcast = exact.create_broadcast("room/alice").unwrap();
8534 assert!(matches!(
8535 exact.create_broadcast("room/alice/cam"),
8536 Err(Error::Unauthorized)
8537 ));
8538 assert!(producer.consume().get_broadcast("room/alice").is_some());
8539 drop(broadcast);
8540 }
8541
8542 #[tokio::test]
8543 async fn wildcard_scope_filters_announcements_and_reports_captures() {
8544 let producer = origin(1).produce();
8545 let consumer = producer
8546 .consume()
8547 .scope("", &Patterns::from("room/*/chat".parse::<Pattern>().unwrap()))
8548 .unwrap();
8549 let mut announced = consumer.announced();
8550
8551 let alice = producer.create_broadcast("room/alice/chat").unwrap();
8552 alice.announce(Route::default()).unwrap();
8553 let update = announced.try_next().expect("alice's chat");
8554 assert_eq!(update.prefix.as_str(), "room/alice/chat");
8555 assert_eq!(update.captures, Some(vec!["alice".parse::<Pattern>().unwrap()]));
8556
8557 let audio = producer.create_broadcast("room/alice/audio").unwrap();
8558 audio.announce(Route::default()).unwrap();
8559 announced.assert_next_wait();
8560
8561 let broad = producer.announce("room", Route::default()).unwrap();
8562 let update = announced.try_next().expect("overlapping broad route");
8563 assert_eq!(update.prefix.as_str(), "room");
8564 assert_eq!(update.captures, None, "an overlap does not pin the wildcard");
8565
8566 drop(broad);
8567 drop(audio);
8568 drop(alice);
8569 }
8570
8571 #[tokio::test]
8572 async fn local_broadcast_wins_announcement_ties() {
8573 let producer = origin(1).produce();
8574 let remote = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
8575 let local = producer.create_broadcast("room/alice").unwrap();
8576 local.announce(Route::default()).unwrap();
8577
8578 let mut announced = producer.consume().announced();
8579 let update = announced.try_next().expect("one winning route");
8580 assert_eq!(update.prefix.as_str(), "room/alice");
8581 assert_eq!(update.route.cost, Cost::default());
8582 announced.assert_next_wait();
8583
8584 drop(local);
8585 drop(remote);
8586 }
8587
8588 #[test]
8593 fn cost_charge_saturates() {
8594 assert_eq!(Cost { warm: 4, cold: 6 }.charged(5), Cost { warm: 9, cold: 11 });
8595 assert_eq!(Cost::new(u64::MAX).charged(10), Cost::new(MAX_COST));
8596
8597 assert_eq!(Cost::UNKNOWN.charged(3).cold, MAX_COST);
8600 }
8601
8602 fn expiring_origin(expiry: Duration) -> Producer {
8604 let pool = cache::Pool::new(cache::Config::default().with_expiry(expiry));
8605 Config {
8606 pool,
8607 ..Config::default()
8608 }
8609 .produce()
8610 }
8611
8612 #[tokio::test(start_paused = true)]
8616 async fn stalled_publisher_open_group_is_reclaimed() {
8617 let expiry = Duration::from_secs(1);
8618 let origin = expiring_origin(expiry);
8619 let broadcast = origin.create_broadcast("test").unwrap();
8620 let track = broadcast.create_track("video", None).unwrap();
8621
8622 let mut stalled = track.append_group().unwrap();
8623 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
8624 let _successor = track.append_group().unwrap();
8628
8629 let mut reading = stalled.consume();
8630 assert!(reading.read_frame().await.unwrap().is_some());
8631
8632 crate::model::clock::advance(expiry * 2);
8634
8635 let reclaimed = tokio::time::timeout(Duration::from_secs(60), reading.read_frame()).await;
8638 assert!(
8639 matches!(reclaimed, Ok(Err(Error::Old))),
8640 "the sweep must reclaim an idle open group and surface the gap, got {reclaimed:?}"
8641 );
8642 }
8643
8644 #[tokio::test(start_paused = true)]
8647 async fn sweep_respects_a_disabled_expiry() {
8648 let origin = Config {
8649 pool: cache::Pool::unbounded(),
8650 ..Config::default()
8651 }
8652 .produce();
8653 let broadcast = origin.create_broadcast("test").unwrap();
8654 let track = broadcast.create_track("video", None).unwrap();
8655
8656 let mut stalled = track.append_group().unwrap();
8657 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
8658 let _successor = track.append_group().unwrap();
8659
8660 let mut reading = stalled.consume();
8661 assert!(reading.read_frame().await.unwrap().is_some());
8662
8663 crate::model::clock::advance(Duration::from_secs(3600));
8664 tokio::time::advance(Duration::from_secs(3600)).await;
8665
8666 assert!(
8667 reading.read_frame().now_or_never().is_none(),
8668 "a pool without an expiry window never reclaims"
8669 );
8670 }
8671
8672 #[test]
8675 fn drain_cost_is_encodable() {
8676 use crate::coding::Encode;
8677
8678 let mut buf = Vec::new();
8679 Cost::DRAIN
8680 .encode(&mut buf, crate::lite::Version::Lite06)
8681 .expect("a draining route is still forwarded, so its cost must encode");
8682 }
8683}