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
131impl Default for Config {
132 fn default() -> Self {
134 let pool = cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY));
135 Self {
136 hop: Hop::random(),
137 pool,
138 cache_duration: Duration::MAX,
139 default_max_age: track::DEFAULT_MAX_AGE,
140 }
141 }
142}
143
144impl Config {
145 pub fn new(hop: Hop) -> Self {
147 Self { hop, ..Self::default() }
148 }
149}
150
151impl From<Hop> for Config {
152 fn from(hop: Hop) -> Self {
154 Self::new(hop)
155 }
156}
157
158impl TryFrom<u64> for Hop {
159 type Error = InvalidHop;
160
161 fn try_from(id: u64) -> Result<Self, Self::Error> {
162 Self::new(id)
163 }
164}
165
166impl fmt::Display for Hop {
167 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
168 self.id.fmt(f)
169 }
170}
171
172impl<V: Copy> Encode<V> for Hop
173where
174 u64: Encode<V>,
175{
176 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
177 self.id.encode(w, version)
178 }
179}
180
181impl<V: Copy> Decode<V> for Hop
182where
183 u64: Decode<V>,
184{
185 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
186 Self::from_wire(u64::decode(r, version)?)
187 }
188}
189
190pub(crate) const MAX_HOPS: usize = 32;
196
197#[derive(Debug, Clone, Default, PartialEq, Eq)]
205pub struct Hops(Vec<Hop>);
206
207#[derive(Debug, Clone, Copy, PartialEq, Eq)]
209#[non_exhaustive]
210pub enum InvalidHop {
211 Range,
214
215 TooMany,
218
219 Duplicate,
223}
224
225impl fmt::Display for InvalidHop {
226 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
227 match self {
228 Self::Range => write!(f, "local hop id must be non-zero and below 2^62"),
229 Self::TooMany => write!(f, "too many hops (max {MAX_HOPS})"),
230 Self::Duplicate => write!(f, "hop already in the chain"),
231 }
232 }
233}
234
235impl std::error::Error for InvalidHop {}
236
237impl From<InvalidHop> for DecodeError {
238 fn from(err: InvalidHop) -> Self {
239 match err {
240 InvalidHop::TooMany => DecodeError::BoundsExceeded,
241 InvalidHop::Range | InvalidHop::Duplicate => DecodeError::InvalidValue,
242 }
243 }
244}
245
246impl Hops {
247 pub fn new() -> Self {
249 Self(Vec::new())
250 }
251
252 pub fn push(&mut self, hop: Hop) -> Result<(), InvalidHop> {
258 if self.0.len() >= MAX_HOPS {
259 return Err(InvalidHop::TooMany);
260 }
261 if hop != Hop::UNKNOWN && self.0.contains(&hop) {
262 return Err(InvalidHop::Duplicate);
263 }
264 self.0.push(hop);
265 Ok(())
266 }
267
268 pub(crate) fn stamp(&mut self, stamp: Hop) -> Result<(), InvalidHop> {
277 match self.0.first() {
278 None => {
279 self.push(stamp)?;
280 self.push(Hop::UNKNOWN)
281 }
282 Some(first) if *first == Hop::UNKNOWN => {
283 if self.0.len() >= MAX_HOPS {
284 return Err(InvalidHop::TooMany);
285 }
286 if self.0.contains(&stamp) {
287 return Err(InvalidHop::Duplicate);
288 }
289 self.0.insert(0, stamp);
290 Ok(())
291 }
292 Some(_) => Ok(()),
293 }
294 }
295
296 pub fn contains(&self, hop: &Hop) -> bool {
298 self.0.contains(hop)
299 }
300
301 pub fn len(&self) -> usize {
303 self.0.len()
304 }
305
306 pub fn is_empty(&self) -> bool {
308 self.0.is_empty()
309 }
310
311 pub fn iter(&self) -> std::slice::Iter<'_, Hop> {
313 self.0.iter()
314 }
315
316 pub fn as_slice(&self) -> &[Hop] {
318 &self.0
319 }
320}
321
322impl TryFrom<Vec<Hop>> for Hops {
323 type Error = InvalidHop;
324
325 fn try_from(v: Vec<Hop>) -> Result<Self, Self::Error> {
326 if v.len() > MAX_HOPS {
327 return Err(InvalidHop::TooMany);
328 }
329 for (i, hop) in v.iter().enumerate() {
331 if *hop != Hop::UNKNOWN && v[i + 1..].contains(hop) {
332 return Err(InvalidHop::Duplicate);
333 }
334 }
335 Ok(Self(v))
336 }
337}
338
339impl<'a> IntoIterator for &'a Hops {
340 type Item = &'a Hop;
341 type IntoIter = std::slice::Iter<'a, Hop>;
342
343 fn into_iter(self) -> Self::IntoIter {
344 self.iter()
345 }
346}
347
348impl<V: Copy> Encode<V> for Hops
349where
350 u64: Encode<V>,
351 Hop: Encode<V>,
352{
353 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
354 (self.0.len() as u64).encode(w, version)?;
355 for origin in &self.0 {
356 origin.encode(w, version)?;
357 }
358 Ok(())
359 }
360}
361
362impl<V: Copy> Decode<V> for Hops
363where
364 u64: Decode<V>,
365 Hop: Decode<V>,
366{
367 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
368 let count = u64::decode(r, version)? as usize;
369 if count > MAX_HOPS {
370 return Err(DecodeError::BoundsExceeded);
371 }
372 let mut list = Self(Vec::with_capacity(count));
375 for _ in 0..count {
376 list.push(Hop::decode(r, version)?)?;
377 }
378 Ok(list)
379 }
380}
381
382const MAX_COST: u64 = (1 << 62) - 1;
390
391#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)]
400pub struct Cost {
401 pub warm: u64,
410
411 pub cold: u64,
418}
419
420impl Cost {
421 pub const fn new(cost: u64) -> Self {
424 Self { warm: cost, cold: cost }
425 }
426
427 pub const MAX: Self = Self::new(MAX_COST);
434
435 pub const DRAIN: Self = Self::MAX;
438
439 pub(crate) const UNKNOWN: Self = Self {
443 warm: 0,
444 cold: MAX_COST,
445 };
446
447 pub(crate) fn charged(self, link_cost: u64) -> Self {
450 Self {
451 warm: self.warm.saturating_add(link_cost).min(MAX_COST),
452 cold: self.cold.saturating_add(link_cost).min(MAX_COST),
453 }
454 }
455
456 pub(crate) fn clamped(self) -> Self {
459 Self {
460 warm: self.warm.min(MAX_COST),
461 cold: self.cold.min(MAX_COST),
462 }
463 }
464}
465
466impl From<u64> for Cost {
467 fn from(cost: u64) -> Self {
468 Self::new(cost)
469 }
470}
471
472#[derive(Clone, Debug, PartialEq, Eq)]
483#[non_exhaustive]
484pub struct Route {
485 pub hops: Hops,
491
492 pub cost: Cost,
497
498 pub(crate) via: Hop,
504
505 pub(crate) source: Source,
507}
508
509impl Default for Route {
510 fn default() -> Self {
511 Self {
512 hops: Hops::new(),
513 cost: Cost::default(),
514 via: Hop::UNKNOWN,
515 source: Source::Local,
516 }
517 }
518}
519
520#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash)]
526pub enum Source {
527 #[default]
529 Local,
530 Peer(Hop),
533}
534
535impl Route {
536 pub fn with_hops(mut self, hops: Hops) -> Self {
538 self.hops = hops;
539 self
540 }
541
542 pub fn with_cost(mut self, cost: impl Into<Cost>) -> Self {
547 self.cost = cost.into();
548 self
549 }
550
551 pub(crate) fn with_via(mut self, via: Hop) -> Self {
556 self.via = via;
557 self
558 }
559
560 pub fn is_anonymous(&self) -> bool {
568 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
569 }
570
571 pub fn source(&self) -> Source {
576 self.source
577 }
578}
579
580static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
581
582#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
583struct ConsumerId(u64);
584
585impl ConsumerId {
586 fn new() -> Self {
587 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
588 }
589}
590
591fn fnv_key(name: &str, origins: impl IntoIterator<Item = Hop>) -> u64 {
600 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
602
603 let mut hash = SEED;
604 for &byte in name.as_bytes() {
605 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
606 }
607 for origin in origins {
608 for &byte in &origin.id().to_le_bytes() {
609 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
610 }
611 }
612
613 hash
614}
615
616fn route_order(prefix: &Path, entry: &RouteEntry) -> (bool, Cost, bool, usize, u64, Reverse<u64>) {
624 (
625 entry.is_anonymous(),
626 entry.cost,
627 !entry.local,
628 entry.hops.len(),
629 fnv_key(prefix.as_str(), entry.hops.iter().copied()),
630 Reverse(entry.id),
631 )
632}
633
634type RouteMeta = (Hops, Cost, Source);
636
637type AnnounceMeta = (RouteMeta, Option<Vec<Pattern>>);
644
645enum PendingUpdate {
646 Announce(AnnounceMeta),
647 Unannounce(AnnounceMeta),
648 UnannounceAnnounce { old: AnnounceMeta, new: AnnounceMeta },
649}
650
651#[derive(Default)]
656struct OriginConsumerState {
657 pending: BTreeMap<PathOwned, PendingUpdate>,
658 delivered: BTreeSet<PathOwned>,
663 ended: bool,
666}
667
668impl OriginConsumerState {
669 fn apply_announce(&mut self, prefix: PathOwned, meta: RouteMeta, captures: Option<Vec<Pattern>>) {
670 let meta = (meta, captures);
671 let new = match self.pending.remove(&prefix) {
672 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(meta),
674 Some(PendingUpdate::Unannounce(old) | PendingUpdate::UnannounceAnnounce { old, .. }) => {
676 PendingUpdate::UnannounceAnnounce { old, new: meta }
677 }
678 };
679 self.pending.insert(prefix, new);
680 }
681
682 fn apply_unannounce(&mut self, prefix: PathOwned, last: RouteMeta, captures: Option<Vec<Pattern>>) {
683 let last = (last, captures);
684 match self.pending.remove(&prefix) {
685 Some(PendingUpdate::Announce(_)) if !self.delivered.contains(&prefix) => {}
688 None | Some(PendingUpdate::Announce(_) | PendingUpdate::Unannounce(_)) => {
691 self.pending.insert(prefix, PendingUpdate::Unannounce(last));
692 }
693 Some(PendingUpdate::UnannounceAnnounce { old, .. }) => {
696 self.pending.insert(prefix, PendingUpdate::Unannounce(old));
697 }
698 }
699 }
700
701 fn take(&mut self) -> Option<AnnounceUpdate> {
703 let prefix = self.pending.keys().next()?.clone();
704 let ((meta, captures), kind) = match self.pending.remove(&prefix).unwrap() {
705 PendingUpdate::Announce(meta) => {
706 let kind = match self.delivered.insert(prefix.clone()) {
708 true => AnnounceKind::Announced,
709 false => AnnounceKind::Updated,
710 };
711 (meta, kind)
712 }
713 PendingUpdate::Unannounce(meta) => {
714 self.delivered.remove(&prefix);
715 (meta, AnnounceKind::Retracted)
716 }
717 PendingUpdate::UnannounceAnnounce { old, new } => {
718 self.delivered.remove(&prefix);
721 self.pending.insert(prefix.clone(), PendingUpdate::Announce(new));
722 (old, AnnounceKind::Retracted)
723 }
724 };
725 Some(AnnounceUpdate {
726 prefix,
727 captures,
728 route: Route {
729 hops: meta.0,
730 cost: meta.1,
731 via: Hop::UNKNOWN,
732 source: meta.2,
733 },
734 kind,
735 })
736 }
737}
738
739struct RouteEntry {
741 id: u64,
742 prefix: PathOwned,
743 scope: Patterns,
746 hops: Hops,
747 cost: Cost,
748 via: Hop,
752 local: bool,
755 peer: bool,
758 server: Option<kio::Shared<ServeState>>,
762 source: Option<broadcast::Consumer>,
766 advertised: bool,
770 stale: bool,
774 claim: Pattern,
780}
781
782impl RouteEntry {
783 fn live(&self) -> bool {
785 self.advertised && !self.stale
786 }
787
788 fn is_anonymous(&self) -> bool {
789 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
790 }
791
792 fn entered(&self) -> Source {
794 match self.peer {
795 true => Source::Peer(self.via),
796 false => Source::Local,
797 }
798 }
799
800 fn serves(&self, path: &Path) -> bool {
804 self.server.is_some() || (self.source.is_some() && self.prefix == *path)
805 }
806
807 fn qualifies(&self, pin: Pin) -> bool {
809 match pin {
810 Pin::Any => true,
811 Pin::Local => self.local,
812 Pin::Publisher(first) => self.hops.iter().next() == Some(&first),
813 Pin::Route(id) => self.id == id && self.hops.iter().next().is_none_or(|first| *first == Hop::UNKNOWN),
815 }
816 }
817
818 fn visible_to(&self, exclude: Option<Hop>) -> bool {
823 match exclude {
824 Some(peer) if peer != Hop::UNKNOWN => self.via != peer && !self.hops.contains(&peer),
825 _ => true,
826 }
827 }
828
829 fn overlaps(&self, allowed: &Patterns) -> bool {
831 self.scope.iter().any(|scope| {
832 scope
833 .intersect(&self.claim)
834 .is_ok_and(|scoped| scoped.iter().any(|restriction| allowed.overlaps(restriction)))
835 })
836 }
837}
838
839fn prefix_claim(prefix: &Path) -> Result<Pattern, InvalidPattern> {
841 if prefix.parts().count() == Path::MAX_PARTS {
842 Pattern::literal(prefix.as_str())
843 } else {
844 Pattern::subtree(prefix.as_str())
845 }
846}
847
848#[derive(Default)]
853struct ServeState {
854 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
857
858 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
864
865 closed: bool,
868}
869
870type FrontKey = (PathOwned, Horizon);
875
876#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash)]
878struct Horizon {
879 exclude: Option<Hop>,
883 local: bool,
885}
886
887impl Horizon {
888 fn admits(&self, entry: &RouteEntry) -> bool {
890 !(self.local && entry.peer) && entry.visible_to(self.exclude)
891 }
892}
893
894#[derive(Clone)]
897struct RemoteFront {
898 request: kio::Producer<PendingBroadcast>,
902 broadcast: broadcast::WeakConsumer,
905 pin: kio::Lock<Pin>,
908}
909
910type CursorRoute = (u64, RouteMeta, bool, Option<Vec<Pattern>>);
912
913impl WeakEntry for RemoteFront {
914 fn is_closed(&self) -> bool {
915 self.broadcast.is_closed()
916 }
917
918 fn same_channel(&self, other: &Self) -> bool {
919 self.broadcast.same_channel(&other.broadcast)
920 }
921}
922
923struct TableCursor {
926 root: PathOwned,
928 allowed: Patterns,
930 heads: Vec<PathOwned>,
934 horizon: Horizon,
936 hidden: Hidden,
938 state: kio::Producer<OriginConsumerState>,
941 under: PathOwned,
944 holes: Vec<PathOwned>,
947 named: Option<(Mount, Patterns)>,
950 current: HashMap<PathOwned, CursorRoute>,
955}
956
957impl TableCursor {
958 fn presented(&self, prefix: &Path, claim: &Pattern) -> Option<PathOwned> {
964 if !self.allowed.overlaps(claim) {
965 return None;
966 }
967
968 if let Some(relative) = prefix.strip_prefix(&self.root) {
969 return Some(relative.to_owned());
970 }
971 self.root.has_prefix(prefix).then(PathOwned::default)
972 }
973
974 fn captures(&self, prefix: &Path) -> Option<Vec<Pattern>> {
977 let (literal, allowed) = match &self.named {
978 Some((mount, allowed)) => (Pattern::literal(mount.name(prefix)?.as_str()).ok()?, allowed),
979 None => (Pattern::literal(prefix.as_str()).ok()?, &self.allowed),
980 };
981 allowed
982 .iter()
983 .filter_map(|allowed| {
984 allowed
985 .captures(&literal)
986 .map(|captures| (allowed.specificity(), captures))
987 })
988 .max_by_key(|(specificity, _)| *specificity)
989 .map(|(_, captures)| captures)
990 }
991
992 fn visible(&self, entry: &RouteEntry) -> bool {
996 entry.live()
997 && self.horizon.admits(entry)
998 && entry.overlaps(&self.allowed)
999 && self.discovers(&entry.prefix)
1000 && !self.holes.iter().any(|hole| entry.prefix.has_prefix(hole))
1001 && self.named.as_ref().is_none_or(|(mount, _)| mount.names(&entry.prefix))
1002 }
1003
1004 fn discovers(&self, prefix: &Path) -> bool {
1006 self.hidden.discovers(&self.heads, prefix)
1007 }
1008}
1009
1010#[derive(Clone)]
1013struct OriginScope {
1014 allowed: Patterns,
1017 mounts: Arc<[Mount]>,
1019}
1020
1021#[derive(Clone, Debug)]
1024struct Mount {
1025 at: PathOwned,
1026 target: PathOwned,
1027}
1028
1029impl Mount {
1030 fn resolve(&self, path: &Path) -> Option<PathOwned> {
1033 let resolved = self.target.join(path.strip_prefix(&self.at)?);
1034 (resolved.parts().count() <= Path::MAX_PARTS).then_some(resolved)
1035 }
1036
1037 fn name(&self, path: &Path) -> Option<PathOwned> {
1040 Some(self.at.join(path.strip_prefix(&self.target)?))
1041 }
1042
1043 fn names(&self, path: &Path) -> bool {
1047 path.strip_prefix(&self.target)
1048 .is_none_or(|rest| self.at.parts().count() + rest.parts().count() <= Path::MAX_PARTS)
1049 }
1050
1051 fn translate(&self, patterns: &Patterns) -> Patterns {
1058 let target = self.target.as_str();
1059 patterns
1060 .rebase(self.at.as_str())
1061 .iter()
1062 .filter_map(|member| {
1063 member
1064 .rooted(target)
1065 .or_else(|_| {
1066 let segments = member
1067 .segments()
1068 .iter()
1069 .filter(|segment| **segment != Segment::Globstar);
1070 Pattern::new(segments.cloned())?.rooted(target)
1071 })
1072 .ok()
1073 })
1074 .collect()
1075 }
1076
1077 fn translate_head(&self, head: &Path) -> Option<PathOwned> {
1080 match head.strip_prefix(&self.at) {
1081 Some(rest) => Some(self.target.join(rest)),
1082 None => self.at.has_prefix(head).then(|| self.target.clone()),
1083 }
1084 }
1085}
1086
1087impl OriginScope {
1088 fn empty() -> Self {
1090 Self {
1091 allowed: Patterns::new(),
1092 mounts: Arc::from([]),
1093 }
1094 }
1095
1096 fn narrow(&self, patterns: &Patterns) -> Option<Self> {
1098 let allowed = self.allowed.intersect(patterns).ok()?;
1099 if allowed.is_empty() {
1100 None
1101 } else {
1102 Some(Self {
1103 allowed,
1104 mounts: self.mounts.clone(),
1105 })
1106 }
1107 }
1108
1109 fn mount(&self, path: &Path) -> Option<&Mount> {
1111 self.mounts.iter().find(|mount| path.has_prefix(&mount.at))
1112 }
1113
1114 fn resolve<'a>(&self, path: &'a Path<'a>) -> Option<Path<'a>> {
1117 match self.mount(path) {
1118 Some(mount) => mount.resolve(path),
1119 None => Some(path.borrow()),
1120 }
1121 }
1122
1123 fn publishes(&self, prefix: &Path) -> bool {
1126 self.mount(prefix).is_none()
1127 }
1128
1129 fn permits(&self, path: &Path) -> bool {
1131 self.allowed.matches(path.as_str())
1132 }
1133
1134 fn relative(&self, root: &Path) -> Patterns {
1136 self.allowed.rebase(root.as_str())
1137 }
1138}
1139
1140impl Default for OriginScope {
1141 fn default() -> Self {
1142 Self {
1143 allowed: Patterns::from(Pattern::all()),
1144 mounts: Arc::from([]),
1145 }
1146 }
1147}
1148
1149pub(crate) fn interest_prefixes(allowed: &Patterns) -> Vec<PathOwned> {
1152 let mut heads: Vec<PathOwned> = allowed
1153 .iter()
1154 .map(|pattern| Path::new(pattern.head()).to_owned())
1155 .collect();
1156 heads.sort();
1157 heads.dedup();
1158 let covered = heads.clone();
1159 heads.retain(|head| !covered.iter().any(|other| other != head && head.has_prefix(other)));
1160 heads
1161}
1162
1163#[derive(Clone, Debug, Default, PartialEq, Eq)]
1170pub(crate) struct Hidden {
1171 include: bool,
1173 from: Option<PathOwned>,
1175 beyond: Option<Vec<PathOwned>>,
1178}
1179
1180impl Hidden {
1181 fn discovers(&self, heads: &[PathOwned], prefix: &Path) -> bool {
1183 (self.include || !hides(self.from.as_ref().map(std::slice::from_ref).unwrap_or(heads), prefix))
1184 && self.beyond.as_ref().is_none_or(|outer| hides(outer, prefix))
1185 }
1186
1187 fn translate(&self, mount: &Mount) -> Self {
1191 let beyond = self.beyond.as_ref().and_then(|outer| {
1192 (!hides(outer, &mount.at)).then(|| outer.iter().filter_map(|head| mount.translate_head(head)).collect())
1193 });
1194 Self {
1195 include: self.include,
1196 from: self.from.as_ref().and_then(|head| mount.translate_head(head)),
1197 beyond,
1198 }
1199 }
1200}
1201
1202fn hides(heads: &[PathOwned], prefix: &Path) -> bool {
1205 heads
1206 .iter()
1207 .any(|head| prefix.strip_prefix(head).is_some_and(|below| below.is_hidden()))
1208}
1209
1210#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1212pub enum AnnounceKind {
1213 Announced,
1215 Updated,
1217 Retracted,
1219}
1220
1221impl AnnounceKind {
1222 pub fn is_active(self) -> bool {
1224 !matches!(self, Self::Retracted)
1225 }
1226}
1227
1228#[derive(Clone, Debug)]
1236pub struct AnnounceUpdate {
1237 pub prefix: PathOwned,
1239 pub captures: Option<Vec<Pattern>>,
1242 pub route: Route,
1245 pub kind: AnnounceKind,
1247}
1248
1249#[derive(Clone)]
1251pub struct Producer {
1252 hop: Hop,
1255
1256 scope: OriginScope,
1258
1259 root: PathOwned,
1261
1262 shared: kio::Shared<OriginState>,
1265
1266 pool: cache::Pool,
1269
1270 cache_duration: Duration,
1273
1274 default_max_age: Duration,
1277
1278 stats: stats::Session,
1282
1283 peer: bool,
1286
1287 tasks: Tasks,
1291
1292 timers: Clock,
1294}
1295
1296impl Producer {
1297 pub fn new(config: Config) -> (Self, Driver) {
1304 let (tasks, set) = TaskSet::new();
1305 let scope = OriginScope::default();
1306 let shared = kio::Shared::<OriginState>::default();
1307 let timers = Clock::default();
1308 let pool = config.pool.clone();
1309 let producer = Self {
1310 hop: config.hop,
1311 scope: scope.clone(),
1312 root: PathOwned::default(),
1313 shared: shared.clone(),
1314 pool: config.pool,
1315 cache_duration: config.cache_duration,
1316 default_max_age: config.default_max_age,
1317 stats: stats::Session::default(),
1318 peer: false,
1319 tasks,
1320 timers: timers.clone(),
1321 };
1322 let driver = Driver {
1323 state: DriverState {
1324 set,
1325 shared,
1326 done: false,
1327 },
1328 timers,
1329 pool,
1330 };
1331 (producer, driver)
1332 }
1333
1334 pub fn with_stats(mut self, session: stats::Session) -> Self {
1338 self.stats = session;
1339 self
1340 }
1341
1342 pub fn peer(mut self) -> Self {
1350 self.peer = true;
1351 self
1352 }
1353
1354 pub fn config(&self) -> Config {
1356 Config {
1357 hop: self.hop,
1358 pool: self.pool.clone(),
1359 cache_duration: self.cache_duration,
1360 default_max_age: self.default_max_age,
1361 }
1362 }
1363
1364 pub fn hop(&self) -> Hop {
1366 self.hop
1367 }
1368
1369 pub(crate) fn default_max_age(&self) -> Duration {
1372 self.default_max_age
1373 }
1374
1375 pub(crate) fn empty(hop: Hop) -> Self {
1380 let (tasks, _) = TaskSet::new();
1383 Self {
1384 hop,
1385 scope: OriginScope::empty(),
1386 root: PathOwned::default(),
1387 shared: kio::Shared::default(),
1388 pool: cache::Pool::default(),
1389 cache_duration: Duration::MAX,
1390 default_max_age: track::DEFAULT_MAX_AGE,
1391 stats: stats::Session::default(),
1392 peer: false,
1393 tasks,
1394 timers: Clock::default(),
1395 }
1396 }
1397
1398 pub fn create_broadcast(&self, path: impl AsPath) -> Result<broadcast::Producer, Error> {
1429 let path = path.as_path();
1430
1431 let full = self.root.join(&path).to_owned();
1432 if !self.scope.permits(&full) || !self.scope.publishes(&full) {
1433 return Err(Error::Unauthorized);
1434 }
1435 if full.parts().count() > Path::MAX_PARTS {
1439 return Err(BoundsExceeded.into());
1440 }
1441 let claim = prefix_claim(&full)?;
1444
1445 let ingress = self.stats.ingress(&full);
1447
1448 let announcing = Announcing {
1453 hop: self.hop,
1454 shared: self.shared.clone(),
1455 requested: full.clone(),
1456 prefixes: vec![(full.clone(), claim)],
1457 scope: self.scope.allowed.clone(),
1458 local: true,
1459 peer: self.peer,
1460 stats: self.stats.clone(),
1461 };
1462 let info = broadcast::Info {
1463 pool: self.pool.clone(),
1464 cache_duration: self.cache_duration,
1465 path: full,
1466 };
1467 let source = info.produce().with_stats(ingress.clone());
1468 let entry = announcing.announce(
1469 Route::default(),
1470 Serving {
1471 server: None,
1472 source: Some(source.consume()),
1473 advertised: false,
1474 },
1475 )?;
1476 Ok(source.with_announcer(Announcer {
1477 entry,
1478 ingress,
1479 _keepalive: self.tasks.keepalive(),
1480 }))
1481 }
1482
1483 pub fn publish(&self, path: impl AsPath, route: Route) -> Result<broadcast::Producer, Error> {
1485 let broadcast = self.create_broadcast(path)?;
1486 broadcast.announce(route)?;
1487 Ok(broadcast)
1488 }
1489
1490 pub(crate) fn create_source(&self, path: impl AsPath) -> broadcast::Producer {
1496 let path = path.as_path();
1497 let full = self.root.join(&path).to_owned();
1498 let ingress = self.stats.ingress(&full);
1499 broadcast::Info {
1500 pool: self.pool.clone(),
1501 cache_duration: self.cache_duration,
1502 path: full,
1503 }
1504 .produce()
1505 .with_stats(ingress)
1506 }
1507
1508 #[cfg(test)]
1517 pub(crate) fn announce(&self, prefix: impl AsPath, route: Route) -> Result<AnnounceProducer, Error> {
1518 Announcing::new(self, prefix)?.announce(
1519 route,
1520 Serving {
1521 server: None,
1522 source: None,
1523 advertised: true,
1524 },
1525 )
1526 }
1527
1528 pub fn dynamic(&self, prefix: impl AsPath, route: Route) -> Result<Dynamic, Error> {
1550 let announcing = Announcing::new(self, prefix)?;
1551 let serve = kio::Shared::<ServeState>::default();
1552 serve.lock().requests.add_handler();
1553 let announcement = announcing.announce(
1554 route,
1555 Serving {
1556 server: Some(serve.clone()),
1557 source: None,
1558 advertised: true,
1559 },
1560 )?;
1561 Ok(Dynamic {
1562 announcement,
1563 state: serve,
1564 _keepalive: self.tasks.keepalive(),
1565 })
1566 }
1567
1568 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Producer, Error> {
1575 let root = self.root.join(root).to_owned();
1576 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
1577 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
1578 Ok(Producer {
1579 hop: self.hop,
1580 scope,
1581 root,
1582 shared: self.shared.clone(),
1583 pool: self.pool.clone(),
1584 cache_duration: self.cache_duration,
1585 default_max_age: self.default_max_age,
1586 stats: self.stats.clone(),
1587 peer: self.peer,
1588 tasks: self.tasks.clone(),
1589 timers: self.timers.clone(),
1590 })
1591 }
1592
1593 pub fn mount(&self, at: impl AsPath, target: impl AsPath) -> Result<Producer, Error> {
1612 let at = self.root.join(at).to_owned();
1613 let target = self.root.join(target).to_owned();
1614 if [&at, &target]
1615 .into_iter()
1616 .any(|path| path.parts().count() > Path::MAX_PARTS)
1617 {
1618 return Err(BoundsExceeded.into());
1619 }
1620 Pattern::literal(at.as_str())?;
1622 if !self
1623 .scope
1624 .allowed
1625 .covers(&Patterns::from(Pattern::subtree(target.as_str())?))
1626 {
1627 return Err(Error::Unauthorized);
1628 }
1629 let overlaps = |a: &Path, b: &Path| a.has_prefix(b) || b.has_prefix(a);
1633 if overlaps(&at, &target)
1634 || self
1635 .scope
1636 .mounts
1637 .iter()
1638 .any(|mount| overlaps(&at, &mount.at) || overlaps(&target, &mount.at) || overlaps(&at, &mount.target))
1639 {
1640 return Err(Error::Duplicate);
1641 }
1642 let mounts = self
1643 .scope
1644 .mounts
1645 .iter()
1646 .cloned()
1647 .chain([Mount { at, target }])
1648 .collect();
1649 Ok(Producer {
1650 scope: OriginScope {
1651 allowed: self.scope.allowed.clone(),
1652 mounts,
1653 },
1654 ..self.clone()
1655 })
1656 }
1657
1658 pub fn consume(&self) -> Consumer {
1663 Consumer::from_producer(self, stats::Session::default())
1666 }
1667
1668 pub fn root(&self) -> &Path<'_> {
1670 &self.root
1671 }
1672
1673 pub fn allowed(&self) -> Patterns {
1675 self.scope.relative(&self.root)
1676 }
1677
1678 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1680 self.root.join(path)
1681 }
1682}
1683
1684struct Announcing {
1688 hop: Hop,
1689 shared: kio::Shared<OriginState>,
1690 requested: PathOwned,
1692 prefixes: Vec<(PathOwned, Pattern)>,
1696 scope: Patterns,
1698 local: bool,
1699 peer: bool,
1701 stats: stats::Session,
1702}
1703
1704impl Announcing {
1705 fn new(producer: &Producer, prefix: impl AsPath) -> Result<Self, Error> {
1707 let requested = producer.root.join(prefix.as_path()).to_owned();
1708 if requested.parts().count() > Path::MAX_PARTS {
1709 return Err(BoundsExceeded.into());
1710 }
1711 let claim = prefix_claim(&requested)?;
1712 if !producer.scope.allowed.overlaps(&claim) || !producer.scope.publishes(&requested) {
1713 return Err(Error::Unauthorized);
1714 }
1715 Ok(Self {
1716 hop: producer.hop,
1717 shared: producer.shared.clone(),
1718 requested: requested.clone(),
1719 prefixes: vec![(requested, claim)],
1720 scope: producer.scope.allowed.clone(),
1721 local: false,
1722 peer: producer.peer,
1723 stats: producer.stats.clone(),
1724 })
1725 }
1726
1727 fn announce(&self, route: Route, serving: Serving) -> Result<AnnounceProducer, Error> {
1728 debug_assert!(
1729 !route.hops.contains(&self.hop),
1730 "announce called with a looping hop chain",
1731 );
1732
1733 let via = route.via;
1734
1735 let mut shared = self.shared.lock();
1736 if shared.closed {
1737 return Err(Error::Closed);
1738 }
1739
1740 let mut entries = Vec::with_capacity(self.prefixes.len());
1741 for (prefix, claim) in &self.prefixes {
1742 shared.reannounced(prefix, &route.hops);
1743 let id = shared.next_route;
1744 shared.next_route += 1;
1745 let stale = shared.withdrawn_through(prefix, &route.hops);
1746 shared.routes.insert(RouteEntry {
1747 id,
1748 prefix: prefix.clone(),
1749 scope: self.scope.clone(),
1750 hops: route.hops.clone(),
1751 cost: route.cost,
1752 via,
1753 local: self.local,
1754 peer: self.peer,
1755 server: serving.server.clone(),
1756 source: serving.source.clone(),
1757 advertised: serving.advertised,
1758 stale,
1759 claim: claim.clone(),
1760 });
1761 shared.sync_route(prefix, claim);
1762 entries.push((prefix.clone(), id));
1763 }
1764 drop(shared);
1765
1766 let guard = serving
1768 .advertised
1769 .then(|| self.stats.ingress(&self.requested).announce());
1770
1771 Ok(AnnounceProducer {
1772 shared: self.shared.clone(),
1773 entries,
1774 guard,
1775 })
1776 }
1777}
1778
1779struct Serving {
1781 server: Option<kio::Shared<ServeState>>,
1782 source: Option<broadcast::Consumer>,
1783 advertised: bool,
1784}
1785
1786pub(crate) struct Announcer {
1793 entry: AnnounceProducer,
1794 ingress: stats::Scope,
1796 _keepalive: Keepalive,
1800}
1801
1802impl Announcer {
1803 pub(crate) fn announce(&mut self, route: Route) -> Result<(), Error> {
1805 self.entry.update(route)?;
1806 if self.entry.guard.is_none() {
1807 self.entry.guard = Some(self.ingress.announce());
1808 }
1809 Ok(())
1810 }
1811
1812 pub(crate) fn withdraw(&mut self) {
1814 self.entry.withdraw();
1815 self.entry.guard = None;
1816 }
1817}
1818
1819#[must_use = "dropping an announcement retracts the route"]
1826pub(crate) struct AnnounceProducer {
1827 shared: kio::Shared<OriginState>,
1828 entries: Vec<(PathOwned, u64)>,
1831 guard: Option<stats::Announce>,
1833}
1834
1835impl AnnounceProducer {
1836 pub fn update(&self, route: Route) -> Result<(), Error> {
1844 let mut shared = self.shared.lock();
1845 if shared.closed {
1846 return Err(Error::Closed);
1847 }
1848 for (prefix, id) in &self.entries {
1849 shared.reannounced(prefix, &route.hops);
1850 let stale = shared.withdrawn_through(prefix, &route.hops);
1851 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1853 return Err(Error::Closed);
1854 };
1855 if entry.hops.iter().next() != route.hops.iter().next()
1858 && let Some(server) = &entry.server
1859 {
1860 drop(std::mem::take(&mut server.lock().served));
1861 }
1862 entry.hops = route.hops.clone();
1863 entry.stale = stale;
1864 entry.cost = route.cost;
1865 entry.via = route.via;
1866 entry.advertised = true;
1867 let claim = entry.claim.clone();
1868 shared.sync_route(prefix, &claim);
1869 shared.prune_withdrawn(prefix);
1870 }
1871 Ok(())
1872 }
1873
1874 fn withdraw(&self) {
1879 let mut shared = self.shared.lock();
1880 for (prefix, id) in &self.entries {
1881 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1882 continue;
1883 };
1884 if !entry.advertised {
1885 continue;
1886 }
1887 entry.advertised = false;
1888 let claim = entry.claim.clone();
1889 shared.sync_route(prefix, &claim);
1890 }
1891 }
1892
1893 fn retract(&self, withdrawn: bool) {
1896 let mut shared = self.shared.lock();
1897 for (prefix, id) in &self.entries {
1898 let Some(entry) = shared.routes.remove(prefix, *id) else {
1899 continue;
1900 };
1901 if let Some(server) = &entry.server {
1904 let mut server = server.lock();
1905 server.closed = true;
1906 for producer in server.requests.drain_all() {
1907 if let Ok(mut request) = producer.write() {
1908 request.resolved.get_or_insert(Err(Error::Unroutable));
1909 }
1910 }
1911 }
1912 if withdrawn
1917 && let Some(&peer) = entry.hops.iter().last()
1918 && peer != Hop::UNKNOWN
1919 && !shared
1920 .routes
1921 .at(prefix)
1922 .any(|other| other.live() && other.hops.iter().last() == Some(&peer))
1923 {
1924 shared.withdrawn.entry(prefix.clone()).or_default().insert(peer);
1925 shared.restale(prefix);
1926 }
1927 shared.sync_route(&entry.prefix, &entry.claim);
1928 shared.prune_withdrawn(&entry.prefix);
1929 }
1930 }
1931}
1932
1933impl Drop for AnnounceProducer {
1934 fn drop(&mut self) {
1935 self.retract(false);
1936 }
1937}
1938
1939#[must_use = "poll the driver or the origin makes no progress"]
1954pub struct Driver {
1955 state: DriverState,
1956 timers: Clock,
1958 pool: cache::Pool,
1961}
1962
1963struct DriverState {
1965 set: TaskSet,
1967 shared: kio::Shared<OriginState>,
1970 done: bool,
1972}
1973
1974impl Driver {
1975 pub fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
1982 self.timers.advance(now);
1983 let result = self.state.poll(waiter);
1984 let gc = self.pool.gc(now);
1985 if result.is_ready() {
1986 return Err(Error::Closed);
1987 }
1988 Ok(self.timers.timeout().into_iter().chain(gc).min())
1989 }
1990}
1991
1992impl crate::time::Driver for Driver {
1993 fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
1994 self.poll(now, waiter)
1995 }
1996}
1997
1998impl DriverState {
1999 fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
2000 if !self.done {
2004 ready!(self.set.poll(waiter));
2005 self.done = true;
2006 }
2007 Poll::Ready(())
2008 }
2009
2010 fn teardown(&mut self) {
2014 drop(std::mem::replace(&mut self.set, TaskSet::owned()));
2017
2018 let (servers, cursors, fronts) = {
2023 let mut shared = self.shared.lock();
2024 shared.closed = true;
2025 shared.routes.poke_all();
2027 let servers: Vec<_> = shared
2028 .routes
2029 .entries()
2030 .filter_map(|entry| entry.server.clone())
2031 .collect();
2032 let cursors: Vec<_> = shared.cursors.values().map(|cursor| cursor.state.clone()).collect();
2033 let fronts: Vec<_> = shared.fronts.values().map(|front| front.request.clone()).collect();
2034 (servers, cursors, fronts)
2035 };
2036
2037 for producer in fronts {
2040 if let Ok(mut request) = producer.write() {
2041 request.resolved.get_or_insert(Err(Error::Dropped));
2042 }
2043 }
2044 for server in servers {
2048 let mut server = server.lock();
2049 server.closed = true;
2050 for producer in server.requests.drain_all() {
2051 if let Ok(mut request) = producer.write() {
2052 request.resolved.get_or_insert(Err(Error::Dropped));
2053 }
2054 }
2055 }
2056
2057 for state in cursors {
2060 if let Ok(mut state) = state.write() {
2061 state.ended = true;
2062 }
2063 }
2064 }
2065}
2066
2067impl Drop for DriverState {
2068 fn drop(&mut self) {
2069 self.teardown();
2070 }
2071}
2072
2073const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
2087
2088struct WarmCopy {
2093 track: track::Producer,
2094 _dynamic: track::Dynamic,
2095 edge: Option<WarmGroup>,
2098}
2099
2100impl Drop for WarmCopy {
2101 fn drop(&mut self) {
2102 let _ = self.track.finish();
2103 }
2104}
2105
2106struct WarmGroup(group::Producer);
2113
2114impl Drop for WarmGroup {
2115 fn drop(&mut self) {
2116 if !self.0.is_finished() {
2117 let _ = self.0.clone().abort(Error::Cancel);
2118 }
2119 }
2120}
2121
2122fn warm_copy(source: &track::Consumer, head: Option<&WarmGroup>) -> Option<WarmCopy> {
2129 let info = source.cached_info()?;
2130 let mut track = track::Producer::new(Arc::new(source.broadcast().clone()), source.name(), info);
2131 let head = head.map(|head| &head.0);
2132 let groups = source.cached_groups();
2133 let latest = groups
2135 .iter()
2136 .filter(|(_, visible)| *visible)
2137 .map(|(group, _)| group.sequence)
2138 .max();
2139 let mut edge = None;
2140 let mut first = None;
2141
2142 if let Some(head) = head
2146 && !groups.iter().any(|(group, _)| group.sequence == head.sequence)
2147 {
2148 let is_latest = latest.is_none_or(|latest| head.sequence >= latest);
2149 if head.is_finished() {
2150 let _ = track.adopt_group(head.clone(), true);
2151 first = Some(head.sequence);
2152 if is_latest {
2153 edge = Some(WarmGroup(head.clone()));
2154 }
2155 } else if is_latest {
2156 edge = warm_rebuild(&track, head, None);
2157 first = edge.as_ref().map(|_| head.sequence);
2158 }
2159 }
2160
2161 for (group, visible) in groups {
2162 let finished = group.is_finished();
2163 let whole = group.live_first_frame() == Some(0);
2164 let is_latest = visible && Some(group.sequence) == latest;
2165 let warm = if finished && whole {
2166 let _ = track.adopt_group(group.clone(), visible);
2167 Some(WarmGroup(group))
2168 } else if (finished || is_latest) && (whole || head.is_some_and(|head| head.sequence == group.sequence)) {
2169 warm_rebuild(&track, &group, head)
2175 } else {
2176 None
2178 };
2179 if visible && let Some(warm) = &warm {
2180 first = Some(first.map_or(warm.0.sequence, |first: u64| first.min(warm.0.sequence)));
2181 }
2182 if is_latest {
2183 edge = warm;
2184 }
2185 }
2186 track.start_at(first).ok()?;
2190 let dynamic = track.dynamic();
2191 Some(WarmCopy {
2192 track,
2193 _dynamic: dynamic,
2194 edge,
2195 })
2196}
2197
2198fn warm_rebuild(track: &track::Producer, live: &group::Producer, head: Option<&group::Producer>) -> Option<WarmGroup> {
2201 let mut start = live.live_first_frame()? as u64;
2203 let mut tail = live.consume();
2204 tail.start_at(start);
2205 let mut frames = Vec::new();
2206 if start > 0
2207 && let Some(head) = head.filter(|head| head.sequence == live.sequence)
2208 && let Some(head_start) = head.live_first_frame()
2209 {
2210 let mut head = head.consume();
2211 head.start_at(head_start as u64);
2212 while head.index() < start {
2213 match head.poll_read_frame(&kio::Waiter::noop()) {
2214 Poll::Ready(Ok(Some(frame))) => frames.push(frame),
2215 _ => break,
2216 }
2217 }
2218 match head.index() == start {
2219 true => start = head_start as u64,
2220 false => frames.clear(),
2222 }
2223 }
2224 while let Poll::Ready(Ok(Some(frame))) = tail.poll_read_frame(&kio::Waiter::noop()) {
2225 frames.push(frame);
2226 }
2227 if frames.is_empty() {
2228 return None;
2229 }
2230
2231 let rebuilt = WarmGroup(
2233 track
2234 .create_group(group::Info {
2235 sequence: live.sequence,
2236 })
2237 .ok()?,
2238 );
2239 let mut writer = rebuilt.0.clone();
2240 if start > 0 {
2241 writer.start_at(start).ok()?;
2242 }
2243 for frame in frames {
2244 writer.write_frame(frame.timestamp, frame.payload).ok()?;
2245 }
2246 if live.is_finished() {
2247 writer.finish().ok()?;
2248 }
2249 Some(rebuilt)
2250}
2251
2252struct FrontTask {
2254 shared: kio::Shared<OriginState>,
2256 broadcast: broadcast::Producer,
2258 path: PathOwned,
2260 horizon: Horizon,
2262 watch: Watch,
2264 request: kio::Producer<PendingBroadcast>,
2266 pin: kio::Lock<Pin>,
2268 timers: Clock,
2269}
2270
2271struct TrackIo {
2274 resume: super::resume::Producer,
2275 staged: Option<(u64, track::Consumer)>,
2277 query: Option<(u64, track::Consumer, track::Querying)>,
2279 copy: Option<(u64, track::Consumer)>,
2281 edge: Option<track::Position>,
2286 warm: Option<WarmCopy>,
2289 head: Option<WarmGroup>,
2292 used: bool,
2294}
2295
2296impl TrackIo {
2297 fn end(&mut self) {
2300 self.staged = None;
2301 self.query = None;
2302 self.copy = None;
2303 self.warm = None;
2304 self.head = None;
2305 }
2306}
2307
2308async fn run_front(task: FrontTask) {
2312 let FrontTask {
2313 shared,
2314 broadcast,
2315 path,
2316 horizon,
2317 watch,
2318 request,
2319 pin,
2320 timers,
2321 } = task;
2322
2323 enum Step {
2325 Assigned(Arc<str>, super::resume::Producer),
2326 Resolved(u64, Result<broadcast::Consumer, Error>),
2327 SourceClosed(u64),
2328 Info(Arc<str>, u64, Result<track::Info, Error>),
2329 Ended(Arc<str>, u64, Result<(), Error>),
2330 Demand(Arc<str>),
2331 Deadline,
2332 Table,
2333 }
2334
2335 let mut front = Front::new(TRACK_IDLE_LINGER);
2336 let mut sources: HashMap<u64, broadcast::Consumer> = HashMap::new();
2337 let mut next_source = 0u64;
2338 let mut upstream: Option<(u64, kio::Consumer<PendingBroadcast>)> = None;
2340 let mut tracks: HashMap<Arc<str>, TrackIo> = HashMap::new();
2341 let mut deadline = crate::runtime::Deadline::new(&timers);
2342 let mut seen = 0;
2344 let mut events: VecDeque<Event> = VecDeque::new();
2345
2346 let select = |front: &mut Front, sources: &HashMap<u64, broadcast::Consumer>, seen: &mut u64| -> Event {
2349 let table = shared.read();
2350 if table.closed {
2351 return Event::Closed;
2352 }
2353 *seen = watch.seen();
2355 front.retain_routes(|route| table.routes.covers(&path.as_path(), route));
2356 let best = table
2357 .best_route(&path.as_path(), horizon, front.pin(), front.refused_routes())
2358 .map(|entry| Candidate {
2359 route: entry.id,
2360 first: entry.hops.iter().next().copied(),
2361 local: entry.local,
2362 });
2363 let serving_closing = front
2364 .serving()
2365 .and_then(|id| sources.get(&id))
2366 .is_some_and(|source| source.is_closing());
2367 Event::Selected { best, serving_closing }
2368 };
2369
2370 events.push_back(select(&mut front, &sources, &mut seen));
2371
2372 loop {
2373 while let Some(event) = events.pop_front() {
2374 for action in front.step(event) {
2375 match action {
2376 Action::Reselect => events.push_back(select(&mut front, &sources, &mut seen)),
2377 Action::Request { route } => {
2378 let found = {
2380 let table = shared.read();
2381 table
2382 .routes
2383 .covering(&path.as_path())
2384 .find(|entry| entry.id == route && entry.live())
2385 .map(|entry| {
2386 (
2387 Candidate {
2388 route,
2389 first: entry.hops.iter().next().copied(),
2390 local: entry.local,
2391 },
2392 entry.source.clone(),
2393 entry.server.clone(),
2394 )
2395 })
2396 };
2397 let Some((candidate, source, server)) = found else {
2398 events.push_back(Event::Resolved {
2399 route,
2400 result: Err(Refusal {
2401 err: Error::Unroutable,
2402 standing: false,
2403 }),
2404 });
2405 continue;
2406 };
2407 front.identify(candidate);
2408 *pin.lock() = front.pin();
2409 if let Some(source) = source {
2410 let id = next_source;
2411 next_source += 1;
2412 sources.insert(id, source);
2413 events.push_back(Event::Resolved { route, result: Ok(id) });
2414 continue;
2415 }
2416 let Some(server) = server else {
2417 events.push_back(Event::Resolved {
2418 route,
2419 result: Err(Refusal {
2420 err: Error::Unroutable,
2421 standing: true,
2422 }),
2423 });
2424 continue;
2425 };
2426 let mut serve = server.lock();
2427 if serve.closed {
2428 drop(serve);
2431 events.push_back(Event::Resolved {
2432 route,
2433 result: Err(Refusal {
2434 err: Error::Unroutable,
2435 standing: true,
2436 }),
2437 });
2438 continue;
2439 }
2440 if let Some(weak) = serve.served.get(&path) {
2443 drop(serve);
2444 let id = next_source;
2445 next_source += 1;
2446 sources.insert(id, weak.consume());
2447 events.push_back(Event::Resolved { route, result: Ok(id) });
2448 continue;
2449 }
2450 let pending = match serve.requests.join(&path) {
2451 Some(producer) => producer.consume(),
2452 None => {
2453 let producer = kio::Producer::<PendingBroadcast>::default();
2454 let consumer = producer.consume();
2455 match serve.requests.insert(path.clone(), producer) {
2456 Ok(()) => consumer,
2457 Err(_) => {
2460 drop(serve);
2461 events.push_back(Event::Resolved {
2462 route,
2463 result: Err(Refusal {
2464 err: Error::Unroutable,
2465 standing: true,
2466 }),
2467 });
2468 continue;
2469 }
2470 }
2471 }
2472 };
2473 upstream = Some((route, pending));
2474 }
2475 Action::Detach { source } => {
2476 sources.remove(&source);
2477 for io in tracks.values_mut() {
2480 if io.copy.as_ref().is_some_and(|(s, _)| *s == source) {
2481 io.copy = None;
2482 }
2483 if io.query.as_ref().is_some_and(|(s, ..)| *s == source) {
2484 io.query = None;
2485 }
2486 if io.staged.as_ref().is_some_and(|(s, _)| *s == source) {
2487 io.staged = None;
2488 }
2489 }
2490 }
2491 Action::Resolve => {
2492 if let Ok(mut pending) = request.write() {
2493 pending.resolved.get_or_insert(Ok(broadcast.consume()));
2494 }
2495 }
2496 Action::Query { track: name, source } => {
2497 let Some(io) = tracks.get_mut(&name) else { continue };
2498 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2499 match sources.get(&source).map(|s| s.track(&name)) {
2500 Some(Ok(copy)) => {
2501 let query = copy.query().into_inner();
2504 io.query = Some((source, copy, query));
2505 }
2506 Some(Err(err)) => events.push_back(Event::TrackInfo {
2507 track: name,
2508 source,
2509 closing,
2510 result: Err(err),
2511 }),
2512 None => {}
2513 }
2514 }
2515 Action::Splice { track: name, source } => {
2516 let Some(io) = tracks.get_mut(&name) else { continue };
2517 let Some((staged, copy)) = io.staged.take() else {
2518 continue;
2519 };
2520 if staged != source {
2521 continue;
2522 }
2523 if let Err(err) = io.resume.takeover(©) {
2524 let _ = io.resume.abort(err);
2528 tracks.remove(&name);
2529 continue;
2530 }
2531 if let Some(head) = io.warm.take().and_then(|mut warm| warm.edge.take()) {
2532 io.head = Some(head);
2533 }
2534 io.edge = io.resume.resume_position();
2537 io.copy = Some((source, copy));
2538 }
2539 Action::Park { track: name } => {
2540 let Some(io) = tracks.get_mut(&name) else { continue };
2541 let Some((_, copy)) = io.copy.take() else { continue };
2542 let warm = warm_copy(©, io.head.as_ref());
2546 drop(copy);
2547 io.head = None;
2548 let parked = match &warm {
2549 Some(warm) => io.resume.park(&warm.track),
2550 None => io.resume.release(),
2551 };
2552 if parked.is_err() {
2553 tracks.remove(&name);
2554 continue;
2555 }
2556 io.warm = warm;
2557 }
2558 Action::Forget { track: name } => {
2559 if let Some(io) = tracks.get_mut(&name)
2564 && !broadcast.forget_spliced(&name, &io.resume)
2565 {
2566 if !io.used {
2567 io.used = true;
2568 events.push_back(Event::Used { track: name });
2569 }
2570 continue;
2571 }
2572 tracks.remove(&name);
2573 events.push_back(Event::Forgotten { track: name });
2574 }
2575 Action::Finish { track: name } => {
2576 if let Some(io) = tracks.get_mut(&name) {
2577 io.end();
2578 let _ = io.resume.finish();
2579 }
2580 }
2581 Action::Abort { track: name, err } => {
2582 if let Some(io) = tracks.get_mut(&name) {
2583 tracing::debug!(name = %name, %err, "aborting track");
2584 io.end();
2585 let _ = io.resume.abort(err);
2586 }
2587 }
2588 Action::Arm { at } => deadline.set(at),
2589 Action::End { err } => {
2590 if let Ok(mut pending) = request.write() {
2591 pending.resolved.get_or_insert(Err(err.clone()));
2592 }
2593 broadcast.close();
2600 broadcast.release_spliced(err.clone());
2601 for (_, mut io) in tracks.drain() {
2602 let waiting = io.staged.take().map(|(_, copy)| copy);
2606 let waiting = waiting.or_else(|| io.query.take().map(|(_, copy, _)| copy));
2607 if let Some(copy) = waiting
2608 && io.resume.is_used()
2609 {
2610 if io.resume.takeover(©).is_err() {
2611 continue;
2612 }
2613 io.warm = None;
2614 }
2615 if !io.resume.is_used() || !io.resume.is_spliced() || io.warm.is_some() {
2617 let _ = io.resume.abort(err.clone());
2618 }
2619 }
2620 return;
2621 }
2622 }
2623 }
2624 }
2625
2626 let step = kio::wait(|waiter| {
2627 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
2628 return Poll::Ready(Step::Assigned(name, resume));
2629 }
2630 if let Some((route, pending)) = &upstream
2631 && let Poll::Ready(result) = pending.poll(waiter, |p| match &p.resolved {
2632 Some(result) => Poll::Ready(result.clone()),
2633 None => Poll::Pending,
2634 }) {
2635 return Poll::Ready(Step::Resolved(
2636 *route,
2637 match result {
2638 Ok(resolved) => resolved,
2639 Err(_closed) => Err(Error::Unroutable),
2642 },
2643 ));
2644 }
2645 if let Some(id) = front.serving()
2646 && let Some(source) = sources.get(&id)
2647 && source.poll_closed(waiter).is_ready()
2648 {
2649 return Poll::Ready(Step::SourceClosed(id));
2650 }
2651 for (name, io) in &tracks {
2652 if let Some((source, _, query)) = &io.query
2653 && let Poll::Ready(result) = query.poll(waiter)
2654 {
2655 return Poll::Ready(Step::Info(name.clone(), *source, result));
2656 }
2657 if let Some((source, copy)) = &io.copy
2658 && let Poll::Ready(result) = copy.poll_complete(waiter)
2659 {
2660 return Poll::Ready(Step::Ended(name.clone(), *source, result));
2661 }
2662 let edge = match io.used {
2664 true => io.resume.poll_unused(waiter),
2665 false => io.resume.poll_used(waiter),
2666 };
2667 if edge.is_ready() {
2668 return Poll::Ready(Step::Demand(name.clone()));
2669 }
2670 }
2671 if deadline.poll(waiter).is_ready() {
2672 return Poll::Ready(Step::Deadline);
2673 }
2674 watch.poll_changed(waiter, seen).map(|()| Step::Table)
2675 })
2676 .await;
2677
2678 let event = match step {
2679 Step::Assigned(name, resume) => {
2680 tracks.insert(
2681 name.clone(),
2682 TrackIo {
2683 resume,
2684 staged: None,
2685 query: None,
2686 copy: None,
2687 edge: None,
2688 warm: None,
2689 head: None,
2690 used: false,
2691 },
2692 );
2693 Event::TrackAssigned {
2694 track: name,
2695 now: timers.now(),
2696 }
2697 }
2698 Step::Resolved(route, result) => {
2699 upstream = None;
2700 match result {
2701 Ok(source) => {
2702 let id = next_source;
2703 next_source += 1;
2704 sources.insert(id, source);
2705 Event::Resolved { route, result: Ok(id) }
2706 }
2707 Err(err) => {
2708 let standing =
2712 !matches!(err, Error::Unroutable) || shared.read().routes.covers(&path.as_path(), route);
2713 Event::Resolved {
2714 route,
2715 result: Err(Refusal { err, standing }),
2716 }
2717 }
2718 }
2719 }
2720 Step::SourceClosed(source) => Event::SourceClosed { source },
2721 Step::Info(name, source, result) => {
2722 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2723 let Some(io) = tracks.get_mut(&name) else { continue };
2724 let Some((_, copy, _)) = io.query.take() else { continue };
2725 let result = match result {
2728 Ok(info) => match copy.poll_complete(&kio::Waiter::noop()) {
2729 Poll::Ready(Err(err)) => Err(err),
2730 _ => Ok(info),
2731 },
2732 Err(err) => Err(err),
2733 };
2734 if result.is_ok() && io.used {
2737 io.staged = Some((source, copy));
2738 }
2739 Event::TrackInfo {
2740 track: name,
2741 source,
2742 closing,
2743 result,
2744 }
2745 }
2746 Step::Ended(name, source, result) => {
2747 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2748 let Some(io) = tracks.get_mut(&name) else { continue };
2749 io.copy = None;
2750 let delivered = io.resume.resume_position() != io.edge;
2751 Event::TrackEnded {
2752 track: name,
2753 source,
2754 closing,
2755 result,
2756 delivered,
2757 }
2758 }
2759 Step::Demand(name) => {
2760 let Some(io) = tracks.get_mut(&name) else { continue };
2761 io.used = io.resume.is_used();
2762 if !io.used {
2763 io.query = None;
2766 io.staged = None;
2767 }
2768 match io.used {
2769 true => Event::Used { track: name },
2770 false => Event::Unused {
2771 track: name,
2772 now: timers.now(),
2773 },
2774 }
2775 }
2776 Step::Deadline => {
2777 deadline.set(None);
2780 Event::Deadline { now: timers.now() }
2781 }
2782 Step::Table => select(&mut front, &sources, &mut seen),
2783 };
2784 events.push_back(event);
2785 }
2786}
2787
2788#[derive(Default)]
2793struct RouteTable {
2794 root: RouteNode,
2795}
2796
2797#[derive(Default)]
2800struct RouteNode {
2801 entries: Vec<RouteEntry>,
2803 cursors: Vec<ConsumerId>,
2805 cursors_below: usize,
2808 watches: Vec<(u64, kio::Producer<Watched>)>,
2811 watches_below: usize,
2814 children: HashMap<String, RouteNode>,
2815}
2816
2817#[derive(Default)]
2820struct Watched {
2821 generation: u64,
2822}
2823
2824struct Watch {
2830 shared: kio::Shared<OriginState>,
2831 path: PathOwned,
2832 id: u64,
2833 signal: kio::Consumer<Watched>,
2834}
2835
2836impl Watch {
2837 fn seen(&self) -> u64 {
2841 self.signal.read().generation
2842 }
2843
2844 fn poll_changed(&self, waiter: &kio::Waiter, seen: u64) -> Poll<()> {
2846 self.signal
2847 .poll(waiter, |watched| match watched.generation != seen {
2848 true => Poll::Ready(()),
2849 false => Poll::Pending,
2850 })
2851 .map(|_| ())
2852 }
2853}
2854
2855impl Drop for Watch {
2856 fn drop(&mut self) {
2857 self.shared.lock().routes.remove_watch(&self.path, self.id);
2858 }
2859}
2860
2861#[derive(Clone, Copy)]
2863struct Below {
2864 cursors: usize,
2865 watches: usize,
2866}
2867
2868impl Below {
2869 const NONE: Self = Self { cursors: 0, watches: 0 };
2870 const CURSOR: Self = Self { cursors: 1, watches: 0 };
2871 const WATCH: Self = Self { cursors: 0, watches: 1 };
2872}
2873
2874impl RouteNode {
2875 fn is_empty(&self) -> bool {
2877 self.entries.is_empty() && self.cursors.is_empty() && self.watches.is_empty() && self.children.is_empty()
2878 }
2879
2880 fn find<'a>(&self, mut parts: impl Iterator<Item = &'a str>) -> Option<&Self> {
2882 match parts.next() {
2883 None => Some(self),
2884 Some(part) => self.children.get(part)?.find(parts),
2885 }
2886 }
2887
2888 fn reach<'a>(&mut self, mut parts: impl Iterator<Item = &'a str>, below: Below) -> &mut Self {
2891 self.cursors_below += below.cursors;
2892 self.watches_below += below.watches;
2893 match parts.next() {
2894 None => self,
2895 Some(part) => self.children.entry(part.to_string()).or_default().reach(parts, below),
2896 }
2897 }
2898
2899 fn edit<'a, R>(
2903 &mut self,
2904 mut parts: impl Iterator<Item = &'a str>,
2905 below: Below,
2906 f: impl FnOnce(&mut Self) -> R,
2907 ) -> Option<R> {
2908 let result = match parts.next() {
2909 None => f(self),
2910 Some(part) => {
2911 let child = self.children.get_mut(part)?;
2912 let result = child.edit(parts, below, f)?;
2913 if child.is_empty() {
2914 self.children.remove(part);
2915 }
2916 result
2917 }
2918 };
2919 self.cursors_below -= below.cursors;
2920 self.watches_below -= below.watches;
2921 Some(result)
2922 }
2923
2924 fn poke(&self) {
2926 for (_, watch) in &self.watches {
2927 if let Ok(mut watched) = watch.write() {
2928 watched.generation += 1;
2929 }
2930 }
2931 }
2932
2933 fn poke_below(&self) {
2936 if self.watches_below == 0 {
2937 return;
2938 }
2939 self.poke();
2940 for child in self.children.values() {
2941 child.poke_below();
2942 }
2943 }
2944
2945 fn walk<'a>(&'a self, visit: &mut impl FnMut(&'a Self)) {
2947 visit(self);
2948 for child in self.children.values() {
2949 child.walk(visit);
2950 }
2951 }
2952
2953 fn collect_cursors(&self, out: &mut Vec<ConsumerId>) {
2955 if self.cursors_below == 0 {
2956 return;
2957 }
2958 out.extend(&self.cursors);
2959 for child in self.children.values() {
2960 child.collect_cursors(out);
2961 }
2962 }
2963}
2964
2965impl RouteTable {
2966 fn split(&self, path: &Path) -> (Vec<&RouteNode>, Option<&RouteNode>) {
2969 let mut above = Vec::new();
2970 let mut node = &self.root;
2971 for part in path.parts() {
2972 above.push(node);
2973 match node.children.get(part) {
2974 Some(child) => node = child,
2975 None => return (above, None),
2976 }
2977 }
2978 (above, Some(node))
2979 }
2980
2981 fn covering(&self, path: &Path) -> impl Iterator<Item = &RouteEntry> {
2983 let (above, at) = self.split(path);
2984 above.into_iter().chain(at).flat_map(|node| node.entries.iter())
2985 }
2986
2987 fn covers(&self, path: &Path, id: u64) -> bool {
2989 self.covering(path).any(|entry| entry.id == id && entry.live())
2990 }
2991
2992 fn at(&self, prefix: &Path) -> impl Iterator<Item = &RouteEntry> {
2994 self.root
2995 .find(prefix.parts())
2996 .into_iter()
2997 .flat_map(|node| node.entries.iter())
2998 }
2999
3000 fn at_mut(&mut self, prefix: &Path) -> impl Iterator<Item = &mut RouteEntry> {
3002 let mut node = Some(&mut self.root);
3003 for part in prefix.parts() {
3004 node = node.and_then(|node| node.children.get_mut(part));
3005 }
3006 node.into_iter().flat_map(|node| node.entries.iter_mut())
3007 }
3008
3009 fn entries(&self) -> impl Iterator<Item = &RouteEntry> {
3011 let mut nodes = Vec::new();
3012 self.root.walk(&mut |node| nodes.push(node));
3013 nodes.into_iter().flat_map(|node| node.entries.iter())
3014 }
3015
3016 fn insert(&mut self, entry: RouteEntry) {
3018 let node = self.root.reach(entry.prefix.parts(), Below::NONE);
3019 node.entries.push(entry);
3020 }
3021
3022 fn entry_mut(&mut self, prefix: &Path, id: u64) -> Option<&mut RouteEntry> {
3024 let mut node = &mut self.root;
3025 for part in prefix.parts() {
3026 node = node.children.get_mut(part)?;
3027 }
3028 node.entries.iter_mut().find(|entry| entry.id == id)
3029 }
3030
3031 fn remove(&mut self, prefix: &Path, id: u64) -> Option<RouteEntry> {
3033 self.root
3034 .edit(prefix.parts(), Below::NONE, |node| {
3035 let index = node.entries.iter().position(|entry| entry.id == id)?;
3036 Some(node.entries.swap_remove(index))
3037 })
3038 .flatten()
3039 }
3040
3041 fn add_cursor(&mut self, head: &Path, id: ConsumerId) {
3043 self.root.reach(head.parts(), Below::CURSOR).cursors.push(id);
3044 }
3045
3046 fn remove_cursor(&mut self, head: &Path, id: ConsumerId) {
3049 self.root.edit(head.parts(), Below::CURSOR, |node| {
3050 node.cursors.retain(|cursor| *cursor != id)
3051 });
3052 }
3053
3054 fn add_watch(&mut self, path: &Path, id: u64) -> kio::Consumer<Watched> {
3056 let producer = kio::Producer::<Watched>::default();
3057 let consumer = producer.consume();
3058 self.root.reach(path.parts(), Below::WATCH).watches.push((id, producer));
3059 consumer
3060 }
3061
3062 fn remove_watch(&mut self, path: &Path, id: u64) {
3065 self.root.edit(path.parts(), Below::WATCH, |node| {
3066 node.watches.retain(|(watch, _)| *watch != id)
3067 });
3068 }
3069
3070 fn poke_below(&self, prefix: &Path) {
3072 if let (_, Some(node)) = self.split(prefix) {
3073 node.poke_below();
3074 }
3075 }
3076
3077 fn poke_all(&self) {
3079 self.root.walk(&mut |node| node.poke());
3080 }
3081
3082 fn cursors_touching(&self, prefix: &Path) -> Vec<ConsumerId> {
3086 let (above, at) = self.split(prefix);
3087 let mut cursors: Vec<ConsumerId> = above.iter().flat_map(|node| node.cursors.iter().copied()).collect();
3088 if let Some(node) = at {
3089 node.collect_cursors(&mut cursors);
3090 }
3091 cursors.sort_unstable();
3093 cursors.dedup();
3094 cursors
3095 }
3096}
3097
3098#[derive(Default)]
3105struct OriginState {
3106 routes: RouteTable,
3109 next_route: u64,
3110 next_watch: u64,
3111
3112 cursors: HashMap<ConsumerId, TableCursor>,
3116
3117 fronts: WeakCache<FrontKey, RemoteFront>,
3125
3126 withdrawn: HashMap<PathOwned, HashSet<Hop>>,
3129
3130 closed: bool,
3133}
3134
3135impl OriginState {
3136 fn withdrawn_through(&self, prefix: &Path, hops: &Hops) -> bool {
3138 if self.withdrawn.is_empty() {
3139 return false;
3140 }
3141 self.withdrawn
3142 .get(prefix)
3143 .is_some_and(|peers| hops.iter().any(|hop| peers.contains(hop)))
3144 }
3145
3146 fn restale(&mut self, prefix: &Path) {
3148 let peers = self.withdrawn.get(prefix);
3149 let mut changed = None;
3150 for entry in self.routes.at_mut(prefix) {
3151 let stale = peers.is_some_and(|peers| entry.hops.iter().any(|hop| peers.contains(hop)));
3152 if entry.stale != stale {
3153 entry.stale = stale;
3154 changed = Some(entry.claim.clone());
3155 }
3156 }
3157 if let Some(claim) = changed {
3158 self.sync_route(prefix, &claim);
3159 }
3160 }
3161
3162 fn reannounced(&mut self, prefix: &PathOwned, hops: &Hops) {
3166 if let Some(sender) = hops.iter().last()
3167 && self.withdrawn.get_mut(prefix).is_some_and(|peers| peers.remove(sender))
3168 {
3169 self.restale(prefix);
3170 self.prune_withdrawn(prefix);
3171 }
3172 }
3173
3174 fn prune_withdrawn(&mut self, prefix: &PathOwned) {
3176 if self.withdrawn.is_empty() {
3177 return;
3178 }
3179 let Some(peers) = self.withdrawn.get_mut(prefix) else {
3180 return;
3181 };
3182 if peers.len() <= 1 {
3183 peers.retain(|peer| self.routes.at(prefix).any(|entry| entry.hops.contains(peer)));
3185 } else {
3186 let mut unreferenced = peers.clone();
3187 for entry in self.routes.at(prefix) {
3188 if unreferenced.is_empty() {
3189 break;
3190 }
3191 for hop in entry.hops.iter() {
3192 unreferenced.remove(hop);
3193 }
3194 }
3195 peers.retain(|peer| !unreferenced.contains(peer));
3196 }
3197 if peers.is_empty() {
3198 self.withdrawn.remove(prefix);
3199 }
3200 }
3201
3202 fn sync_route(&mut self, prefix: &Path, claim: &Pattern) {
3207 let routes = &self.routes;
3209 for id in routes.cursors_touching(prefix) {
3210 let Some(cursor) = self.cursors.get_mut(&id) else {
3211 continue;
3212 };
3213 if let Some(presented) = cursor.presented(prefix, claim) {
3214 Self::sync_cursor(routes, cursor, &presented);
3215 }
3216 }
3217 routes.poke_below(prefix);
3219 }
3220
3221 fn watch(&mut self, shared: &kio::Shared<OriginState>, path: &Path) -> Watch {
3223 let id = self.next_watch;
3224 self.next_watch += 1;
3225 let signal = self.routes.add_watch(path, id);
3226 Watch {
3227 shared: shared.clone(),
3228 path: path.to_owned(),
3229 id,
3230 signal,
3231 }
3232 }
3233
3234 fn sync_cursor(routes: &RouteTable, cursor: &mut TableCursor, presented: &PathOwned) {
3237 let candidates: Vec<&RouteEntry> = match presented.is_empty() {
3243 true => routes
3244 .covering(&cursor.root)
3245 .filter(|entry| cursor.visible(entry))
3246 .collect(),
3247 false => {
3248 let absolute = cursor.root.join(presented);
3249 routes.at(&absolute).filter(|entry| cursor.visible(entry)).collect()
3250 }
3251 };
3252 let most = candidates.iter().map(|entry| entry.prefix.len()).max();
3253 let best = most.and_then(|most| {
3254 candidates
3255 .into_iter()
3256 .filter(|entry| entry.prefix.len() == most)
3257 .min_by_key(|entry| route_order(&entry.prefix, entry))
3258 });
3259
3260 match best {
3261 Some(entry) => {
3262 let meta = (entry.hops.clone(), entry.cost, entry.entered());
3263 let served = entry.server.is_some();
3264 let captures = cursor.captures(&entry.prefix);
3265 let previous = cursor
3266 .current
3267 .insert(presented.clone(), (entry.id, meta.clone(), served, captures.clone()));
3268 match previous {
3269 Some((_, prev, prev_served, prev_captures))
3276 if prev == meta && prev_served == served && prev_captures == captures => {}
3277 Some((_, prev, _, prev_captures)) if prev_captures != captures => {
3280 if let Ok(mut state) = cursor.state.write() {
3281 state.apply_unannounce(cursor.under.join(presented), prev, prev_captures);
3282 state.apply_announce(cursor.under.join(presented), meta, captures);
3283 }
3284 }
3285 _ => {
3286 if let Ok(mut state) = cursor.state.write() {
3287 state.apply_announce(cursor.under.join(presented), meta, captures);
3288 }
3289 }
3290 }
3291 }
3292 None => {
3293 if let Some((_, last, _, captures)) = cursor.current.remove(presented)
3294 && let Ok(mut state) = cursor.state.write()
3295 {
3296 state.apply_unannounce(cursor.under.join(presented), last, captures);
3297 }
3298 }
3299 }
3300 }
3301
3302 fn register_cursor(&mut self, id: ConsumerId, mut cursor: TableCursor) {
3304 let mut presented: BTreeSet<PathOwned> = BTreeSet::new();
3307 for head in &cursor.heads {
3308 let (above, at) = self.routes.split(head);
3309 let mut nodes = above;
3310 if let Some(node) = at {
3311 node.walk(&mut |node| nodes.push(node));
3312 }
3313 for entry in nodes.into_iter().flat_map(|node| node.entries.iter()) {
3314 if let Some(p) = cursor.presented(&entry.prefix, &entry.claim) {
3315 presented.insert(p);
3316 }
3317 }
3318 }
3319 for p in &presented {
3320 Self::sync_cursor(&self.routes, &mut cursor, p);
3321 }
3322 for head in &cursor.heads {
3323 self.routes.add_cursor(head, id);
3324 }
3325 self.cursors.insert(id, cursor);
3326 }
3327
3328 fn best_route(&self, path: &Path, horizon: Horizon, pin: Pin, refused: &HashSet<u64>) -> Option<&RouteEntry> {
3343 let (above, at) = self.routes.split(path);
3347 let mut best = None;
3348 for node in above.into_iter().chain(at) {
3349 let mut candidates = node
3350 .entries
3351 .iter()
3352 .filter(|entry| entry.live())
3353 .filter(|entry| entry.scope.matches(path.as_str()))
3354 .filter(|entry| horizon.admits(entry))
3355 .filter(|entry| entry.qualifies(pin))
3356 .filter(|entry| !refused.contains(&entry.id))
3357 .peekable();
3358 if candidates.peek().is_some() {
3359 best = candidates
3360 .filter(|entry| entry.serves(path))
3361 .min_by_key(|entry| route_order(&entry.prefix, entry));
3362 }
3363 }
3364 best
3365 }
3366}
3367
3368#[derive(Default)]
3375struct PendingBroadcast {
3376 resolved: Option<Result<broadcast::Consumer, Error>>,
3377}
3378
3379#[must_use = "dropping an origin::Dynamic retracts the route"]
3394pub struct Dynamic {
3395 announcement: AnnounceProducer,
3397 state: kio::Shared<ServeState>,
3398 _keepalive: Keepalive,
3402}
3403
3404impl Dynamic {
3405 pub(crate) fn withdrawn(self) {
3408 self.announcement.retract(true);
3409 }
3410
3411 pub fn update(&self, route: Route) -> Result<(), Error> {
3421 self.announcement.update(route)
3422 }
3423
3424 pub fn poll_requested_broadcast(&self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
3429 let mut state = ready!(self.state.poll(waiter, |state| {
3430 if state.closed || state.requests.has_queued() {
3431 Poll::Ready(())
3432 } else {
3433 Poll::Pending
3434 }
3435 }));
3436
3437 if state.closed {
3439 return Poll::Ready(Err(Error::Closed));
3440 }
3441
3442 let path = state.requests.pop().expect("predicate guaranteed a request");
3443 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
3449 Poll::Ready(Ok(Request {
3450 path,
3451 producer,
3452 home: self.state.clone(),
3453 }))
3454 }
3455
3456 pub async fn requested_broadcast(&self) -> Result<Request, Error> {
3462 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
3463 }
3464}
3465
3466impl ServeState {
3467 fn resolve(
3478 shared: &kio::Shared<Self>,
3479 path: &PathOwned,
3480 producer: &kio::Producer<PendingBroadcast>,
3481 result: Result<broadcast::Consumer, Error>,
3482 ) {
3483 let mut state = shared.lock();
3484 if state.closed {
3485 return;
3486 }
3487 let resolved = match result {
3488 Ok(broadcast) => {
3489 let existing = state.served.insert(path.clone(), broadcast.weak());
3493 Ok(existing.map(|weak| weak.consume()).unwrap_or(broadcast))
3494 }
3495 Err(err) => Err(err),
3496 };
3497 state.requests.remove_if(path, |p| p.same_channel(producer));
3498 if let Ok(mut pending) = producer.write() {
3499 pending.resolved.get_or_insert(resolved);
3500 drop(state);
3501 }
3502 }
3503
3504 fn forget(shared: &kio::Shared<Self>, path: &PathOwned, producer: &kio::Producer<PendingBroadcast>) {
3506 shared.lock().requests.remove_if(path, |p| p.same_channel(producer));
3507 }
3508}
3509
3510pub struct Request {
3517 path: PathOwned,
3519
3520 producer: kio::Producer<PendingBroadcast>,
3523
3524 home: kio::Shared<ServeState>,
3527}
3528
3529impl Request {
3530 pub fn path(&self) -> &Path<'_> {
3532 &self.path
3533 }
3534
3535 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
3541 let broadcast = broadcast.consume();
3542 ServeState::resolve(&self.home, &self.path, &self.producer, Ok(broadcast));
3543 }
3545
3546 pub fn reject(self, err: Error) {
3548 ServeState::resolve(&self.home, &self.path, &self.producer, Err(err));
3549 }
3550}
3551
3552impl Drop for Request {
3553 fn drop(&mut self) {
3554 ServeState::forget(&self.home, &self.path, &self.producer);
3562 }
3563}
3564
3565pub struct Requesting {
3572 inner: RequestState,
3573 path: PathOwned,
3577 stats: stats::Scope,
3580}
3581
3582enum RequestState {
3583 Failed(Error),
3586 Pending(kio::Consumer<PendingBroadcast>),
3588}
3589
3590impl Requesting {
3591 fn failed(error: Error) -> Self {
3592 Self::new(RequestState::Failed(error))
3593 }
3594
3595 fn queued(consumer: kio::Consumer<PendingBroadcast>) -> Self {
3596 Self::new(RequestState::Pending(consumer))
3597 }
3598
3599 pub fn is_queued(&self) -> bool {
3608 matches!(self.inner, RequestState::Pending(_))
3609 }
3610
3611 fn new(inner: RequestState) -> Self {
3612 Self {
3613 inner,
3614 path: PathOwned::default(),
3615 stats: stats::Scope::default(),
3616 }
3617 }
3618
3619 fn with_path(mut self, path: PathOwned) -> Self {
3620 self.path = path;
3621 self
3622 }
3623
3624 fn with_stats(mut self, scope: stats::Scope) -> Self {
3626 self.stats = scope;
3627 self
3628 }
3629
3630 fn hand_out(&self, broadcast: broadcast::Consumer) -> broadcast::Consumer {
3632 broadcast.with_path(self.path.clone()).with_stats(self.stats.clone())
3633 }
3634
3635 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
3637 match &self.inner {
3638 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
3639 RequestState::Pending(consumer) => Poll::Ready(
3640 match ready!(consumer.poll(waiter, |state| match &state.resolved {
3641 Some(result) => Poll::Ready(result.clone()),
3642 None => Poll::Pending,
3643 })) {
3644 Ok(result) => result.map(|broadcast| self.hand_out(broadcast)),
3645 Err(_closed) => Err(Error::Unroutable),
3647 },
3648 ),
3649 }
3650 }
3651}
3652
3653impl kio::Pollable for Requesting {
3654 type Output = Result<broadcast::Consumer, Error>;
3655
3656 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
3657 self.poll_ok(waiter)
3658 }
3659}
3660
3661pub trait Consume<T> {
3669 fn consume(&self) -> T;
3671}
3672
3673impl<T, U: Consume<T>> Consume<T> for &U {
3674 fn consume(&self) -> T {
3675 (**self).consume()
3676 }
3677}
3678
3679impl Consume<Consumer> for Producer {
3680 fn consume(&self) -> Consumer {
3681 Consumer::from_producer(self, stats::Session::default())
3685 }
3686}
3687
3688impl Consume<Consumer> for Consumer {
3689 fn consume(&self) -> Consumer {
3690 self.clone()
3691 }
3692}
3693
3694impl Consume<broadcast::Consumer> for broadcast::Producer {
3695 fn consume(&self) -> broadcast::Consumer {
3696 self.consume()
3698 }
3699}
3700
3701impl Consume<broadcast::Consumer> for broadcast::Consumer {
3702 fn consume(&self) -> broadcast::Consumer {
3703 self.clone()
3704 }
3705}
3706
3707impl Consume<track::Consumer> for track::Producer {
3708 fn consume(&self) -> track::Consumer {
3709 self.consume()
3710 }
3711}
3712
3713impl Consume<track::Consumer> for track::Consumer {
3714 fn consume(&self) -> track::Consumer {
3715 self.clone()
3716 }
3717}
3718
3719#[derive(Clone)]
3725pub struct Consumer {
3726 hop: Hop,
3728 scope: OriginScope,
3729
3730 root: PathOwned,
3732
3733 shared: kio::Shared<OriginState>,
3736
3737 stats: stats::Session,
3741
3742 horizon: Horizon,
3747
3748 hidden: Hidden,
3750
3751 pool: cache::Pool,
3754 cache_duration: Duration,
3755
3756 tasks: TasksWeak,
3760
3761 timers: Clock,
3764}
3765
3766impl Consumer {
3767 fn from_producer(producer: &Producer, stats: stats::Session) -> Self {
3768 Self {
3769 hop: producer.hop,
3770 scope: producer.scope.clone(),
3771 root: producer.root.clone(),
3772 shared: producer.shared.clone(),
3773 stats,
3774 horizon: Horizon::default(),
3775 hidden: Hidden::default(),
3776 pool: producer.pool.clone(),
3777 cache_duration: producer.cache_duration,
3778 tasks: producer.tasks.downgrade(),
3779 timers: producer.timers.clone(),
3780 }
3781 }
3782
3783 pub fn hop(&self) -> Hop {
3785 self.hop
3786 }
3787
3788 pub(crate) fn excluding(mut self, peer: Hop) -> Self {
3796 self.horizon.exclude = Some(peer);
3797 self
3798 }
3799
3800 pub fn local(mut self) -> Self {
3807 self.horizon.local = true;
3808 self
3809 }
3810
3811 pub fn with_hidden(mut self, hidden: bool) -> Self {
3816 self.hidden.include = hidden;
3817 self
3818 }
3819
3820 pub(crate) fn beyond(mut self, outer: &Consumer) -> Self {
3823 self.hidden.beyond = Some(
3824 outer
3825 .hidden
3826 .from
3827 .clone()
3828 .map(|head| vec![head])
3829 .unwrap_or_else(|| interest_prefixes(&outer.scope.allowed)),
3830 );
3831 self
3832 }
3833
3834 pub(crate) fn discovery(mut self, hidden: bool) -> Self {
3836 self.hidden.include = hidden;
3837 self.hidden.from = Some(self.root.clone());
3838 self
3839 }
3840
3841 pub(crate) fn includes_hidden(&self) -> bool {
3843 self.hidden.include
3844 }
3845
3846 pub fn with_stats(mut self, session: stats::Session) -> Self {
3850 self.stats = session;
3851 self
3852 }
3853
3854 fn untagged(&self) -> Self {
3858 Self {
3859 stats: stats::Session::default(),
3860 ..self.clone()
3861 }
3862 }
3863
3864 pub(crate) fn empty(&self) -> Self {
3869 Self {
3870 scope: OriginScope::empty(),
3871 ..self.clone()
3872 }
3873 }
3874
3875 pub fn announced(&self) -> AnnounceConsumer {
3884 let state = kio::Producer::<OriginConsumerState>::default();
3885 let cursor = |root: PathOwned,
3886 allowed: Patterns,
3887 hidden: Hidden,
3888 mount: Option<&Mount>,
3889 under: PathOwned,
3890 holes: Vec<PathOwned>| TableCursor {
3891 root,
3892 heads: interest_prefixes(&allowed),
3893 allowed,
3894 horizon: self.horizon,
3895 hidden,
3896 state: state.clone(),
3897 under,
3898 holes,
3899 named: mount.map(|mount| (mount.clone(), self.scope.allowed.clone())),
3900 current: HashMap::new(),
3901 };
3902
3903 if let Some(mount) = self.scope.mount(&self.root) {
3906 let Some(root) = mount.resolve(&self.root) else {
3908 return AnnounceConsumer::new(self.root.clone(), Vec::new(), state, self.stats.clone(), &self.shared);
3909 };
3910 let cursors = vec![cursor(
3911 root,
3912 mount.translate(&self.scope.allowed),
3913 self.hidden.translate(mount),
3914 Some(mount),
3915 PathOwned::default(),
3916 Vec::new(),
3917 )];
3918 return AnnounceConsumer::new(self.root.clone(), cursors, state, self.stats.clone(), &self.shared);
3919 }
3920
3921 let heads = interest_prefixes(&self.scope.allowed);
3924 let mut holes = Vec::new();
3925 let mut cursors = Vec::new();
3926 for mount in self.scope.mounts.iter() {
3927 let Some(under) = mount.at.strip_prefix(&self.root) else {
3928 continue;
3929 };
3930 holes.push(mount.at.clone());
3931 let allowed = mount.translate(&self.scope.allowed);
3932 if allowed.is_empty()
3935 || !(self.hidden.include
3936 || !hides(
3937 self.hidden.from.as_ref().map(std::slice::from_ref).unwrap_or(&heads),
3938 &mount.at,
3939 )) {
3940 continue;
3941 }
3942 cursors.push(cursor(
3943 mount.target.clone(),
3944 allowed,
3945 self.hidden.translate(mount),
3946 Some(mount),
3947 under.to_owned(),
3948 Vec::new(),
3949 ));
3950 }
3951 cursors.insert(
3952 0,
3953 cursor(
3954 self.root.clone(),
3955 self.scope.allowed.clone(),
3956 self.hidden.clone(),
3957 None,
3958 PathOwned::default(),
3959 holes,
3960 ),
3961 );
3962 AnnounceConsumer::new(self.root.clone(), cursors, state, self.stats.clone(), &self.shared)
3963 }
3964
3965 pub fn consume(&self) -> Self {
3967 self.clone()
3968 }
3969
3970 #[cfg(test)]
3973 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
3974 let full = self.root.join(path).to_owned();
3975 if !self.scope.permits(&full) {
3976 return None;
3977 }
3978 let full = self.scope.resolve(&full)?;
3979 let table = self.shared.lock();
3980 table
3981 .routes
3982 .at(&full)
3983 .filter(|entry| entry.local)
3984 .min_by_key(|entry| route_order(&entry.prefix, entry))
3985 .and_then(|entry| entry.source.clone())
3986 }
3987
3988 pub async fn routed(&self, path: impl AsPath) -> Option<Route> {
3998 let path = path.as_path();
3999
4000 let consumer = match Pattern::subtree(path.as_str()) {
4005 Ok(subtree) => self.scope("", &Patterns::from(subtree)).ok()?,
4006 Err(InvalidPattern::TooManySegments) => self.clone(),
4007 Err(_) => return None,
4008 };
4009
4010 if !consumer.allowed().matches(path.as_str()) {
4014 return None;
4015 }
4016
4017 let mut announced = consumer.untagged().with_hidden(true).announced();
4021 loop {
4022 let update = announced.next().await?;
4023 if update.kind.is_active() && path.has_prefix(&update.prefix) {
4024 return Some(update.route);
4025 }
4026 }
4027 }
4028
4029 pub async fn routed_broadcast(&self, path: impl AsPath) -> Result<broadcast::Consumer, Error> {
4043 let path = path.as_path();
4044
4045 if !self.allowed().matches(path.as_str()) {
4048 return Err(Error::Unauthorized);
4049 }
4050 loop {
4051 let (watch, seen) = {
4060 let mut table = self.shared.lock();
4061 if table.closed {
4062 return Err(Error::Closed);
4063 }
4064 let named = self.root.join(&path);
4065 let resolved = self.scope.resolve(&named).ok_or(BoundsExceeded)?;
4066 let watch = table.watch(&self.shared, &resolved);
4067 let seen = watch.seen();
4068 (watch, seen)
4069 };
4070 match self.request_broadcast(&path).await {
4071 Ok(broadcast) => return Ok(broadcast),
4072 Err(Error::Unroutable) => {
4073 kio::wait(|waiter| watch.poll_changed(waiter, seen)).await;
4074 }
4075 Err(Error::Dropped) if self.shared.lock().closed => return Err(Error::Closed),
4078 Err(err) => return Err(err),
4079 }
4080 }
4081 }
4082
4083 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Consumer, Error> {
4090 let root = self.root.join(root).to_owned();
4091 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
4092 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
4093 Ok(Consumer {
4094 scope,
4095 root,
4096 ..self.clone()
4097 })
4098 }
4099
4100 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
4123 let path = path.as_path();
4124
4125 let named = self.root.join(&path).to_owned();
4129 let scope = self.stats.egress(&named);
4130 let requested = path.to_owned();
4134
4135 if !self.scope.permits(&named) {
4137 return kio::Pending::new(Requesting::failed(Error::Unauthorized));
4138 }
4139
4140 let Some(absolute) = self.scope.resolve(&named).map(|path| path.to_owned()) else {
4143 return kio::Pending::new(Requesting::failed(BoundsExceeded.into()));
4144 };
4145
4146 let mut state = self.shared.lock();
4147
4148 if state.closed {
4150 return kio::Pending::new(Requesting::failed(Error::Closed));
4151 }
4152
4153 if state
4157 .best_route(&absolute.as_path(), self.horizon, Pin::Any, &HashSet::new())
4158 .is_none()
4159 {
4160 return kio::Pending::new(Requesting::failed(Error::Unroutable));
4161 }
4162
4163 let key = (absolute.clone(), self.horizon);
4171 if let Some(front) = state.fronts.get(&key) {
4172 let pin = *front.pin.lock();
4173 let current = state
4174 .best_route(&absolute.as_path(), self.horizon, Pin::Any, &HashSet::new())
4175 .is_some_and(|entry| entry.qualifies(pin));
4176 if current {
4177 let pending = Requesting::queued(front.request.consume())
4178 .with_path(requested)
4179 .with_stats(scope);
4180 return kio::Pending::new(pending);
4181 }
4182 state.fronts.remove(&key);
4183 }
4184
4185 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
4190 pool: self.pool.clone(),
4191 cache_duration: self.cache_duration,
4192 path: absolute.clone(),
4193 });
4194 let request = kio::Producer::<PendingBroadcast>::default();
4195 let consumer = request.consume();
4196 let watch = state.watch(&self.shared, &absolute);
4197 let pin = kio::Lock::new(Pin::Any);
4198 state.fronts.insert(
4199 key,
4200 RemoteFront {
4201 request: request.clone(),
4202 broadcast: broadcast.consume().weak(),
4203 pin: pin.clone(),
4204 },
4205 );
4206 drop(state);
4209 self.tasks.push(run_front(FrontTask {
4210 shared: self.shared.clone(),
4211 broadcast,
4212 path: absolute,
4213 horizon: self.horizon,
4214 watch,
4215 request,
4216 pin,
4217 timers: self.timers.clone(),
4218 }));
4219 kio::Pending::new(Requesting::queued(consumer).with_path(requested).with_stats(scope))
4220 }
4221
4222 pub fn root(&self) -> &Path<'_> {
4224 &self.root
4225 }
4226
4227 pub fn allowed(&self) -> Patterns {
4229 self.scope.relative(&self.root)
4230 }
4231
4232 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
4234 self.root.join(path)
4235 }
4236}
4237
4238pub struct AnnounceConsumer {
4243 ids: Vec<ConsumerId>,
4246 shared: kio::Shared<OriginState>,
4247 root: PathOwned,
4248
4249 state: kio::Producer<OriginConsumerState>,
4252
4253 stats: stats::Session,
4256
4257 guards: HashMap<PathOwned, stats::Announce>,
4261
4262 park: kio::Park,
4265}
4266
4267impl AnnounceConsumer {
4268 fn new(
4269 root: PathOwned,
4270 cursors: Vec<TableCursor>,
4271 state: kio::Producer<OriginConsumerState>,
4272 stats: stats::Session,
4273 shared: &kio::Shared<OriginState>,
4274 ) -> Self {
4275 let mut ids = Vec::with_capacity(cursors.len());
4276 {
4277 let mut table = shared.lock();
4278 if table.closed {
4279 if let Ok(mut state) = state.write() {
4281 state.ended = true;
4282 }
4283 } else {
4284 for cursor in cursors {
4285 let id = ConsumerId::new();
4286 table.register_cursor(id, cursor);
4287 ids.push(id);
4288 }
4289 }
4290 }
4291
4292 Self {
4293 ids,
4294 shared: shared.clone(),
4295 root,
4296 state,
4297 stats,
4298 guards: HashMap::new(),
4299 park: kio::Park::default(),
4300 }
4301 }
4302
4303 fn hand_out(&mut self, update: AnnounceUpdate) -> AnnounceUpdate {
4305 let absolute = self.root.join(&update.prefix).to_owned();
4306 if update.kind.is_active() {
4307 let scope = self.stats.egress(&absolute);
4308 self.guards
4309 .entry(update.prefix.clone())
4310 .or_insert_with(|| scope.announce());
4311 } else {
4312 self.guards.remove(&update.prefix);
4313 }
4314 update
4315 }
4316
4317 pub async fn next(&mut self) -> Option<AnnounceUpdate> {
4325 kio::wait(|waiter| self.poll_next(waiter)).await
4326 }
4327
4328 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<AnnounceUpdate>> {
4334 let update = {
4335 let mut state = match ready!(self.state.poll(waiter, |state| {
4336 if state.pending.is_empty() && !state.ended {
4337 Poll::Pending
4338 } else {
4339 Poll::Ready(())
4340 }
4341 })) {
4342 Ok(state) => state,
4343 Err(_) => return Poll::Ready(None),
4345 };
4346 match state.take() {
4347 Some(update) => update,
4348 None => {
4349 state.close();
4352 return Poll::Ready(None);
4353 }
4354 }
4355 };
4356 Poll::Ready(Some(self.hand_out(update)))
4357 }
4358
4359 pub fn try_next(&mut self) -> Option<AnnounceUpdate> {
4364 let update = self.state.write().ok()?.take()?;
4365 Some(self.hand_out(update))
4366 }
4367
4368 pub fn is_closed(&self) -> bool {
4370 let state = self.state.read();
4371 state.is_closed() || state.ended
4372 }
4373
4374 pub fn root(&self) -> &Path<'_> {
4376 &self.root
4377 }
4378
4379 pub fn absolute(&self, prefix: impl AsPath) -> Path<'_> {
4381 self.root.join(prefix)
4382 }
4383}
4384
4385impl futures::Stream for AnnounceConsumer {
4386 type Item = AnnounceUpdate;
4387
4388 fn poll_next(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> {
4389 let this = self.get_mut();
4390 let waiter = this.park.hold(cx).clone();
4391 this.poll_next(&waiter)
4392 }
4393}
4394
4395impl Drop for AnnounceConsumer {
4396 fn drop(&mut self) {
4397 let mut shared = self.shared.lock();
4398 for id in &self.ids {
4399 if let Some(cursor) = shared.cursors.remove(id) {
4400 for head in &cursor.heads {
4401 shared.routes.remove_cursor(head, *id);
4402 }
4403 }
4404 }
4405 }
4406}
4407
4408#[cfg(test)]
4409use futures::FutureExt;
4410
4411#[cfg(test)]
4412#[allow(missing_docs)] impl AnnounceConsumer {
4414 pub fn assert_next_active(&mut self, expected: impl AsPath) -> Route {
4416 let expected = expected.as_path();
4417 let update = self.next().now_or_never().expect("next blocked").expect("no next");
4418 assert_eq!(update.prefix, expected, "wrong prefix");
4419 assert!(update.kind.is_active(), "should be an active route");
4420 update.route
4421 }
4422
4423 pub fn assert_try_next_active(&mut self, expected: impl AsPath) -> Route {
4425 let expected = expected.as_path();
4426 let update = self.try_next().expect("no next");
4427 assert_eq!(update.prefix, expected, "wrong prefix");
4428 assert!(update.kind.is_active(), "should be an active route");
4429 update.route
4430 }
4431
4432 pub fn assert_next_ended(&mut self, expected: impl AsPath) {
4434 let expected = expected.as_path();
4435 let update = self.next().now_or_never().expect("next blocked").expect("no next");
4436 assert_eq!(update.prefix, expected, "wrong prefix");
4437 assert_eq!(update.kind, AnnounceKind::Retracted, "should be a retraction");
4438 }
4439
4440 pub fn assert_next_wait(&mut self) {
4441 if let Some(res) = self.next().now_or_never() {
4442 panic!("next should block: got {:?}", res.map(|u| u.prefix));
4443 }
4444 }
4445}
4446
4447#[cfg(test)]
4451pub(crate) trait ProduceTest {
4452 fn produce(self) -> Producer;
4453}
4454
4455#[cfg(test)]
4456impl ProduceTest for Config {
4457 fn produce(self) -> Producer {
4458 let (producer, driver) = Producer::new(self);
4459 if tokio::runtime::Handle::try_current().is_ok() {
4460 tokio::spawn(crate::time::run(driver));
4461 } else {
4462 std::mem::forget(driver);
4465 }
4466 producer
4467 }
4468}
4469
4470#[cfg(test)]
4471impl ProduceTest for Hop {
4472 fn produce(self) -> Producer {
4473 Config::new(self).produce()
4474 }
4475}
4476
4477#[cfg(test)]
4478mod tests {
4479 use super::*;
4480 use futures::FutureExt;
4481
4482 fn origin(id: u64) -> Hop {
4483 Hop::new(id).unwrap()
4484 }
4485
4486 fn hops(ids: &[u64]) -> Hops {
4487 let mut list = Hops::new();
4488 for &id in ids {
4489 list.push(if id == 0 { Hop::UNKNOWN } else { origin(id) }).unwrap();
4490 }
4491 list
4492 }
4493
4494 fn scopes(prefixes: &[&str]) -> Patterns {
4496 prefixes
4497 .iter()
4498 .map(|prefix| Pattern::subtree(prefix).unwrap())
4499 .collect()
4500 }
4501
4502 #[test]
4503 fn default_config_mints_a_real_hop() {
4504 let config = Config::default();
4505 assert_ne!(config.hop, Hop::UNKNOWN);
4506 let (producer, _driver) = Producer::new(config.clone());
4507 assert_eq!(producer.hop(), config.hop);
4508 assert_eq!(producer.consume().hop(), config.hop);
4509 }
4510
4511 #[test]
4512 fn random_hops_fit_legacy_lite_clients() {
4513 for _ in 0..32 {
4514 assert!(Hop::random().id() < 1u64 << 53);
4515 }
4516 }
4517
4518 async fn settle(mut check: impl FnMut() -> bool) {
4521 for _ in 0..100 {
4522 if check() {
4523 return;
4524 }
4525 tokio::task::yield_now().await;
4526 }
4527 panic!("condition never settled");
4528 }
4529
4530 async fn queued(server: &Dynamic) -> Request {
4532 let mut request = None;
4533 settle(|| match server.poll_requested_broadcast(&kio::Waiter::noop()) {
4534 Poll::Ready(Ok(popped)) => {
4535 request = Some(popped);
4536 true
4537 }
4538 _ => false,
4539 })
4540 .await;
4541 request.unwrap()
4542 }
4543
4544 async fn next_group(subscription: &mut crate::track::Subscriber) -> Result<Option<crate::group::Consumer>, Error> {
4546 let mut next = None;
4547 settle(|| match subscription.poll_recv_group(&kio::Waiter::noop()) {
4548 Poll::Ready(result) => {
4549 next = Some(result);
4550 true
4551 }
4552 Poll::Pending => false,
4553 })
4554 .await;
4555 next.unwrap()
4556 }
4557
4558 #[tokio::test]
4559 async fn announce_and_retract() {
4560 let producer = origin(1).produce();
4561 let consumer = producer.consume();
4562 let mut announced = consumer.announced();
4563 announced.assert_next_wait();
4564
4565 let announcement = producer.announce("room/alice", Route::default()).unwrap();
4566 let route = announced.assert_next_active("room/alice");
4567 assert!(route.hops.is_empty());
4568 assert_eq!(route.cost, Cost::default());
4569 announced.assert_next_wait();
4570
4571 drop(announcement);
4572 announced.assert_next_ended("room/alice");
4573 announced.assert_next_wait();
4574 }
4575
4576 #[tokio::test]
4579 async fn hidden_routes_need_an_opt_in() {
4580 let producer = origin(1).produce();
4581 let consumer = producer.consume();
4582 let _visible = producer.announce("room/alice", Route::default()).unwrap();
4583 let _stats = producer.announce(".stats/node", Route::default()).unwrap();
4584 let _nested = producer.announce("room/.internal", Route::default()).unwrap();
4585 let _suffix = producer.announce("room/catalog.pro", Route::default()).unwrap();
4587
4588 let mut announced = consumer.announced();
4589 announced.assert_next_active("room/alice");
4590 announced.assert_next_active("room/catalog.pro");
4591 announced.assert_next_wait();
4592
4593 let mut announced = consumer.clone().with_hidden(true).announced();
4594 announced.assert_next_active(".stats/node");
4595 announced.assert_next_active("room/.internal");
4596 announced.assert_next_active("room/alice");
4597 announced.assert_next_active("room/catalog.pro");
4598 announced.assert_next_wait();
4599
4600 let mut announced = consumer
4602 .scope(".stats", &Patterns::from(Pattern::all()))
4603 .unwrap()
4604 .announced();
4605 announced.assert_next_active("node");
4606 announced.assert_next_wait();
4607 let mut announced = consumer.scope("", &scopes(&["room/.internal"])).unwrap().announced();
4608 announced.assert_next_active("room/.internal");
4609 announced.assert_next_wait();
4610
4611 let mut announced = consumer.clone().with_hidden(true).beyond(&consumer).announced();
4614 announced.assert_next_active(".stats/node");
4615 announced.assert_next_active("room/.internal");
4616 announced.assert_next_wait();
4617 let room = consumer.scope("", &scopes(&["room"])).unwrap().beyond(&consumer);
4618 let mut announced = room.announced();
4619 announced.assert_next_wait();
4620 let stats = consumer.scope("", &scopes(&[".stats"])).unwrap().beyond(&consumer);
4621 let mut announced = stats.announced();
4622 announced.assert_next_active(".stats/node");
4623 announced.assert_next_wait();
4624 }
4625
4626 #[tokio::test]
4628 async fn hidden_broadcast_resolves_by_path() {
4629 let producer = origin(1).produce();
4630 let consumer = producer.consume();
4631 let broadcast = producer.create_broadcast(".stats/node").unwrap();
4632 broadcast.announce(Route::default()).unwrap();
4633
4634 consumer.announced().assert_next_wait();
4635 let resolved = consumer.request_broadcast(".stats/node").await.expect("resolves");
4636 assert_eq!(resolved.info().path.as_str(), ".stats/node");
4637 }
4638
4639 fn mounted(producer: &Producer, pid: &str, patterns: &[&str]) -> Consumer {
4641 let patterns: Patterns = patterns.iter().map(|pattern| pattern.parse().unwrap()).collect();
4642 producer
4643 .mount(format!("{pid}/.svc"), format!(".svc/{pid}"))
4644 .unwrap()
4645 .scope(pid, &patterns)
4646 .unwrap()
4647 .consume()
4648 }
4649
4650 #[tokio::test]
4653 async fn mount_resolves_on_the_target_front() {
4654 let producer = origin(1).produce();
4655 let server = producer.dynamic(".svc", Route::default()).unwrap();
4656 let project = mounted(&producer, "p1", &["**"]);
4657
4658 let through = project.request_broadcast(".svc/foo");
4659 let direct = producer.consume().request_broadcast(".svc/p1/foo");
4660
4661 let request = queued(&server).await;
4662 assert_eq!(request.path().as_str(), ".svc/p1/foo");
4663 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
4664 let source = broadcast::Info::new().produce();
4665 request.accept(&source);
4666
4667 let through = through.await.expect("resolves");
4668 let direct = direct.await.expect("resolves");
4669 assert!(through.is_clone(&direct));
4670 assert_eq!(through.info().path.as_str(), ".svc/foo");
4672 }
4673
4674 #[tokio::test]
4677 async fn mount_presents_target_routes_under_the_mount() {
4678 let producer = origin(1).produce();
4679 let _claim = producer.announce(".svc", Route::default()).unwrap();
4680 let foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4681 let _other = producer.publish(".svc/p2/bar", Route::default()).unwrap();
4682 let _cam = producer.publish("p1/cam", Route::default()).unwrap();
4683 let _shadowed = producer.publish("p1/.svc/forged", Route::default()).unwrap();
4684 let project = mounted(&producer, "p1", &["**"]);
4685
4686 let mut announced = project.announced();
4688 announced.assert_next_active("cam");
4689 announced.assert_next_wait();
4690
4691 let mut announced = project.clone().with_hidden(true).announced();
4692 announced.assert_next_active(".svc");
4693 announced.assert_next_active(".svc/foo");
4694 announced.assert_next_active("cam");
4695 announced.assert_next_wait();
4696
4697 let mut inside = project
4699 .scope(".svc", &Patterns::from(Pattern::all()))
4700 .unwrap()
4701 .announced();
4702 inside.assert_next_active("");
4703 inside.assert_next_active("foo");
4704 inside.assert_next_wait();
4705
4706 drop(foo);
4707 announced.assert_next_ended(".svc/foo");
4708 inside.assert_next_ended("foo");
4709 announced.assert_next_wait();
4710
4711 let err = project.request_broadcast(".svc/forged").await.err().unwrap();
4713 assert!(matches!(err, Error::NotFound | Error::Unroutable), "{err:?}");
4714 }
4715
4716 #[tokio::test]
4719 async fn mount_authorizes_the_named_path() {
4720 let producer = origin(1).produce();
4721 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4722 let _bar = producer.publish(".svc/p1/bar", Route::default()).unwrap();
4723
4724 let granted = mounted(&producer, "p1", &["foo", ".svc/foo"]);
4725 granted.request_broadcast(".svc/foo").await.expect("granted");
4726 let refused = granted
4727 .request_broadcast(".svc/bar")
4728 .now_or_never()
4729 .expect("refused at once");
4730 assert!(matches!(refused, Err(Error::Unauthorized)));
4731 let mut announced = granted.clone().with_hidden(true).announced();
4732 announced.assert_next_active(".svc/foo");
4733 announced.assert_next_wait();
4734
4735 let narrower = mounted(&producer, "p1", &["foo"]);
4736 let refused = narrower
4737 .request_broadcast(".svc/foo")
4738 .now_or_never()
4739 .expect("refused at once");
4740 assert!(matches!(refused, Err(Error::Unauthorized)));
4741 narrower.clone().with_hidden(true).announced().assert_next_wait();
4742
4743 let other = mounted(&producer, "p2", &["**"]);
4745 other.clone().with_hidden(true).announced().assert_next_wait();
4746 let err = other.request_broadcast(".svc/foo").await.err().unwrap();
4747 assert!(matches!(err, Error::Unroutable));
4748 }
4749
4750 #[tokio::test]
4753 async fn mount_captures_the_named_path() {
4754 let producer = origin(1).produce();
4755 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4756 let project = mounted(&producer, "p1", &["**"]).with_hidden(true);
4757
4758 let update = project.announced().try_next().expect("foo");
4759 assert_eq!(update.prefix.as_str(), ".svc/foo");
4760 assert_eq!(update.captures, Some(vec![".svc/foo".parse::<Pattern>().unwrap()]));
4761
4762 let inside = project.scope(".svc", &Patterns::from(Pattern::all())).unwrap();
4763 let update = inside.announced().try_next().expect("foo");
4764 assert_eq!(update.prefix.as_str(), "foo");
4765 assert_eq!(update.captures, Some(vec!["foo".parse::<Pattern>().unwrap()]));
4766 }
4767
4768 #[tokio::test]
4771 async fn mount_keeps_a_max_depth_target() {
4772 let producer = origin(1).produce();
4773 let deep = vec!["d"; Path::MAX_PARTS].join("/");
4774 let _leaf = producer.publish(deep.as_str(), Route::default()).unwrap();
4775 let project = producer
4776 .mount("p1/.svc", deep.as_str())
4777 .unwrap()
4778 .scope("p1", &Patterns::from(Pattern::all()))
4779 .unwrap()
4780 .consume()
4781 .with_hidden(true);
4782
4783 let mut announced = project.announced();
4784 announced.assert_next_active(".svc");
4785 announced.assert_next_wait();
4786
4787 let err = project.request_broadcast(".svc/x").await.err().unwrap();
4789 assert!(matches!(err, Error::BoundsExceeded(_)), "{err:?}");
4790 let inside = project.scope(".svc/x", &Patterns::from(Pattern::all())).unwrap();
4791 inside.announced().assert_next_wait();
4792 }
4793
4794 #[tokio::test]
4797 async fn mount_bounds_the_named_path() {
4798 let producer = origin(1).produce();
4799 let _near = producer.publish("t/x", Route::default()).unwrap();
4800 let deep = format!("t/{}", vec!["d"; Path::MAX_PARTS - 1].join("/"));
4801 let _deep = producer.publish(deep.as_str(), Route::default()).unwrap();
4802 let project = producer
4803 .mount("p1/a/b", "t")
4804 .unwrap()
4805 .scope("p1", &Patterns::from(Pattern::all()))
4806 .unwrap()
4807 .consume();
4808
4809 let mut announced = project.announced();
4810 announced.assert_next_active("a/b/x");
4811 announced.assert_next_wait();
4812
4813 let over = vec!["d"; Path::MAX_PARTS + 1].join("/");
4814 assert!(matches!(
4815 producer.mount(over.as_str(), "t"),
4816 Err(Error::BoundsExceeded(_))
4817 ));
4818 }
4819
4820 #[tokio::test]
4822 async fn mount_is_read_only() {
4823 let producer = origin(1).produce();
4824 let project = producer
4825 .mount("p1/.svc", ".svc/p1")
4826 .unwrap()
4827 .scope("p1", &Patterns::from(Pattern::all()))
4828 .unwrap();
4829 assert!(matches!(project.create_broadcast(".svc/foo"), Err(Error::Unauthorized)));
4830 assert!(matches!(
4831 project.dynamic(".svc", Route::default()),
4832 Err(Error::Unauthorized)
4833 ));
4834 assert!(matches!(
4835 project.dynamic(".svc/foo", Route::default()),
4836 Err(Error::Unauthorized)
4837 ));
4838 project.publish("cam", Route::default()).unwrap();
4839 }
4840
4841 #[test]
4843 fn mount_never_widens_a_scope() {
4844 let (producer, _driver) = Producer::new(Config::new(origin(1)));
4845 let project = producer.scope("", &scopes(&["p1"])).unwrap();
4846 assert!(matches!(project.mount("p1/.svc", ".svc/p1"), Err(Error::Unauthorized)));
4847
4848 let mounted = producer.mount("p1/.svc", ".svc/p1").unwrap();
4849 assert!(matches!(mounted.mount("p1/.svc/x", ".other"), Err(Error::Duplicate)));
4850 assert!(matches!(mounted.mount("p1", ".other"), Err(Error::Duplicate)));
4851 assert!(matches!(mounted.mount("p2", "p1/.svc/x"), Err(Error::Duplicate)));
4853 assert!(matches!(mounted.mount("p2", "p1"), Err(Error::Duplicate)));
4854 assert!(matches!(mounted.mount(".svc/p1/x", ".other"), Err(Error::Duplicate)));
4856 assert!(matches!(mounted.mount(".svc", ".other"), Err(Error::Duplicate)));
4857 mounted.mount("p1/.other", ".other/p1").unwrap();
4858 mounted.mount("p2/.svc", ".svc/p1").unwrap();
4860 for (at, target) in [("a", "a/b"), ("a/b", "a"), ("a", "a")] {
4862 assert!(
4863 matches!(producer.mount(at, target), Err(Error::Duplicate)),
4864 "{at} -> {target}"
4865 );
4866 }
4867 }
4868
4869 #[test]
4871 fn mount_refuses_a_wildcard_mount_point() {
4872 let (producer, _driver) = Producer::new(Config::new(origin(1)));
4873 for at in ["*", "p1/*", "p1/**", "p1/a*"] {
4874 assert!(
4875 matches!(producer.mount(at, ".svc/p1"), Err(Error::InvalidPath(_))),
4876 "{at}"
4877 );
4878 }
4879 }
4880
4881 #[tokio::test]
4884 async fn mount_egress_counts_under_the_named_path() {
4885 let registry = stats::Registry::new(stats::Config::new());
4886 let producer = origin(1).produce();
4887 let _foo = producer.publish(".svc/p1/foo", Route::default()).unwrap();
4888 let project = mounted(&producer, "p1", &["**"])
4889 .with_stats(registry.tier(stats::Tier::default()).session("p1"))
4890 .with_hidden(true);
4891
4892 let mut announced = project.announced();
4893 announced.assert_next_active(".svc/foo");
4894 project.request_broadcast(".svc/foo").await.expect("resolves");
4895
4896 let mut report = stats::Report::default();
4897 registry.report(&mut report);
4898 let paths: Vec<_> = report
4899 .traffic
4900 .iter()
4901 .map(|entry| entry.path.as_str().to_string())
4902 .collect();
4903 assert_eq!(paths, ["p1/.svc/foo"]);
4904 }
4905
4906 #[tokio::test]
4908 async fn hidden_route_announced_later_stays_hidden() {
4909 let producer = origin(1).produce();
4910 let consumer = producer.consume();
4911 let mut announced = consumer.announced();
4912 let mut opted = consumer.clone().with_hidden(true).announced();
4913
4914 let hidden = producer.announce(".stats/node", Route::default()).unwrap();
4915 announced.assert_next_wait();
4916 opted.assert_next_active(".stats/node");
4917
4918 drop(hidden);
4919 announced.assert_next_wait();
4920 opted.assert_next_ended(".stats/node");
4921 }
4922
4923 #[tokio::test]
4924 async fn broadcast_announces_its_own_path() {
4925 let producer = origin(1).produce();
4926 let consumer = producer.consume();
4927 let mut announced = consumer.announced();
4928 let mut peer = consumer.clone().excluding(Hop::UNKNOWN).announced();
4929
4930 let broadcast = producer.create_broadcast("room/alice").unwrap();
4932 announced.assert_next_wait();
4933 peer.assert_next_wait();
4934
4935 broadcast.announce(Route::default().with_cost(3)).unwrap();
4936 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(3));
4937 assert_eq!(peer.assert_next_active("room/alice").cost, Cost::new(3));
4938
4939 broadcast.announce(Route::default().with_cost(1)).unwrap();
4941 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(1));
4942 assert_eq!(peer.assert_next_active("room/alice").cost, Cost::new(1));
4943
4944 broadcast.unannounce();
4946 announced.assert_next_ended("room/alice");
4947 peer.assert_next_ended("room/alice");
4948 broadcast.unannounce();
4949 announced.assert_next_wait();
4950 let err = consumer.request_broadcast("room/alice").await.err().unwrap();
4951 assert!(matches!(err, Error::Unroutable));
4952
4953 broadcast.announce(Route::default()).unwrap();
4955 announced.assert_next_active("room/alice");
4956 peer.assert_next_active("room/alice");
4957 broadcast.close();
4958 announced.assert_next_ended("room/alice");
4959 peer.assert_next_ended("room/alice");
4960 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
4961 announced.assert_next_wait();
4962 }
4963
4964 #[tokio::test]
4965 async fn broadcast_announcement_retracts_with_the_last_producer() {
4966 let producer = origin(1).produce();
4967 let consumer = producer.consume();
4968 let mut announced = consumer.announced();
4969
4970 let broadcast = producer.create_broadcast("room/alice").unwrap();
4971 let clone = broadcast.clone();
4972 broadcast.announce(Route::default()).unwrap();
4973 announced.assert_next_active("room/alice");
4974
4975 drop(broadcast);
4977 announced.assert_next_wait();
4978 drop(clone);
4979 announced.assert_next_ended("room/alice");
4980 }
4981
4982 #[tokio::test]
4983 async fn publish_creates_and_announces_together() {
4984 let producer = origin(1).produce();
4985 let mut announced = producer.consume().announced();
4986 let _broadcast = producer.publish("room/alice", Route::default()).unwrap();
4987 announced.assert_next_active("room/alice");
4988 }
4989
4990 #[tokio::test]
4991 async fn standalone_broadcast_cannot_announce() {
4992 let broadcast = broadcast::Info::new().produce();
4993 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
4994 broadcast.unannounce();
4996 }
4997
4998 #[tokio::test]
4999 async fn announce_replays_to_late_cursor() {
5000 let producer = origin(1).produce();
5001 let _a = producer.announce("room/alice", Route::default()).unwrap();
5002 let _b = producer.announce("room/bob", Route::default()).unwrap();
5003
5004 let mut announced = producer.consume().announced();
5005 announced.assert_next_active("room/alice");
5007 announced.assert_next_active("room/bob");
5008 announced.assert_next_wait();
5009 }
5010
5011 #[tokio::test]
5012 async fn announce_keeps_its_prefix_under_a_producer_scope() {
5013 let producer = origin(1).produce();
5014 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
5015
5016 let _a = scoped.announce("", Route::default()).unwrap();
5018 let mut announced = producer.consume().announced();
5019 announced.assert_next_active("");
5020
5021 assert!(matches!(
5023 scoped.announce("other", Route::default()),
5024 Err(Error::Unauthorized)
5025 ));
5026 }
5027
5028 #[tokio::test]
5029 async fn cursor_keeps_an_overlapping_prefix_above_its_scope() {
5030 let producer = origin(1).produce();
5031 let _a = producer.announce("", Route::default()).unwrap();
5032
5033 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
5034 let mut announced = consumer.announced();
5035 announced.assert_next_active("");
5036 }
5037
5038 #[tokio::test]
5039 async fn cursor_root_strips_prefix() {
5040 let producer = origin(1).produce();
5041 let _a = producer.announce("room/alice", Route::default()).unwrap();
5042
5043 let consumer = producer
5044 .consume()
5045 .scope("room", &Patterns::from(Pattern::all()))
5046 .unwrap();
5047 let mut announced = consumer.announced();
5048 announced.assert_next_active("alice");
5049 }
5050
5051 #[tokio::test]
5052 async fn best_route_wins_and_fails_over() {
5053 let producer = origin(1).produce();
5054 let mut announced = producer.consume().announced();
5055
5056 let expensive = producer
5057 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5058 .unwrap();
5059 let route = announced.assert_next_active("room");
5060 assert_eq!(route.cost, Cost::new(5));
5061
5062 let cheap = producer
5064 .announce("room", Route::default().with_hops(hops(&[20])).with_cost(1))
5065 .unwrap();
5066 let route = announced.assert_next_active("room");
5067 assert_eq!(route.cost, Cost::new(1));
5068
5069 drop(cheap);
5071 let route = announced.assert_next_active("room");
5072 assert_eq!(route.cost, Cost::new(5));
5073
5074 drop(expensive);
5076 announced.assert_next_ended("room");
5077 }
5078
5079 #[tokio::test]
5080 async fn identical_reannounce_is_invisible() {
5081 let producer = origin(1).produce();
5082 let mut announced = producer.consume().announced();
5083
5084 let old = producer
5085 .announce("room", Route::default().with_hops(hops(&[10])))
5086 .unwrap();
5087 let first = announced.assert_next_active("room");
5088 assert_eq!(first.hops.as_slice(), hops(&[10]).as_slice());
5089
5090 let _new = producer
5094 .announce("room", Route::default().with_hops(hops(&[10])))
5095 .unwrap();
5096 announced.assert_next_wait();
5097
5098 drop(old);
5100 announced.assert_next_wait();
5101 }
5102
5103 #[tokio::test]
5104 async fn exclude_hides_routes_through_the_peer() {
5105 let producer = origin(1).produce();
5106 let _a = producer
5107 .announce("room", Route::default().with_hops(hops(&[7])))
5108 .unwrap();
5109
5110 let mut hidden = producer.consume().excluding(origin(7)).announced();
5111 hidden.assert_next_wait();
5112
5113 let mut visible = producer.consume().excluding(origin(8)).announced();
5114 visible.assert_next_active("room");
5115 }
5116
5117 #[tokio::test]
5118 async fn exclude_matches_via_when_the_chain_is_anonymous() {
5119 let producer = origin(1).produce();
5120 let assigned = origin(777);
5121 let _echoed = producer
5122 .announce("echoed", Route::default().with_hops(hops(&[0])).with_via(assigned))
5123 .unwrap();
5124 let _local = producer
5125 .announce("local", Route::default().with_hops(hops(&[10])))
5126 .unwrap();
5127
5128 let mut hidden = producer.consume().excluding(assigned).announced();
5129 hidden.assert_next_active("local");
5130 hidden.assert_next_wait();
5131 }
5132
5133 #[tokio::test]
5134 async fn anonymous_route_loses_to_identified_at_any_cost() {
5135 let producer = origin(1).produce();
5136 let mut announced = producer.consume().announced();
5137
5138 let _anonymous = producer
5139 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(1))
5140 .unwrap();
5141 let route = announced.assert_next_active("room");
5142 assert!(route.is_anonymous());
5143 assert_eq!(route.cost, Cost::new(1));
5144
5145 let _identified = producer
5146 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5147 .unwrap();
5148 let route = announced.assert_next_active("room");
5149 assert!(!route.is_anonymous());
5150 assert_eq!(route.cost, Cost::new(5));
5151 }
5152
5153 #[tokio::test]
5156 async fn a_stamped_route_ranks_below_an_identified_one() {
5157 let producer = origin(1).produce();
5158 let mut announced = producer.consume().announced();
5159
5160 let mut stamped = Hops::new();
5161 stamped.stamp(origin(5)).unwrap();
5162 assert_eq!(stamped.as_slice(), &[origin(5), Hop::UNKNOWN]);
5163
5164 let _legacy = producer
5165 .announce("room", Route::default().with_hops(stamped).with_cost(0))
5166 .unwrap();
5167 let route = announced.assert_next_active("room");
5168 assert!(route.is_anonymous());
5169
5170 let _identified = producer
5171 .announce("room", Route::default().with_hops(hops(&[10, 11])).with_cost(5))
5172 .unwrap();
5173 let route = announced.assert_next_active("room");
5174 assert_eq!(route.hops.as_slice(), hops(&[10, 11]).as_slice());
5175 }
5176
5177 #[test]
5178 fn stamping_keeps_the_leading_zero_and_names_nothing_else() {
5179 let mut chain = hops(&[0, 7]);
5180 chain.stamp(origin(5)).unwrap();
5181 assert_eq!(chain.as_slice(), hops(&[5, 0, 7]).as_slice());
5182
5183 let mut named = hops(&[7, 0]);
5185 named.stamp(origin(5)).unwrap();
5186 assert_eq!(named.as_slice(), hops(&[7, 0]).as_slice());
5187
5188 let mut full = Hops::try_from(vec![Hop::UNKNOWN; MAX_HOPS]).unwrap();
5190 assert_eq!(full.stamp(origin(5)), Err(InvalidHop::TooMany));
5191 }
5192
5193 #[tokio::test]
5194 async fn anonymous_routes_order_by_cost() {
5195 let producer = origin(1).produce();
5196 let mut announced = producer.consume().announced();
5197
5198 let expensive = producer
5199 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(5))
5200 .unwrap();
5201 let route = announced.assert_next_active("room");
5202 assert_eq!(route.cost, Cost::new(5));
5203
5204 let _cheap = producer
5205 .announce("room", Route::default().with_hops(hops(&[0, 7])).with_cost(1))
5206 .unwrap();
5207 let route = announced.assert_next_active("room");
5208 assert!(route.is_anonymous());
5209 assert_eq!(route.cost, Cost::new(1));
5210
5211 drop(expensive);
5212 announced.assert_next_wait();
5213 }
5214
5215 #[tokio::test]
5216 async fn anonymous_chain_from_identified_peer_still_ranks_last() {
5217 let producer = origin(1).produce();
5218 let mut announced = producer.consume().announced();
5219
5220 let _anonymous = producer
5221 .announce(
5222 "room",
5223 Route::default()
5224 .with_hops(hops(&[0, 7]))
5225 .with_cost(1)
5226 .with_via(origin(7)),
5227 )
5228 .unwrap();
5229 announced.assert_next_active("room");
5230
5231 let _identified = producer
5232 .announce("room", Route::default().with_hops(hops(&[10, 20])).with_cost(5))
5233 .unwrap();
5234 let route = announced.assert_next_active("room");
5235 assert!(!route.is_anonymous());
5236 assert_eq!(route.cost, Cost::new(5));
5237 }
5238
5239 #[tokio::test]
5240 async fn request_prefers_identified_over_cheaper_anonymous() {
5241 let producer = origin(1).produce();
5242 let consumer = producer.consume();
5243
5244 let anonymous = producer
5245 .dynamic("room", Route::default().with_hops(hops(&[0])).with_cost(1))
5246 .unwrap();
5247 let identified = producer
5248 .dynamic("room", Route::default().with_hops(hops(&[10])).with_cost(5))
5249 .unwrap();
5250
5251 let _pending = consumer.request_broadcast("room/alice");
5252 let request = queued(&identified).await;
5253 assert_eq!(request.path().as_str(), "room/alice");
5254 assert!(
5255 anonymous.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
5256 "the cheaper anonymous route must not serve"
5257 );
5258 }
5259
5260 #[tokio::test]
5261 async fn update_reprices_in_place() {
5262 let producer = origin(1).produce();
5263 let mut announced = producer.consume().announced();
5264
5265 let announcement = producer.announce("room", Route::default()).unwrap();
5266 announced.assert_next_active("room");
5267
5268 announcement.update(Route::default().with_cost(9)).unwrap();
5269 let route = announced.assert_next_active("room");
5270 assert_eq!(route.cost, Cost::new(9));
5271 }
5272
5273 #[tokio::test]
5274 async fn retract_after_undelivered_reprice_still_delivered() {
5275 let producer = origin(1).produce();
5276 let mut announced = producer.consume().announced();
5277
5278 let announcement = producer.announce("room", Route::default()).unwrap();
5279 announced.assert_next_active("room");
5280
5281 announcement.update(Route::default().with_cost(9)).unwrap();
5285 drop(announcement);
5286 announced.assert_next_ended("room");
5287 announced.assert_next_wait();
5288 }
5289
5290 #[tokio::test]
5291 async fn scoped_cursor_advertises_most_specific_covering_route() {
5292 let producer = origin(1).produce();
5293 let _broad = producer.announce("room", Route::default().with_cost(1)).unwrap();
5296 let _narrow = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
5297
5298 let consumer = producer
5299 .consume()
5300 .scope("room/alice", &Patterns::from(Pattern::all()))
5301 .unwrap();
5302 let mut announced = consumer.announced();
5303 let route = announced.assert_next_active("");
5304 assert_eq!(route.cost, Cost::new(9));
5305 announced.assert_next_wait();
5306 }
5307
5308 #[tokio::test]
5309 async fn capture_change_retracts_before_reannouncing_a_presented_prefix() {
5310 let producer = origin(1).produce();
5311 let _broad = producer.announce("room", Route::default()).unwrap();
5312 let exact = producer.announce("room/alice", Route::default()).unwrap();
5313 let consumer = producer
5314 .consume()
5315 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
5316 .unwrap()
5317 .scope("room/alice", &Patterns::from(Pattern::all()))
5318 .unwrap();
5319 let mut announced = consumer.announced();
5320
5321 let first = announced.next().now_or_never().expect("next").expect("announce");
5322 assert_eq!(first.prefix.as_str(), "");
5323 assert_eq!(first.kind, AnnounceKind::Announced);
5324 assert_eq!(first.captures, Some(Vec::new()));
5325
5326 drop(exact);
5327 let retracted = announced.next().now_or_never().expect("next").expect("retract");
5328 assert_eq!(retracted.prefix.as_str(), "");
5329 assert_eq!(retracted.kind, AnnounceKind::Retracted);
5330 assert_eq!(retracted.captures, Some(Vec::new()));
5331 let replacement = announced.next().now_or_never().expect("next").expect("announce");
5332 assert_eq!(replacement.prefix.as_str(), "");
5333 assert_eq!(replacement.kind, AnnounceKind::Announced);
5334 assert_eq!(replacement.captures, None);
5335 }
5336
5337 #[tokio::test]
5338 async fn routed_broadcast_resolves_once_announced() {
5339 let producer = origin(1).produce();
5340 let consumer = producer.consume();
5341
5342 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
5344 assert!((&mut resolving).now_or_never().is_none());
5345
5346 let broadcast = producer.create_broadcast("room/alice").unwrap();
5348 for _ in 0..20 {
5349 tokio::task::yield_now().await;
5350 }
5351 assert!((&mut resolving).now_or_never().is_none());
5352
5353 broadcast.announce(Route::default()).unwrap();
5354 let resolved = resolving.await.expect("resolves once announced");
5355 assert_eq!(resolved.info().path.as_str(), "room/alice");
5356 drop(broadcast);
5357 }
5358
5359 #[tokio::test]
5362 async fn cheaper_remote_route_beats_a_local_broadcast() {
5363 let producer = origin(1).produce();
5364 let consumer = producer.consume();
5365 let mut announced = consumer.announced();
5366
5367 let _local = producer.publish("room/alice", Route::default().with_cost(5)).unwrap();
5368 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::new(5));
5369
5370 let server = producer
5371 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(1))
5372 .unwrap();
5373 let route = announced.assert_next_active("room/alice");
5374 assert_eq!(route.cost, Cost::new(1));
5375 assert_eq!(route.hops, hops(&[10]));
5376
5377 let pending = consumer.request_broadcast("room/alice");
5379 let request = queued(&server).await;
5380 let upstream = broadcast::Info::new().produce();
5381 request.accept(&upstream);
5382 pending.await.expect("resolves through the cheaper route");
5383 }
5384
5385 #[tokio::test]
5389 async fn cheaper_route_after_a_front_wins_new_requests() {
5390 let producer = origin(1).produce();
5391 let consumer = producer.consume();
5392
5393 let _local = producer.publish("room/alice", Route::default().with_cost(5)).unwrap();
5394 let first = consumer
5395 .request_broadcast("room/alice")
5396 .await
5397 .expect("resolves locally");
5398
5399 let server = producer
5400 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(1))
5401 .unwrap();
5402
5403 let pending = consumer.request_broadcast("room/alice");
5404 let request = queued(&server).await;
5405 let upstream = broadcast::Info::new().produce();
5406 request.accept(&upstream);
5407 let second = pending.await.expect("resolves through the cheaper route");
5408
5409 assert!(!first.is_closed(), "the old front must keep serving its readers");
5410 assert!(!first.is_clone(&second), "the newcomer must not join the old front");
5411 }
5412
5413 #[tokio::test]
5417 async fn local_broadcast_wins_a_tie_with_a_hopless_route() {
5418 let producer = origin(1).produce();
5419 let consumer = producer.consume();
5420
5421 let _local = producer.publish("room/alice", Route::default()).unwrap();
5422 let server = producer.dynamic("room/alice", Route::default()).unwrap();
5423
5424 let resolved = tokio::time::timeout(Duration::from_secs(1), consumer.request_broadcast("room/alice"))
5426 .await
5427 .expect("the newer hopless route won the tie")
5428 .expect("resolves");
5429 assert_eq!(resolved.info().path.as_str(), "room/alice");
5430 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
5431 }
5432
5433 #[tokio::test]
5436 async fn announce_stats_follow_the_advertisement() {
5437 let registry = stats::Registry::new(stats::Config::new());
5438 let producer = origin(1)
5439 .produce()
5440 .with_stats(registry.tier(stats::Tier::default()).session("root"));
5441 let announces = || {
5442 registry
5443 .snapshot()
5444 .traffic()
5445 .into_iter()
5446 .find(|(_, role, _)| *role == stats::Role::Subscriber)
5447 .map(|(_, _, traffic)| (traffic.announces_started, traffic.announces_ended))
5448 .unwrap_or_default()
5449 };
5450
5451 let broadcast = producer.create_broadcast("room/alice").unwrap();
5452 assert_eq!(announces(), (0, 0), "a hidden broadcast is not announced");
5453 broadcast.announce(Route::default()).unwrap();
5454 broadcast
5455 .announce(Route {
5456 cost: Cost::new(3),
5457 ..Route::default()
5458 })
5459 .unwrap();
5460 assert_eq!(announces(), (1, 0), "a re-price is not another announce");
5461 broadcast.unannounce();
5462 assert_eq!(announces(), (1, 1));
5463 broadcast.announce(Route::default()).unwrap();
5464 drop(broadcast);
5465 assert_eq!(announces(), (2, 2));
5466 }
5467
5468 #[tokio::test]
5470 async fn local_broadcast_wins_a_cost_tie() {
5471 let producer = origin(1).produce();
5472 let consumer = producer.consume();
5473 let mut announced = consumer.announced();
5474
5475 let server = producer
5476 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(2))
5477 .unwrap();
5478 announced.assert_next_active("room/alice");
5479 let _local = producer.publish("room/alice", Route::default().with_cost(2)).unwrap();
5480 assert!(announced.assert_next_active("room/alice").hops.is_empty());
5481
5482 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5483 assert_eq!(resolved.info().path.as_str(), "room/alice");
5484 for _ in 0..20 {
5485 tokio::task::yield_now().await;
5486 }
5487 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
5488 }
5489
5490 #[tokio::test]
5493 async fn unannounce_ends_the_front() {
5494 let producer = origin(1).produce();
5495 let consumer = producer.consume();
5496
5497 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5498 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5499
5500 broadcast.unannounce();
5501 let err = consumer.request_broadcast("room/alice").await.err().unwrap();
5502 assert!(matches!(err, Error::Unroutable), "joined a retracted front: {err}");
5503 settle(|| resolved.is_closed()).await;
5504 assert!(!broadcast.consume().is_closed(), "the broadcast itself lives on");
5505
5506 broadcast.announce(Route::default()).unwrap();
5508 let again = consumer.request_broadcast("room/alice").await.expect("resolves again");
5509 assert!(!again.is_clone(&resolved));
5510 }
5511
5512 #[tokio::test]
5516 async fn unannounce_keeps_a_track_awaiting_its_info() {
5517 let producer = origin(1).produce();
5518 let consumer = producer.consume();
5519
5520 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5521 let mut dynamic = broadcast.dynamic();
5522 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5523 let track = resolved.track("video").unwrap();
5524 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5525 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5526 .await
5527 .expect("the front asked the source")
5528 .expect("request");
5529
5530 broadcast.unannounce();
5531 settle(|| resolved.is_closed()).await;
5532
5533 let source = request.accept(None);
5534 let mut group = source.append_group().unwrap();
5535 group.write_frame(crate::Timestamp::ZERO, b"late".as_ref()).unwrap();
5536 group.finish().unwrap();
5537 source.finish().unwrap();
5538
5539 let mut subscription = subscribing.await.unwrap().expect("subscribe survives the retraction");
5540 let mut group = subscription.recv_group().await.unwrap().expect("the source's group");
5541 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"late");
5542 assert!(matches!(subscription.recv_group().await, Ok(None)), "ends cleanly");
5543 }
5544
5545 async fn served_front() -> (Dynamic, broadcast::Producer, broadcast::Dynamic, broadcast::Consumer) {
5548 let producer = origin(1).produce();
5549 let consumer = producer.consume();
5550 let server = producer
5551 .dynamic("room/alice", Route::default().with_hops(hops(&[10])))
5552 .unwrap();
5553 let pending = consumer.request_broadcast("room/alice");
5554 let upstream = broadcast::Info::new().produce();
5555 let dynamic = upstream.dynamic();
5556 queued(&server).await.accept(&upstream);
5557 let resolved = pending.await.expect("resolves");
5558 (server, upstream, dynamic, resolved)
5559 }
5560
5561 #[tokio::test]
5566 async fn returning_reader_skips_a_warm_cache_the_copy_resolved_past() {
5567 let ms = |v: u64| crate::Timestamp::from_millis(v).unwrap();
5568 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5569 let budget = track::Subscription::default().with_max_age(Duration::from_millis(100));
5570
5571 let track = resolved.track("audio").unwrap();
5572 let b = budget.clone();
5573 let subscribing = tokio::spawn(async move { track.subscribe(b).await });
5574 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5575 .await
5576 .expect("the front asked the source")
5577 .expect("request");
5578 let source = request.resolving_start().accept(None);
5579 for seq in 0..4u64 {
5580 let mut group = source.create_group(seq.into()).unwrap();
5581 group.write_frame(ms(seq * 20), b"old".as_ref()).unwrap();
5582 group.finish().unwrap();
5583 }
5584 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5585 subscription.recv_group().await.unwrap().expect("the live group");
5586 drop(subscription);
5587 tokio::time::timeout(Duration::from_secs(1), source.unused())
5588 .await
5589 .expect("parked")
5590 .expect("source open");
5591 drop(source);
5592
5593 let track = resolved.track("audio").unwrap();
5594 let subscribing = tokio::spawn(async move { track.subscribe(budget).await });
5595 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5596 .await
5597 .expect("the front asked the source again")
5598 .expect("request");
5599 let mut source = request.resolving_start().accept(None);
5600 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
5601
5602 assert!(
5604 tokio::time::timeout(Duration::from_millis(50), subscription.recv_group())
5605 .await
5606 .is_err(),
5607 "the warm cache was served before the copy resolved its start"
5608 );
5609
5610 source.start_at(20).unwrap();
5612 let mut group = source.create_group(20u64.into()).unwrap();
5613 group.write_frame(ms(2000), b"new".as_ref()).unwrap();
5614 group.finish().unwrap();
5615 let group = subscription.recv_group().await.unwrap().expect("the live group");
5616 assert_eq!(group.sequence, 20, "a stale warm group was served");
5617 }
5618
5619 #[tokio::test]
5624 async fn returning_reader_replays_a_current_warm_cache() {
5625 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5626
5627 let track = resolved.track("catalog").unwrap();
5628 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5629 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5630 .await
5631 .expect("the front asked the source")
5632 .expect("request");
5633 let source = request.resolving_start().accept(None);
5634 let mut group = source.create_group(0u64.into()).unwrap();
5635 group.write_frame(crate::Timestamp::ZERO, b"snapshot".as_ref()).unwrap();
5636 group.finish().unwrap();
5637 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5638 subscription.recv_group().await.unwrap().expect("the catalog");
5639 drop(subscription);
5640 tokio::time::timeout(Duration::from_secs(1), source.unused())
5641 .await
5642 .expect("parked")
5643 .expect("source open");
5644 drop(source);
5645
5646 let track = resolved.track("catalog").unwrap();
5647 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5648 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5649 .await
5650 .expect("the front asked the source again")
5651 .expect("request");
5652 let mut source = request.resolving_start().accept(None);
5653 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
5654
5655 let reading = tokio::spawn(async move {
5658 let mut group = subscription.recv_group().await.unwrap().expect("the catalog");
5659 assert_eq!(group.sequence, 0);
5660 group.read_frame().await.unwrap().expect("the snapshot").payload
5661 });
5662 tokio::task::yield_now().await;
5663 assert_eq!(
5664 source.subscription().and_then(|sub| sub.start),
5665 Some(track::Position { group: 0, frame: 1 }),
5666 "the re-splice asked past the cached catalog"
5667 );
5668 source.start_at(0).unwrap();
5669 let mut tail = source.create_group(0u64.into()).unwrap();
5670 tail.start_at(1).unwrap();
5671 tail.finish().unwrap();
5672 let payload = tokio::time::timeout(Duration::from_secs(1), reading)
5673 .await
5674 .expect("the returning reader never got the catalog")
5675 .unwrap();
5676 assert_eq!(&payload[..], b"snapshot");
5677 }
5678
5679 #[tokio::test(start_paused = true)]
5683 async fn finished_track_is_forgotten_after_the_linger() {
5684 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5685
5686 let track = resolved.track("catalog").unwrap();
5687 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5688 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5689 .await
5690 .expect("the front asked the source")
5691 .expect("request");
5692 let source = request.resolving_start().accept(None);
5693 let mut group = source.create_group(0u64.into()).unwrap();
5694 group.write_frame(crate::Timestamp::ZERO, b"snapshot".as_ref()).unwrap();
5695 group.finish().unwrap();
5696 source.finish().unwrap();
5697 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5698 assert_eq!(
5699 next_group(&mut subscription)
5700 .await
5701 .unwrap()
5702 .expect("the catalog")
5703 .sequence,
5704 0
5705 );
5706 assert!(next_group(&mut subscription).await.unwrap().is_none());
5707 drop(subscription);
5708 drop(source);
5709
5710 let mut subscription = resolved.track("catalog").unwrap().subscribe(None).await.unwrap();
5712 assert_eq!(
5713 next_group(&mut subscription)
5714 .await
5715 .unwrap()
5716 .expect("the catalog")
5717 .sequence,
5718 0
5719 );
5720 assert!(next_group(&mut subscription).await.unwrap().is_none());
5721 drop(subscription);
5722 assert!(
5723 tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5724 .await
5725 .is_err(),
5726 "a finished track within the linger asked the source again"
5727 );
5728
5729 tokio::time::sleep(TRACK_IDLE_LINGER).await;
5731
5732 let track = resolved.track("catalog").unwrap();
5733 let _subscribing = tokio::spawn(async move { track.subscribe(None).await });
5734 tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5735 .await
5736 .expect("the finished track outlived the linger")
5737 .expect("request");
5738 }
5739
5740 #[tokio::test]
5745 async fn returning_reader_continues_an_open_warm_group() {
5746 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
5747
5748 async fn read(group: &mut group::Consumer) -> Vec<u8> {
5749 let frame = tokio::time::timeout(Duration::from_secs(1), group.read_frame())
5750 .await
5751 .expect("frame")
5752 .unwrap()
5753 .expect("group ended");
5754 frame.payload.to_vec()
5755 }
5756
5757 let mut expect: Vec<&[u8]> = Vec::new();
5758 let mut floor: Option<track::Position> = None;
5759 for (round, payload) in [b"a".as_ref(), b"b", b"c"].into_iter().enumerate() {
5760 let track = resolved.track("log").unwrap();
5761 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5762 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5763 .await
5764 .expect("the front asked the source")
5765 .expect("request");
5766 let mut source = request.resolving_start().accept(None);
5767 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5768
5769 source.start_at(0).unwrap();
5771 let mut group = source.create_group(0u64.into()).unwrap();
5772 if let Some(floor) = floor {
5773 group.start_at(floor.frame).unwrap();
5774 }
5775 group.write_frame(crate::Timestamp::ZERO, payload).unwrap();
5776 expect.push(payload);
5777 source
5778 .insert_datagram(10, crate::Timestamp::ZERO, b"datagram".as_ref())
5779 .unwrap();
5780
5781 let mut reading = tokio::time::timeout(Duration::from_secs(1), subscription.recv_group())
5782 .await
5783 .expect("group 0")
5784 .unwrap()
5785 .expect("track ended");
5786 assert_eq!(reading.sequence, 0);
5787 for frame in &expect {
5788 assert_eq!(read(&mut reading).await, *frame, "round {round}");
5789 }
5790 assert_eq!(
5791 source.subscription().and_then(|sub| sub.start),
5792 floor,
5793 "round {round} asked for the wrong continuation"
5794 );
5795
5796 drop(reading);
5797 drop(subscription);
5798 tokio::time::timeout(Duration::from_secs(1), source.unused())
5799 .await
5800 .expect("parked")
5801 .expect("source open");
5802 drop(group);
5803 drop(source);
5804 floor = Some(track::Position {
5805 group: 0,
5806 frame: expect.len() as u64,
5807 });
5808 }
5809 }
5810
5811 #[tokio::test]
5812 async fn warm_head_survives_another_takeover_before_park() {
5813 tokio::time::pause();
5814 let producer = origin(1).produce();
5815 let _server = producer
5816 .dynamic("room/alice", Route::default().with_hops(hops(&[10])).with_cost(5))
5817 .unwrap();
5818 let pending = producer.consume().request_broadcast("room/alice");
5819 let upstream = broadcast::Info::new().produce();
5820 let mut dynamic = upstream.dynamic();
5821 queued(&_server).await.accept(&upstream);
5822 let resolved = pending.await.unwrap();
5823
5824 let track = resolved.track("log").unwrap();
5825 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5826 let source = dynamic.requested_track().await.unwrap().accept(None);
5827 let mut group = source.create_group(0u64.into()).unwrap();
5828 group.write_frame(crate::Timestamp::ZERO, b"a".as_ref()).unwrap();
5829 group.write_frame(crate::Timestamp::ZERO, b"b".as_ref()).unwrap();
5830 let mut subscription = subscribing.await.unwrap().unwrap();
5831 let mut reading = subscription.recv_group().await.unwrap().unwrap();
5832 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"a");
5833 drop(reading);
5834 drop(subscription);
5835 source.unused().await.unwrap();
5836 drop(group);
5837 drop(source);
5838
5839 let track = resolved.track("log").unwrap();
5840 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5841 let mut resumed = dynamic.requested_track().await.unwrap().resolving_start().accept(None);
5842 resumed.start_at(0).unwrap();
5843 let mut subscription = subscribing.await.unwrap().unwrap();
5844 let mut reading = subscription.recv_group().await.unwrap().unwrap();
5845 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"a");
5846
5847 let replacement_server = producer
5848 .dynamic("room/alice", Route::default().with_hops(hops(&[10])))
5849 .unwrap();
5850 let replacement = broadcast::Info::new().produce();
5851 let mut replacement_dynamic = replacement.dynamic();
5852 queued(&replacement_server).await.accept(&replacement);
5853 let mut source = replacement_dynamic
5854 .requested_track()
5855 .await
5856 .unwrap()
5857 .resolving_start()
5858 .accept(None);
5859 source.start_at(0).unwrap();
5860 let mut group = source.create_group(0u64.into()).unwrap();
5861 group.start_at(2).unwrap();
5862 group.write_frame(crate::Timestamp::ZERO, b"c".as_ref()).unwrap();
5863 let mut continuation = reading.clone();
5864 continuation.start_at(2);
5865 let frame = tokio::time::timeout(Duration::from_secs(1), continuation.read_frame())
5866 .await
5867 .expect("replacement frame")
5868 .unwrap()
5869 .unwrap();
5870 assert_eq!(&frame.payload[..], b"c");
5871 drop(continuation);
5872 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"b");
5873 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], b"c");
5874 drop(reading);
5875 drop(subscription);
5876 source.unused().await.unwrap();
5877 drop(group);
5878 drop(source);
5879
5880 let track = resolved.track("log").unwrap();
5881 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5882 let mut next = replacement_dynamic
5883 .requested_track()
5884 .await
5885 .unwrap()
5886 .resolving_start()
5887 .accept(None);
5888 next.start_at(0).unwrap();
5889 let mut subscription = subscribing.await.unwrap().unwrap();
5890 let mut reading = subscription.recv_group().await.unwrap().unwrap();
5891 for expected in [b"a", b"b", b"c"] {
5892 assert_eq!(&reading.read_frame().await.unwrap().unwrap().payload[..], expected);
5893 }
5894 }
5895
5896 #[tokio::test]
5899 async fn unannounce_keeps_a_returning_reader_awaiting_its_info() {
5900 let producer = origin(1).produce();
5901 let consumer = producer.consume();
5902
5903 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5904 let mut dynamic = broadcast.dynamic();
5905 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5906 let track = resolved.track("video").unwrap();
5907 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5908 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5909 .await
5910 .expect("the front asked the source")
5911 .expect("request");
5912 let source = request.accept(None);
5913 let mut group = source.append_group().unwrap();
5914 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
5915 group.finish().unwrap();
5916 let mut subscription = subscribing.await.unwrap().expect("subscribe");
5917 subscription.recv_group().await.unwrap().expect("the cached group");
5918 drop(subscription);
5919
5920 tokio::time::timeout(Duration::from_secs(1), source.unused())
5923 .await
5924 .expect("parked")
5925 .expect("source open");
5926 drop(source);
5927 let track = resolved.track("video").unwrap();
5928 let subscribing = tokio::spawn(async move { track.subscribe(None).await });
5929 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
5930 .await
5931 .expect("the front asked the source again")
5932 .expect("request");
5933
5934 broadcast.unannounce();
5935 settle(|| resolved.is_closed()).await;
5936
5937 let source = request.accept(None);
5939 let mut group = source.create_group(1u64.into()).unwrap();
5940 group.write_frame(crate::Timestamp::ZERO, b"late".as_ref()).unwrap();
5941 group.finish().unwrap();
5942 source.finish().unwrap();
5943
5944 let mut subscription = subscribing.await.unwrap().expect("subscribe survives the retraction");
5945 let mut payloads = Vec::new();
5946 while let Some(mut group) = subscription.recv_group().await.expect("ends cleanly") {
5947 payloads.push(group.read_frame().await.unwrap().unwrap().payload);
5948 }
5949 assert_eq!(payloads.last().map(|p| &p[..]), Some(&b"late"[..]));
5950 }
5951
5952 #[tokio::test]
5955 async fn reannounce_before_the_front_acts_keeps_it() {
5956 let producer = origin(1).produce();
5957 let consumer = producer.consume();
5958
5959 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
5960 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5961
5962 broadcast.unannounce();
5963 broadcast.announce(Route::default()).unwrap();
5964 for _ in 0..20 {
5965 tokio::task::yield_now().await;
5966 }
5967 assert!(!resolved.is_closed(), "the front ended across a reannouncement");
5968 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
5969 assert!(again.is_clone(&resolved));
5970 }
5971
5972 #[tokio::test]
5973 async fn local_broadcast_resolves_once_announced() {
5974 let producer = origin(1).produce();
5975 let consumer = producer.consume();
5976
5977 let broadcast = producer.create_broadcast("room/alice").unwrap();
5979 let err = consumer
5980 .request_broadcast("room/alice")
5981 .now_or_never()
5982 .expect("unroutable is synchronous")
5983 .err()
5984 .unwrap();
5985 assert!(matches!(err, Error::Unroutable));
5986
5987 broadcast.announce(Route::default()).unwrap();
5988 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5989 assert_eq!(resolved.info().path.as_str(), "room/alice");
5990 drop(broadcast);
5991
5992 let err = consumer
5994 .request_broadcast("room/bob")
5995 .now_or_never()
5996 .expect("unroutable is synchronous")
5997 .err()
5998 .unwrap();
5999 assert!(matches!(err, Error::Unroutable));
6000 }
6001
6002 #[test]
6003 fn create_broadcast_accepts_a_max_depth_path() {
6004 let producer = origin(1).produce();
6005 let path = vec!["a"; Path::MAX_PARTS].join("/");
6006 let _broadcast = producer.create_broadcast(path.as_str()).expect("max depth is allowed");
6007 let deeper = vec!["a"; Path::MAX_PARTS + 1].join("/");
6008 assert!(matches!(
6009 producer.create_broadcast(deeper.as_str()),
6010 Err(Error::BoundsExceeded(_))
6011 ));
6012 }
6013
6014 #[tokio::test]
6015 async fn duplicate_routes_aggregate_until_the_last_leaves() {
6016 let producer = origin(1).produce();
6017 let first = producer.dynamic("live", Route::default().with_cost(3)).unwrap();
6018 let second = producer.dynamic("live", Route::default().with_cost(1)).unwrap();
6019
6020 let mut announced = producer.consume().announced();
6021 let update = announced.next().now_or_never().expect("next").expect("no next");
6022 assert_eq!(update.prefix.as_str(), "live");
6023 assert_eq!(update.kind, AnnounceKind::Announced);
6024 assert_eq!(update.route.cost, Cost::new(1));
6025 announced.assert_next_wait();
6026
6027 drop(second);
6028 let update = announced.next().now_or_never().expect("next").expect("no next");
6029 assert_eq!(update.prefix.as_str(), "live");
6030 assert_eq!(update.kind, AnnounceKind::Updated);
6031 assert_eq!(update.route.cost, Cost::new(3));
6032
6033 drop(first);
6034 announced.assert_next_ended("live");
6035 announced.assert_next_wait();
6036 }
6037
6038 #[test]
6039 fn dynamic_may_cover_a_scope_but_disjoint_prefixes_are_refused() {
6040 let producer = origin(1).produce();
6041 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
6042 let _broad = scoped
6043 .dynamic("", Route::default())
6044 .expect("an overlapping prefix is accepted");
6045
6046 let _ok = scoped
6047 .dynamic("room/alice", Route::default())
6048 .expect("a contained prefix is accepted");
6049 assert!(matches!(
6050 scoped.dynamic("other", Route::default()),
6051 Err(Error::Unauthorized)
6052 ));
6053 }
6054
6055 #[tokio::test]
6056 async fn dynamic_route_keeps_its_producer_scope() {
6057 let producer = origin(1).produce();
6058 let scope = Patterns::from("*/chat".parse::<Pattern>().unwrap());
6059 let scoped = producer.scope("", &scope).unwrap();
6060 let dynamic = scoped.dynamic("", Route::default()).unwrap();
6061
6062 let mut matching = producer
6063 .consume()
6064 .scope("", &scopes(&["room/chat"]))
6065 .unwrap()
6066 .announced();
6067 matching.assert_next_active("");
6068 let mut outside = producer
6069 .consume()
6070 .scope("", &scopes(&["room/video"]))
6071 .unwrap()
6072 .announced();
6073 outside.assert_next_wait();
6074
6075 let refused = producer
6076 .consume()
6077 .request_broadcast("room/video")
6078 .now_or_never()
6079 .expect("an out-of-scope request must be refused synchronously");
6080 assert!(matches!(refused, Err(Error::Unroutable)));
6081 assert!(dynamic.requested_broadcast().now_or_never().is_none());
6082
6083 let _pending = producer.consume().request_broadcast("room/chat");
6084 let request = queued(&dynamic).await;
6085 assert_eq!(request.path().as_str(), "room/chat");
6086 }
6087
6088 #[tokio::test]
6089 async fn scoped_cursor_selects_among_the_routes_it_can_see() {
6090 let producer = origin(1).produce();
6091 let scoped = |pattern: &str| {
6092 producer
6093 .scope("", &Patterns::from(pattern.parse::<Pattern>().unwrap()))
6094 .unwrap()
6095 };
6096 let _chat = scoped("*/chat").dynamic("", Route::default().with_cost(1)).unwrap();
6097 let _video = scoped("*/video").dynamic("", Route::default().with_cost(5)).unwrap();
6098
6099 let mut video = producer
6100 .consume()
6101 .scope("", &scopes(&["room/video"]))
6102 .unwrap()
6103 .announced();
6104 assert_eq!(video.assert_next_active("").cost, Cost::new(5));
6105 }
6106
6107 #[tokio::test]
6108 async fn dynamic_accepts_a_max_depth_prefix() {
6109 let producer = origin(1).produce();
6110 let path = (0..Path::MAX_PARTS)
6111 .map(|i| format!("s{i}"))
6112 .collect::<Vec<_>>()
6113 .join("/");
6114 let mut announced = producer.consume().announced();
6115
6116 let dynamic = producer.dynamic(&path, Route::default()).expect("max depth is allowed");
6117 announced.assert_next_active(&path);
6118
6119 let _pending = producer.consume().request_broadcast(&path);
6120 let request = queued(&dynamic).await;
6121 assert_eq!(request.path().as_str(), path);
6122 }
6123
6124 #[tokio::test]
6125 async fn dynamic_exclusion_skips_routes_through_the_subscriber() {
6126 let producer = origin(1).produce();
6127 let _server = producer
6128 .dynamic("live", Route::default().with_hops(hops(&[7])))
6129 .unwrap();
6130
6131 let mut excluded = producer.consume().excluding(origin(7)).announced();
6132 excluded.assert_next_wait();
6133
6134 let mut clean = producer.consume().excluding(origin(8)).announced();
6135 clean.assert_next_active("live");
6136 }
6137
6138 #[tokio::test]
6140 async fn announce_consumer_is_a_stream() {
6141 use futures::StreamExt;
6142 let producer = origin(1).produce();
6143 let server = producer.dynamic("live", Route::default()).unwrap();
6144 let mut announced = producer.consume().announced();
6145 let update = StreamExt::next(&mut announced)
6146 .now_or_never()
6147 .expect("next")
6148 .expect("no next");
6149 assert_eq!(update.prefix.as_str(), "live");
6150 assert_eq!(update.kind, AnnounceKind::Announced);
6151 assert!(StreamExt::next(&mut announced).now_or_never().is_none());
6152 drop(server);
6153 let update = StreamExt::next(&mut announced)
6154 .now_or_never()
6155 .expect("next")
6156 .expect("no next");
6157 assert_eq!(update.kind, AnnounceKind::Retracted);
6158 }
6159
6160 #[tokio::test]
6161 async fn dynamic_retracts() {
6162 let producer = origin(1).produce();
6163 let server = producer.dynamic("live", Route::default()).unwrap();
6164 let mut announced = producer.consume().announced();
6165 announced.assert_next_active("live");
6166
6167 drop(server);
6168 announced.assert_next_ended("live");
6169 }
6170
6171 #[test]
6172 fn charged_wildcard_cost_accumulates_across_hops() {
6173 let first = Cost::new(4).charged(1);
6174 let second = first.charged(2);
6175 assert_eq!(second, Cost { warm: 7, cold: 7 });
6176 }
6177
6178 #[tokio::test]
6179 async fn local_broadcast_is_invisible_until_announced() {
6180 let producer = origin(1).produce();
6181 let mut local = producer.consume().announced();
6182 let mut peer = producer.consume().excluding(Hop::UNKNOWN).announced();
6183 let broadcast = producer.create_broadcast("room/alice").unwrap();
6184 local.assert_next_wait();
6185 peer.assert_next_wait();
6186
6187 broadcast.announce(Route::default()).unwrap();
6188 local.assert_next_active("room/alice");
6189 peer.assert_next_active("room/alice");
6190
6191 drop(broadcast);
6192 local.assert_next_ended("room/alice");
6193 peer.assert_next_ended("room/alice");
6194 }
6195
6196 #[tokio::test]
6197 async fn served_route_materializes_on_demand() {
6198 let producer = origin(1).produce();
6199 let consumer = producer.consume();
6200
6201 let server = producer.dynamic("room", Route::default()).unwrap();
6202
6203 let pending = consumer.request_broadcast("room/alice");
6204 let request = queued(&server).await;
6205 assert_eq!(request.path().as_str(), "room/alice");
6206
6207 let source = broadcast::Info::new().produce();
6208 request.accept(&source);
6209
6210 let resolved = pending.await.expect("resolves");
6211 assert_eq!(resolved.info().path.as_str(), "room/alice");
6213
6214 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
6216 assert!(again.is_clone(&resolved));
6217 }
6218
6219 #[tokio::test]
6220 async fn served_requests_coalesce() {
6221 let producer = origin(1).produce();
6222 let consumer = producer.consume();
6223 let server = producer.dynamic("room", Route::default()).unwrap();
6224
6225 let first = consumer.request_broadcast("room/alice");
6226 let second = consumer.request_broadcast("room/alice");
6227
6228 let request = queued(&server).await;
6229 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
6231
6232 let source = broadcast::Info::new().produce();
6233 request.accept(&source);
6234
6235 let first = first.await.expect("resolves");
6236 let second = second.await.expect("resolves");
6237 assert!(first.is_clone(&second));
6238 }
6239
6240 #[tokio::test]
6241 async fn retract_rejects_pending_requests() {
6242 let producer = origin(1).produce();
6243 let consumer = producer.consume();
6244 let server = producer.dynamic("room", Route::default()).unwrap();
6245
6246 let pending = consumer.request_broadcast("room/alice");
6247 drop(server);
6248
6249 let err = pending.await.err().unwrap();
6250 assert!(matches!(err, Error::Unroutable));
6251
6252 let err = consumer
6254 .request_broadcast("room/alice")
6255 .now_or_never()
6256 .expect("unroutable")
6257 .err()
6258 .unwrap();
6259 assert!(matches!(err, Error::Unroutable));
6260 }
6261
6262 #[tokio::test]
6263 async fn routed_broadcast_survives_serving_route_retraction() {
6264 let producer = origin(1).produce();
6265 let consumer = producer.consume();
6266
6267 let standby_server = producer.dynamic("room", Route::default()).unwrap();
6270 let second_server = producer.dynamic("room", Route::default()).unwrap();
6271 let incumbent_server = producer.dynamic("room", Route::default()).unwrap();
6272
6273 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
6274 assert!((&mut resolving).now_or_never().is_none());
6275
6276 drop(incumbent_server);
6282 assert!((&mut resolving).now_or_never().is_none());
6283 drop(second_server);
6284 assert!((&mut resolving).now_or_never().is_none());
6285
6286 let request = queued(&standby_server).await;
6287 let source = broadcast::Info::new().produce();
6288 request.accept(&source);
6289
6290 let resolved = resolving.await.expect("resolves via the standby");
6291 assert_eq!(resolved.info().path.as_str(), "room/alice");
6292 }
6293
6294 #[tokio::test]
6295 async fn split_horizon_skips_routes_through_the_requester() {
6296 let producer = origin(1).produce();
6297 let _server = producer
6298 .dynamic("room", Route::default().with_hops(hops(&[7])))
6299 .unwrap();
6300
6301 let excluded = producer.consume().excluding(origin(7));
6303 let err = excluded
6304 .request_broadcast("room/alice")
6305 .now_or_never()
6306 .expect("unroutable")
6307 .err()
6308 .unwrap();
6309 assert!(matches!(err, Error::Unroutable));
6310
6311 let clean = producer.consume().excluding(origin(8));
6313 let pending = clean.request_broadcast("room/alice");
6314 assert!(pending.now_or_never().is_none());
6315 }
6316
6317 #[tokio::test]
6318 async fn routes_report_where_they_entered() {
6319 let producer = origin(1).produce();
6320 let peer = producer.clone().peer();
6321 let mut announced = producer.consume().announced();
6322
6323 let _ingest = producer
6324 .dynamic("client", Route::default().with_hops(hops(&[5])).with_via(origin(5)))
6325 .unwrap();
6326 let _gateway = producer.publish("gateway", Route::default()).unwrap();
6327 let _forwarded = peer
6328 .dynamic(
6329 "forwarded",
6330 Route::default().with_hops(hops(&[5, 7])).with_via(origin(7)),
6331 )
6332 .unwrap();
6333
6334 assert_eq!(announced.assert_next_active("client").source(), Source::Local);
6335 assert_eq!(
6336 announced.assert_next_active("forwarded").source(),
6337 Source::Peer(origin(7))
6338 );
6339 assert_eq!(announced.assert_next_active("gateway").source(), Source::Local);
6340
6341 let scoped = peer.scope("room", &Patterns::from(Pattern::all())).unwrap();
6343 let _nested = scoped.dynamic("x", Route::default().with_via(origin(8))).unwrap();
6344 assert_eq!(announced.assert_next_active("room/x").source(), Source::Peer(origin(8)));
6345 }
6346
6347 #[tokio::test]
6351 async fn withdrawn_peer_hides_routes_through_it() {
6352 let producer = origin(1).produce();
6353 let peer = producer.clone().peer();
6354 let mut announced = producer.consume().announced();
6355
6356 let direct = || Route::default().with_hops(hops(&[9, 2])).with_via(origin(2));
6358 let first = peer.dynamic("room", direct()).unwrap();
6359 let relayed = peer
6360 .dynamic("room", Route::default().with_hops(hops(&[9, 2, 3])).with_via(origin(3)))
6361 .unwrap();
6362 announced.assert_next_active("room");
6363 announced.assert_next_wait();
6364
6365 first.withdrawn();
6367 announced.assert_next_ended("room");
6368 announced.assert_next_wait();
6369
6370 let second = peer.dynamic("room", direct()).unwrap();
6372 announced.assert_next_active("room");
6373 assert!(producer.shared.lock().withdrawn.is_empty());
6374
6375 second.withdrawn();
6377 drop(relayed);
6378 announced.assert_next_ended("room");
6379 assert!(producer.shared.lock().withdrawn.is_empty());
6380 }
6381
6382 #[tokio::test]
6384 async fn old_session_withdrawal_keeps_newer_route() {
6385 for restart in [false, true] {
6386 let producer = origin(1).produce();
6387 let peer = producer.clone().peer();
6388 let mut announced = producer.consume().announced();
6389 let route = Route::default().with_hops(hops(&[9, 2]));
6390 let first = peer.dynamic("room", route.clone()).unwrap();
6391 let second = peer.dynamic("room", route.clone()).unwrap();
6392 announced.assert_next_active("room");
6393 announced.assert_next_wait();
6394 let remaining = if restart {
6395 first.update(route).unwrap();
6397 second.withdrawn();
6398 first
6399 } else {
6400 first.withdrawn();
6401 second
6402 };
6403 announced.assert_next_wait();
6404 assert!(producer.shared.lock().withdrawn.is_empty());
6405 remaining.withdrawn();
6406 announced.assert_next_ended("room");
6407 assert!(producer.shared.lock().withdrawn.is_empty());
6408 }
6409 }
6410
6411 #[tokio::test]
6412 async fn withdrawn_route_no_longer_covers_requests() {
6413 let producer = origin(1).produce();
6414 let peer = producer.clone().peer();
6415 let direct = peer.dynamic("room", Route::default().with_hops(hops(&[9, 2]))).unwrap();
6416 let _relayed = peer
6417 .dynamic("room", Route::default().with_hops(hops(&[9, 2, 3])))
6418 .unwrap();
6419 let path = Path::new("room/video");
6420 let id = producer
6421 .shared
6422 .read()
6423 .routes
6424 .covering(&path)
6425 .find(|entry| entry.hops.iter().last() == Some(&origin(3)))
6426 .unwrap()
6427 .id;
6428 assert!(producer.shared.read().routes.covers(&path, id));
6429 direct.withdrawn();
6430 assert!(!producer.shared.read().routes.covers(&path, id));
6431 }
6432
6433 #[tokio::test]
6436 async fn source_change_is_an_update() {
6437 let producer = origin(1).produce();
6438 let peer = producer.clone().peer();
6439 let mut announced = producer.consume().announced();
6440
6441 let route = Route::default().with_hops(hops(&[7])).with_via(origin(7));
6442 let _forwarded = peer.dynamic("room", route.clone()).unwrap();
6443 assert_eq!(announced.assert_next_active("room").source(), Source::Peer(origin(7)));
6444
6445 let local = producer.dynamic("room", route).unwrap();
6447 let update = announced.next().now_or_never().expect("next blocked").expect("no next");
6448 assert_eq!(update.kind, AnnounceKind::Updated);
6449 assert_eq!(update.route.source(), Source::Local);
6450
6451 drop(local);
6452 assert_eq!(announced.assert_next_active("room").source(), Source::Peer(origin(7)));
6453 }
6454
6455 #[tokio::test]
6456 async fn local_view_hides_peer_routes() {
6457 let producer = origin(1).produce();
6458 let peer = producer.clone().peer();
6459 let mut local = producer.consume().local().announced();
6460
6461 let _forwarded = peer
6462 .dynamic("remote", Route::default().with_hops(hops(&[7])).with_via(origin(7)))
6463 .unwrap();
6464 local.assert_next_wait();
6465
6466 let _shadow = peer
6470 .dynamic("both", Route::default().with_hops(hops(&[7])).with_via(origin(7)))
6471 .unwrap();
6472 let ingest = producer
6473 .dynamic(
6474 "both",
6475 Route::default().with_hops(hops(&[5])).with_via(origin(5)).with_cost(9),
6476 )
6477 .unwrap();
6478 assert_eq!(local.assert_next_active("both").source(), Source::Local);
6479 drop(ingest);
6480 local.assert_next_ended("both");
6481
6482 let err = producer
6485 .consume()
6486 .local()
6487 .request_broadcast("remote/alice")
6488 .now_or_never()
6489 .expect("unroutable")
6490 .err()
6491 .unwrap();
6492 assert!(matches!(err, Error::Unroutable));
6493 assert!(
6494 producer
6495 .consume()
6496 .request_broadcast("remote/alice")
6497 .now_or_never()
6498 .is_none()
6499 );
6500 }
6501
6502 #[tokio::test]
6506 async fn handler_rejection_is_final() {
6507 let producer = origin(1).produce();
6508 let consumer = producer.consume();
6509 let server = producer.dynamic("room", Route::default()).unwrap();
6510
6511 let pending = consumer.request_broadcast("room/alice");
6512 let request = queued(&server).await;
6513 request.reject(Error::Unroutable);
6514 let err = tokio::time::timeout(Duration::from_secs(5), pending)
6515 .await
6516 .expect("the front must give up, not spin")
6517 .err()
6518 .unwrap();
6519 assert!(matches!(err, Error::Unroutable));
6520
6521 let pending = consumer.request_broadcast("room/bob");
6523 let request = queued(&server).await;
6524 assert_eq!(request.path().as_str(), "room/bob");
6525 let served = broadcast::Info::new().produce();
6526 request.accept(&served);
6527 pending.await.expect("resolves");
6528 }
6529
6530 #[tokio::test]
6533 async fn routed_broadcast_waits_out_a_rejection() {
6534 let producer = origin(1).produce();
6535 let consumer = producer.consume();
6536 let server = producer.dynamic("room", Route::default()).unwrap();
6537
6538 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
6539 assert!((&mut resolving).now_or_never().is_none());
6540 let request = queued(&server).await;
6541 request.reject(Error::Unroutable);
6542
6543 for _ in 0..20 {
6545 tokio::task::yield_now().await;
6546 }
6547 assert!((&mut resolving).now_or_never().is_none());
6548 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
6549
6550 server.update(Route::default().with_cost(2)).unwrap();
6552 assert!((&mut resolving).now_or_never().is_none());
6553 let request = queued(&server).await;
6554 let served = broadcast::Info::new().produce();
6555 request.accept(&served);
6556 resolving.await.expect("resolves");
6557 }
6558
6559 #[tokio::test]
6562 async fn routed_broadcast_reports_teardown_as_closed() {
6563 let (producer, driver) = Producer::new(Config::new(origin(1)));
6564 let consumer = producer.consume();
6565 let _server = producer.dynamic("room", Route::default()).unwrap();
6566
6567 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
6569 assert!((&mut resolving).now_or_never().is_none());
6570
6571 drop(driver);
6572
6573 let err = tokio::time::timeout(Duration::from_secs(5), resolving)
6574 .await
6575 .expect("teardown resolves the wait")
6576 .err()
6577 .unwrap();
6578 assert!(matches!(err, Error::Closed), "unexpected end: {err}");
6579 }
6580
6581 #[tokio::test]
6584 async fn routed_broadcast_wakes_for_a_local_broadcast() {
6585 let producer = origin(1).produce();
6586 let consumer = producer.consume();
6587 let server = producer.dynamic("room", Route::default()).unwrap();
6588
6589 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
6590 assert!((&mut resolving).now_or_never().is_none());
6591 queued(&server).await.reject(Error::Unroutable);
6592 for _ in 0..20 {
6593 tokio::task::yield_now().await;
6594 }
6595 assert!((&mut resolving).now_or_never().is_none());
6596
6597 let _local = producer.publish("room/alice", Route::default()).unwrap();
6599 let resolved = resolving.await.expect("resolves locally");
6600 assert_eq!(resolved.info().path.as_str(), "room/alice");
6601 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
6602 }
6603
6604 #[tokio::test]
6607 async fn late_track_on_a_served_front_replays() {
6608 let producer = origin(1).produce();
6609 let consumer = producer.consume();
6610 let server = producer.dynamic("room", Route::default()).unwrap();
6611
6612 let source = broadcast::Info::new().produce();
6613 for name in ["a", "b"] {
6614 let track = source.create_track(name, None).unwrap();
6615 let mut group = track.append_group().unwrap();
6616 group.write_frame(crate::Timestamp::ZERO, name.as_bytes()).unwrap();
6617 group.finish().unwrap();
6618 std::mem::forget(track);
6620 }
6621
6622 let pending = consumer.request_broadcast("room/alice");
6623 queued(&server).await.accept(&source);
6624 let resolved = pending.await.expect("resolves");
6625
6626 let budget = track::Subscription::default().with_max_age(Duration::from_secs(3600));
6627 for name in ["a", "b"] {
6628 let mut subscription = resolved
6629 .track(name)
6630 .unwrap()
6631 .subscribe(budget.clone())
6632 .await
6633 .expect("subscribe");
6634 let mut group = tokio::time::timeout(Duration::from_secs(5), subscription.recv_group())
6635 .await
6636 .expect("the late track must replay, not park")
6637 .expect("recv group")
6638 .expect("track ended early");
6639 let frame = group.read_frame().await.expect("read frame").expect("frame");
6640 assert_eq!(&frame.payload[..], name.as_bytes());
6641 }
6642 }
6643
6644 #[tokio::test]
6645 async fn most_specific_prefix_shadows() {
6646 let producer = origin(1).produce();
6647 let consumer = producer.consume();
6648
6649 let broad_server = producer.dynamic("", Route::default()).unwrap();
6650 let _narrow = producer.announce(".dash", Route::default()).unwrap();
6653
6654 let err = consumer
6655 .request_broadcast(".dash/pid")
6656 .now_or_never()
6657 .expect("unroutable")
6658 .err()
6659 .unwrap();
6660 assert!(matches!(err, Error::Unroutable));
6661
6662 let _pending = consumer.request_broadcast("room/alice");
6664 let request = queued(&broad_server).await;
6665 assert_eq!(request.path().as_str(), "room/alice");
6666 }
6667
6668 #[tokio::test]
6669 async fn root_dynamic_serves_any_path() {
6670 let producer = origin(1).produce();
6671 let consumer = producer.consume();
6672 let mut announced = consumer.announced();
6673 let dynamic = producer.dynamic("", Route::default()).unwrap();
6674 announced.assert_next_active("");
6676
6677 let pending = consumer.request_broadcast("anything/at/all");
6678 let request = queued(&dynamic).await;
6679 assert_eq!(request.path().as_str(), "anything/at/all");
6680
6681 let source = broadcast::Info::new().produce();
6682 request.accept(&source);
6683 let resolved = pending.await.expect("resolves");
6684 assert_eq!(resolved.info().path.as_str(), "anything/at/all");
6685
6686 drop(dynamic);
6688 announced.assert_next_ended("");
6689 let err = consumer
6690 .request_broadcast("something/else")
6691 .now_or_never()
6692 .expect("unroutable")
6693 .err()
6694 .unwrap();
6695 assert!(matches!(err, Error::Unroutable));
6696 }
6697
6698 #[tokio::test]
6704 async fn out_of_scope_request_never_reaches_the_dynamic_handler() {
6705 let producer = origin(1).produce();
6706 let dynamic = producer.dynamic("", Route::default()).unwrap();
6707 let scoped = producer.consume().scope("", &scopes(&["tenant-a"])).unwrap();
6708
6709 for path in ["tenant-b/live", "tenant-a-other/live"] {
6712 let refused = scoped
6713 .request_broadcast(path)
6714 .now_or_never()
6715 .expect("an out-of-scope request must be refused synchronously, not queued");
6716 assert!(matches!(refused, Err(Error::Unauthorized)));
6717 assert!(
6718 dynamic.requested_broadcast().now_or_never().is_none(),
6719 "the dynamic handler was asked to create a broadcast the requester may not read"
6720 );
6721 }
6722 }
6723
6724 #[tokio::test]
6725 async fn routed_waits_for_coverage() {
6726 let producer = origin(1).produce();
6727 let consumer = producer.consume();
6728
6729 let mut fut = consumer.routed("room/alice").boxed();
6730 assert!((&mut fut).now_or_never().is_none());
6731
6732 let _a = producer.announce("room", Route::default().with_cost(3)).unwrap();
6734 let route = fut.now_or_never().expect("covered").expect("routed");
6735 assert_eq!(route.cost, Cost::new(3));
6736
6737 consumer
6739 .routed("room/alice/cam")
6740 .now_or_never()
6741 .expect("covered")
6742 .expect("routed");
6743 }
6744
6745 #[tokio::test]
6746 async fn routed_ignores_deeper_routes() {
6747 let producer = origin(1).produce();
6748 let consumer = producer.consume();
6749
6750 let _deep = producer.announce("room/alice/cam", Route::default()).unwrap();
6752 let mut fut = consumer.routed("room/alice").boxed();
6753 assert!((&mut fut).now_or_never().is_none());
6754
6755 let _exact = producer.announce("room/alice", Route::default()).unwrap();
6756 fut.now_or_never().expect("covered").expect("routed");
6757 }
6758
6759 #[tokio::test]
6760 async fn routed_accepts_a_max_depth_path() {
6761 let producer = origin(1).produce();
6762 let consumer = producer.consume();
6763 let path = (0..Path::MAX_PARTS)
6764 .map(|i| format!("s{i}"))
6765 .collect::<Vec<_>>()
6766 .join("/");
6767 assert_eq!(Path::new(&path).parts().count(), Path::MAX_PARTS);
6768
6769 assert!(consumer.allowed().matches(&path));
6770
6771 let mut fut = consumer.routed(&path).boxed();
6772 assert!((&mut fut).now_or_never().is_none());
6773
6774 let _a = producer.announce("", Route::default()).unwrap();
6776 fut.now_or_never().expect("covered").expect("routed");
6777 }
6778
6779 #[tokio::test]
6780 async fn teardown_ends_everything() {
6781 let (producer, driver) = Producer::new(Config::new(origin(1)));
6782 let consumer = producer.consume();
6783 let _announcement = producer.announce("room", Route::default()).unwrap();
6784 let mut announced = consumer.announced();
6785 announced.assert_next_active("room");
6786
6787 let _server = producer.dynamic("served", Route::default()).unwrap();
6788 let pending = consumer.request_broadcast("served/path");
6789
6790 drop(driver);
6791
6792 announced.assert_next_active("served");
6794 assert!(announced.next().now_or_never().expect("ended").is_none());
6795
6796 assert!(pending.now_or_never().expect("rejected").is_err());
6798 assert!(matches!(producer.announce("x", Route::default()), Err(Error::Closed)));
6799 assert!(matches!(producer.create_broadcast("x"), Err(Error::Closed)));
6800 let err = consumer
6801 .request_broadcast("y")
6802 .now_or_never()
6803 .expect("closed")
6804 .err()
6805 .unwrap();
6806 assert!(matches!(err, Error::Closed));
6807
6808 let mut late = consumer.announced();
6810 assert!(late.next().now_or_never().expect("ended").is_none());
6811 }
6812
6813 struct ResumeRig {
6816 producer: Producer,
6817 resolved: broadcast::Consumer,
6818 subscription: track::Subscriber,
6819 incumbent_track: track::Producer,
6822 }
6823
6824 impl ResumeRig {
6825 async fn new(first: &[u64]) -> (Self, Dynamic, broadcast::Producer) {
6828 let producer = origin(1).produce();
6829 let consumer = producer.consume();
6830
6831 let server = producer
6832 .dynamic("room", Route::default().with_hops(hops(first)))
6833 .unwrap();
6834
6835 let pending = consumer.request_broadcast("room/alice");
6836 let request = queued(&server).await;
6837 let source = broadcast::Info::new().produce();
6838 let track = source.create_track("video", None).unwrap();
6839 let mut group = track.append_group().unwrap();
6840 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
6841 group.finish().unwrap();
6842 request.accept(&source);
6843
6844 let resolved = pending.await.expect("resolves");
6845 let mut subscription = resolved
6846 .track("video")
6847 .unwrap()
6848 .subscribe(None)
6849 .await
6850 .expect("subscribe");
6851 let mut group = subscription
6852 .recv_group()
6853 .await
6854 .expect("recv group")
6855 .expect("track ended early");
6856 let frame = group.read_frame().await.expect("read frame").expect("frame");
6857 assert_eq!(&frame.payload[..], b"before");
6858
6859 (
6860 Self {
6861 producer,
6862 resolved,
6863 subscription,
6864 incumbent_track: track,
6865 },
6866 server,
6867 source,
6868 )
6869 }
6870
6871 fn standby(&self, first: &[u64]) -> Dynamic {
6874 self.producer
6875 .dynamic("room", Route::default().with_hops(hops(first)))
6876 .unwrap()
6877 }
6878 }
6879
6880 async fn assert_resumes(rig: &mut ResumeRig, server: &Dynamic) {
6885 let request = queued(server).await;
6886 let replacement = broadcast::Info::new().produce();
6887 let track = replacement.create_track("video", None).unwrap();
6888 let mut group = track.append_group().unwrap();
6891 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
6892 group.finish().unwrap();
6893 request.accept(&replacement);
6894
6895 let mut group = track.append_group().unwrap();
6896 group.write_frame(crate::Timestamp::ZERO, b"resumed".as_ref()).unwrap();
6897 group.finish().unwrap();
6898
6899 let mut group = rig
6900 .subscription
6901 .recv_group()
6902 .await
6903 .expect("subscription survives the failover")
6904 .expect("track ended early");
6905 let frame = group.read_frame().await.expect("read frame").expect("frame");
6906 assert_eq!(&frame.payload[..], b"resumed");
6907 }
6908
6909 #[tokio::test]
6912 async fn driver_resolves_with_live_consumers() {
6913 let (producer, driver) = Producer::new(Config::new(origin(1)));
6914 let consumer = producer.consume();
6915 let run = crate::time::run(driver);
6916 drop(producer);
6917 tokio::time::timeout(Duration::from_secs(5), run)
6918 .await
6919 .expect("driver must finish once the producers are gone");
6920 drop(consumer);
6921 }
6922
6923 #[tokio::test]
6924 async fn remote_source_resumes_through_same_first_hop() {
6925 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
6926 let standby_server = rig.standby(&[10, 20]);
6927
6928 drop(incumbent);
6930 drop(source);
6931
6932 assert_resumes(&mut rig, &standby_server).await;
6934 }
6935
6936 #[tokio::test]
6940 async fn incompatible_successor_is_refused() {
6941 for replacement in [
6942 track::Info::default().with_timescale(crate::Timescale::MICRO),
6943 track::Info::default().with_priority(7),
6944 track::Info::default().with_max_age(Duration::from_secs(7)),
6945 ] {
6946 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
6947 let standby_server = rig.standby(&[10, 20]);
6948 drop(incumbent);
6949 drop(source);
6950
6951 let request = queued(&standby_server).await;
6954 let successor = broadcast::Info::new().produce();
6955 let track = successor.create_track("video", replacement).unwrap();
6956 let mut group = track.append_group().unwrap();
6957 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
6958 group.finish().unwrap();
6959 request.accept(&successor);
6960
6961 assert!(
6962 matches!(rig.subscription.recv_group().await, Err(Error::Unsupported)),
6963 "the subscription must abort rather than resume onto incompatible metadata"
6964 );
6965
6966 let reopened = rig.resolved.track("video").unwrap();
6968 assert!(matches!(reopened.query().await, Err(Error::Unsupported)));
6969 assert!(matches!(reopened.subscribe(None).await, Err(Error::Unsupported)));
6970 }
6971 }
6972
6973 #[tokio::test]
6974 async fn different_first_hop_ends_the_subscription() {
6975 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
6976 let rival_server = rig.standby(&[11]);
6978
6979 drop(incumbent);
6982 drop(source);
6983 rig.incumbent_track.abort(Error::Dropped).unwrap();
6984
6985 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
6987 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
6988
6989 let consumer = rig.producer.consume();
6991 let pending = consumer.request_broadcast("room/alice");
6992 let request = queued(&rival_server).await;
6993 let replacement = broadcast::Info::new().produce();
6994 request.accept(&replacement);
6995 pending.await.expect("re-request resolves through the rival");
6996 }
6997
6998 #[tokio::test]
7003 async fn a_first_hop_update_drains_the_old_publisher() {
7004 for first in [&[][..], &[0][..], &[10][..]] {
7005 let (mut rig, server, _source) = ResumeRig::new(first).await;
7006 server.update(Route::default().with_hops(hops(&[11]))).unwrap();
7007
7008 let mut group = rig.incumbent_track.append_group().unwrap();
7010 group.write_frame(crate::Timestamp::ZERO, b"draining".as_ref()).unwrap();
7011 group.finish().unwrap();
7012 let mut group = next_group(&mut rig.subscription)
7013 .await
7014 .expect("the in-flight subscription survives the update")
7015 .expect("track ended early");
7016 let frame = group.read_frame().await.expect("read frame").expect("frame");
7017 assert_eq!(&frame.payload[..], b"draining", "first hop {first:?}");
7018
7019 let consumer = rig.producer.consume();
7022 let pending = consumer.request_broadcast("room/alice");
7023 let request = queued(&server).await;
7024 let replacement = broadcast::Info::new().produce();
7025 let track = replacement
7026 .create_track("video", track::Info::default().with_priority(7))
7027 .unwrap();
7028 let mut group = track.append_group().unwrap();
7029 group.write_frame(crate::Timestamp::ZERO, b"new".as_ref()).unwrap();
7030 group.finish().unwrap();
7031 request.accept(&replacement);
7032
7033 let resolved = pending.await.expect("resolves through the new publisher");
7034 assert!(
7035 !resolved.is_clone(&rig.resolved),
7036 "first hop {first:?} joined the old front"
7037 );
7038 let mut fresh = resolved
7039 .track("video")
7040 .unwrap()
7041 .subscribe(None)
7042 .await
7043 .expect("the old publisher's track info does not apply");
7044 let mut group = next_group(&mut fresh)
7045 .await
7046 .expect("recv group")
7047 .expect("track ended early");
7048 let frame = group.read_frame().await.expect("read frame").expect("frame");
7049 assert_eq!(&frame.payload[..], b"new");
7050
7051 rig.incumbent_track.finish().unwrap();
7054 let end = next_group(&mut rig.subscription).await;
7055 assert!(
7056 !matches!(end, Ok(Some(_))),
7057 "first hop {first:?} spliced the new publisher into a live subscription"
7058 );
7059 }
7060 }
7061
7062 #[tokio::test]
7065 async fn a_first_hop_update_resumes_through_the_same_publisher() {
7066 let (mut rig, incumbent, _source) = ResumeRig::new(&[10]).await;
7067 let standby_server = rig.standby(&[10, 20]);
7068
7069 incumbent.update(Route::default().with_hops(hops(&[11]))).unwrap();
7070
7071 assert_resumes(&mut rig, &standby_server).await;
7072 assert!(
7073 incumbent.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
7074 "the front never asks the new publisher"
7075 );
7076 }
7077
7078 #[tokio::test]
7079 async fn anonymous_routes_never_resume() {
7080 let (mut rig, incumbent, source) = ResumeRig::new(&[]).await;
7083 let _twin_server = rig.standby(&[]);
7084
7085 drop(incumbent);
7086 drop(source);
7087 rig.incumbent_track.abort(Error::Dropped).unwrap();
7088
7089 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
7090 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
7091 }
7092
7093 #[tokio::test]
7100 async fn anonymous_handoff_serves_the_newcomer_immediately() {
7101 let producer = origin(1).produce();
7102
7103 let server_a = producer
7105 .dynamic("room", Route::default().with_hops(hops(&[10])))
7106 .unwrap();
7107
7108 let consumer = producer.consume().excluding(origin(30));
7111 let pending = consumer.request_broadcast("room/alice");
7112 let request = queued(&server_a).await;
7113 let source_a = broadcast::Info::new().produce();
7114 let track_a = source_a.create_track("video", None).unwrap();
7115 let mut group = track_a.append_group().unwrap();
7116 group.write_frame(crate::Timestamp::ZERO, b"from-a".as_ref()).unwrap();
7117 group.finish().unwrap();
7118 request.accept(&source_a);
7119
7120 let resolved_a = pending.await.expect("resolves");
7121 let mut sub_a = resolved_a
7122 .track("video")
7123 .unwrap()
7124 .subscribe(None)
7125 .await
7126 .expect("subscribe");
7127 let mut group = sub_a
7128 .recv_group()
7129 .await
7130 .expect("recv group")
7131 .expect("track ended early");
7132 assert_eq!(
7133 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
7134 b"from-a"
7135 );
7136
7137 drop(track_a);
7140 drop(source_a);
7141 drop(server_a);
7142
7143 let err = sub_a.recv_group().await.err().expect("front closed");
7145 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
7146
7147 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
7151 settle(|| {
7152 matches!(
7153 consumer.request_broadcast("room/alice").now_or_never(),
7154 Some(Err(Error::Unroutable))
7155 )
7156 })
7157 .await;
7158
7159 let server_b = producer
7161 .dynamic("room", Route::default().with_hops(hops(&[20])))
7162 .unwrap();
7163 let pending = consumer.request_broadcast("room/alice");
7164 let request = queued(&server_b).await;
7165 let source_b = broadcast::Info::new().produce();
7166 let track_b = source_b.create_track("video", None).unwrap();
7167 let mut group = track_b.append_group().unwrap();
7168 group.write_frame(crate::Timestamp::ZERO, b"from-b".as_ref()).unwrap();
7169 group.finish().unwrap();
7170 request.accept(&source_b);
7171
7172 let resolved_b = pending.await.expect("B's front is served immediately");
7173 assert!(
7174 !resolved_b.is_clone(&resolved_a),
7175 "B must not splice into A's closed front"
7176 );
7177
7178 let mut sub_b = resolved_b
7179 .track("video")
7180 .unwrap()
7181 .subscribe(None)
7182 .await
7183 .expect("subscribe");
7184 let mut group = sub_b
7185 .recv_group()
7186 .await
7187 .expect("recv group")
7188 .expect("track ended early");
7189 assert_eq!(
7190 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
7191 b"from-b"
7192 );
7193 }
7194
7195 #[tokio::test]
7196 async fn reprice_is_invisible_to_the_subscription() {
7197 let (rig, incumbent, source) = ResumeRig::new(&[10]).await;
7198
7199 incumbent
7202 .update(Route::default().with_hops(hops(&[10])).with_cost(9))
7203 .unwrap();
7204
7205 let track = source.create_track("audio", None).unwrap();
7206 let mut group = track.append_group().unwrap();
7207 group.write_frame(crate::Timestamp::ZERO, b"steady".as_ref()).unwrap();
7208 group.finish().unwrap();
7209
7210 let mut audio = rig
7211 .resolved
7212 .track("audio")
7213 .unwrap()
7214 .subscribe(None)
7215 .await
7216 .expect("subscribe survives the reprice");
7217 let mut group = audio
7218 .recv_group()
7219 .await
7220 .expect("recv group")
7221 .expect("track ended early");
7222 let frame = group.read_frame().await.expect("read frame").expect("frame");
7223 assert_eq!(&frame.payload[..], b"steady");
7224 }
7225
7226 #[tokio::test]
7227 async fn drain_reprice_migrates_before_the_session_dies() {
7228 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
7229 let standby_server = rig.standby(&[10, 20]);
7230
7231 incumbent
7235 .update(Route::default().with_hops(hops(&[10])).with_cost(Cost::DRAIN))
7236 .unwrap();
7237
7238 assert_resumes(&mut rig, &standby_server).await;
7239
7240 drop(incumbent);
7242 drop(source);
7243 }
7244
7245 #[tokio::test]
7246 async fn local_sources_splice_newest_first() {
7247 let producer = origin(1).produce();
7248 let consumer = producer.consume();
7249
7250 let first = producer.publish("room/alice", Route::default()).unwrap();
7251 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7252
7253 let second = producer.publish("room/alice", Route::default()).unwrap();
7255 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
7256 assert!(again.is_clone(&resolved));
7257
7258 first.close();
7260 settle(|| consumer.get_broadcast("room/alice").is_some()).await;
7261 second.close();
7262 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
7263
7264 let _third = producer.publish("room/alice", Route::default()).unwrap();
7266 assert!(consumer.get_broadcast("room/alice").is_some());
7267 }
7268
7269 #[tokio::test]
7277 async fn a_finished_broadcast_concludes_in_flight_subscriptions() {
7278 let producer = origin(1).produce();
7279 let consumer = producer.consume();
7280
7281 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
7282 let track = broadcast.create_track("video", None).unwrap();
7283
7284 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7285 let mut subscription = resolved
7286 .track("video")
7287 .unwrap()
7288 .subscribe(None)
7289 .await
7290 .expect("subscribe");
7291 let mut group = track.append_group().unwrap();
7293 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
7294 group.finish().unwrap();
7295 track.finish().unwrap();
7296 drop(track);
7297 broadcast.close();
7298
7299 let mut group = next_group(&mut subscription)
7300 .await
7301 .expect("a cleanly finished track was served as an error")
7302 .expect("the track ended before its last group");
7303 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
7304 drop(group);
7305
7306 let end = next_group(&mut subscription)
7307 .await
7308 .expect("a cleanly finished track ended as an error");
7309 assert!(end.is_none(), "a group followed the final one");
7310 }
7311
7312 #[tokio::test]
7319 async fn a_retracted_route_concludes_in_flight_subscriptions() {
7320 let producer = origin(1).produce();
7321 let consumer = producer.consume();
7322 let server = producer
7323 .dynamic("room", Route::default().with_hops(hops(&[10])))
7324 .unwrap();
7325
7326 let pending = consumer.request_broadcast("room/alice");
7327 let request = queued(&server).await;
7328 let source = broadcast::Info::new().produce();
7329 let track = source.create_track("video", None).unwrap();
7330 request.accept(&source);
7331
7332 let resolved = pending.await.expect("resolves");
7333 let mut subscription = resolved
7334 .track("video")
7335 .unwrap()
7336 .subscribe(None)
7337 .await
7338 .expect("subscribe");
7339
7340 source.close();
7343 drop(server);
7344 settle(|| resolved.is_closed()).await;
7345 let mut group = track.append_group().unwrap();
7346 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
7347 group.finish().unwrap();
7348 track.finish().unwrap();
7349 drop(track);
7350
7351 let mut group = next_group(&mut subscription)
7352 .await
7353 .expect("a retracted route's track was served as an error")
7354 .expect("the track ended before its last group");
7355 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
7356 drop(group);
7357
7358 let end = next_group(&mut subscription)
7359 .await
7360 .expect("a cleanly finished track ended as an error");
7361 assert!(end.is_none(), "a group followed the final one");
7362 }
7363
7364 #[tokio::test]
7367 async fn a_closed_source_is_not_requested_again_from_its_standing_route() {
7368 let producer = origin(1).produce();
7369 let server = producer
7370 .dynamic("room", Route::default().with_hops(hops(&[10])))
7371 .unwrap();
7372 let pending = producer.consume().request_broadcast("room/alice");
7373 let source = broadcast::Info::new().produce();
7374 queued(&server).await.accept(&source);
7375 let resolved = pending.await.unwrap();
7376
7377 source.close();
7378 settle(|| resolved.is_closed()).await;
7379 assert!(
7380 server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
7381 "the closed source was requested again"
7382 );
7383 }
7384
7385 #[tokio::test]
7390 async fn origin_front_drops_the_source_when_unused() {
7391 let producer = origin(1).produce();
7392 let consumer = producer.consume();
7393
7394 let broadcast = producer.publish("room/alice", Route::default()).unwrap();
7395 let track = broadcast.create_track("video", None).unwrap();
7396 let mut group = track.append_group().unwrap();
7397 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
7398 group.finish().unwrap();
7399
7400 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7401 let mut subscription = resolved
7402 .track("video")
7403 .unwrap()
7404 .subscribe(None)
7405 .await
7406 .expect("subscribe");
7407 let mut group = subscription.recv_group().await.unwrap().unwrap();
7408 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
7409 drop(group);
7410 drop(subscription);
7411
7412 tokio::time::timeout(Duration::from_secs(1), track.unused())
7413 .await
7414 .expect("source unused should resolve far below TRACK_IDLE_LINGER")
7415 .expect("source closed");
7416
7417 let mut again = resolved
7420 .track("video")
7421 .unwrap()
7422 .subscribe(track::Subscription::default().with_max_age(Duration::from_secs(3600)))
7423 .await
7424 .expect("resubscribe");
7425 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
7426 .await
7427 .expect("cached group is still on the front")
7428 .expect("recv group")
7429 .expect("track ended early");
7430 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
7431
7432 tokio::time::timeout(Duration::from_secs(1), track.used())
7433 .await
7434 .expect("returning reader re-splices the source")
7435 .expect("source closed");
7436
7437 let mut group = track.append_group().unwrap();
7438 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
7439 group.finish().unwrap();
7440 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
7441 .await
7442 .expect("groups past the cached edge come from the re-splice")
7443 .expect("recv group")
7444 .expect("track ended early");
7445 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
7446 }
7447
7448 #[tokio::test]
7458 async fn resumed_reader_skips_warm_groups_behind_the_new_edge() {
7459 let ms = |v: u64| crate::Timestamp::from_millis(v).unwrap();
7460 let (_server, _upstream, mut dynamic, resolved) = served_front().await;
7461 let budget = track::Subscription::default().with_max_age(Duration::from_millis(100));
7462 let write = |source: &track::Producer, sequence: u64, millis: u64| {
7463 let mut group = source.create_group(sequence.into()).unwrap();
7464 group.write_frame(ms(millis), b"x".as_ref()).unwrap();
7465 group.finish().unwrap();
7466 };
7467 let drain = |subscription: &mut track::Subscriber| {
7468 let mut sequences = Vec::new();
7469 while let Poll::Ready(group) = subscription.poll_recv_group(&kio::Waiter::noop()) {
7470 sequences.push(group.unwrap().expect("track ended").sequence);
7471 }
7472 sequences
7473 };
7474
7475 let track = resolved.track("video").unwrap();
7476 let first = budget.clone();
7477 let subscribing = tokio::spawn(async move { track.subscribe(first).await });
7478 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
7479 .await
7480 .expect("the front asked the source")
7481 .expect("request");
7482 let source = request.resolving_start().accept(None);
7483 for (sequence, millis) in [(0, 0), (1, 20), (2, 40), (3, 60)] {
7484 write(&source, sequence, millis);
7485 }
7486 let mut subscription = subscribing.await.unwrap().expect("subscribe");
7487 next_group(&mut subscription).await.unwrap().expect("a cached group");
7488 drain(&mut subscription);
7489 drop(subscription);
7490
7491 tokio::time::timeout(Duration::from_secs(1), source.unused())
7492 .await
7493 .expect("parked")
7494 .expect("source open");
7495 drop(source);
7496
7497 let track = resolved.track("video").unwrap();
7498 let second = budget.clone();
7499 let subscribing = tokio::spawn(async move { track.subscribe(second).await });
7500 let request = tokio::time::timeout(Duration::from_secs(1), dynamic.requested_track())
7501 .await
7502 .expect("the front asked the source again")
7503 .expect("request");
7504 let mut source = request.resolving_start().accept(None);
7505 for (sequence, millis) in [(20, 400), (21, 420), (22, 440)] {
7506 write(&source, sequence, millis);
7507 }
7508 source.start_at(3).unwrap();
7511 let mut subscription = subscribing.await.unwrap().expect("resubscribe");
7512 settle(|| subscription.latest() == Some(22)).await;
7513
7514 assert_eq!(drain(&mut subscription), [3, 20, 21, 22]);
7517 }
7518
7519 #[tokio::test]
7524 async fn chained_front_drops_the_source_when_unused() {
7525 let leaf = origin(1).produce();
7526 let leaf_consumer = leaf.consume();
7527
7528 let broadcast = leaf.publish("room/alice", Route::default()).unwrap();
7529 let track = broadcast.create_track("video", None).unwrap();
7530 let mut group = track.append_group().unwrap();
7531 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
7532 group.finish().unwrap();
7533
7534 let leaf_front = leaf_consumer.request_broadcast("room/alice").await.expect("resolves");
7537
7538 let mid = origin(2).produce();
7539 let mid_server = mid.dynamic("room", Route::default().with_hops(hops(&[10]))).unwrap();
7540 let mid_pending = mid.consume().request_broadcast("room/alice");
7541 queued(&mid_server).await.accept(&leaf_front);
7542 let mid_resolved = mid_pending.await.expect("mid resolves");
7543
7544 let edge = origin(3).produce();
7545 let edge_server = edge.dynamic("room", Route::default().with_hops(hops(&[20]))).unwrap();
7546 let edge_pending = edge.consume().request_broadcast("room/alice");
7547 queued(&edge_server).await.accept(&mid_resolved);
7548 let edge_resolved = edge_pending.await.expect("edge resolves");
7549
7550 let mut subscription = edge_resolved
7551 .track("video")
7552 .unwrap()
7553 .subscribe(None)
7554 .await
7555 .expect("subscribe");
7556 let mut group = subscription.recv_group().await.unwrap().unwrap();
7557 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
7558 drop(group);
7559 drop(subscription);
7560
7561 tokio::time::timeout(Duration::from_secs(5), track.unused())
7562 .await
7563 .expect("chained unused should resolve far below TRACK_IDLE_LINGER")
7564 .expect("source closed");
7565
7566 let cached = edge_resolved.track("video").unwrap().cached_groups();
7567 assert_eq!(
7568 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
7569 vec![0],
7570 "every front keeps the delivered groups after releasing its source"
7571 );
7572
7573 let mut subscription = edge_resolved
7575 .track("video")
7576 .unwrap()
7577 .subscribe(track::Subscription::default().with_max_age(Duration::from_secs(3600)))
7578 .await
7579 .expect("resubscribe");
7580 tokio::time::timeout(Duration::from_secs(5), track.used())
7581 .await
7582 .expect("resubscribe should reach the leaf")
7583 .expect("source open");
7584 let mut group = track.append_group().unwrap();
7585 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
7586 group.finish().unwrap();
7587 let mut group = subscription.recv_group().await.unwrap().unwrap();
7588 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
7589 drop(group);
7590 let mut group = subscription.recv_group().await.unwrap().unwrap();
7591 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
7592 drop(group);
7593 drop(subscription);
7594
7595 tokio::time::timeout(Duration::from_secs(5), track.unused())
7596 .await
7597 .expect("second chained unused should resolve far below TRACK_IDLE_LINGER")
7598 .expect("source closed");
7599
7600 let cached = edge_resolved.track("video").unwrap().cached_groups();
7601 assert_eq!(
7602 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
7603 vec![0, 1],
7604 "repeated demand keeps every complete group while releasing its source"
7605 );
7606
7607 let fetch = edge_resolved.track("video").unwrap().fetch_group(2, None);
7608 let mut fetch = std::pin::pin!(fetch);
7609 assert!(futures::poll!(fetch.as_mut()).is_pending(), "fetch should re-splice");
7610 tokio::time::timeout(Duration::from_secs(5), track.used())
7611 .await
7612 .expect("fetch should reach the leaf")
7613 .expect("source open");
7614 let mut group = track.append_group().unwrap();
7615 group.write_frame(crate::Timestamp::ZERO, b"fetched".as_ref()).unwrap();
7616 group.finish().unwrap();
7617 let mut group = tokio::time::timeout(Duration::from_secs(5), fetch)
7618 .await
7619 .expect("re-spliced source should answer the fetch")
7620 .expect("fetch succeeds");
7621 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"fetched");
7622 }
7623
7624 #[tokio::test]
7628 async fn incompatible_local_source_keeps_the_incumbent() {
7629 let producer = origin(1).produce();
7630 let consumer = producer.consume();
7631
7632 let first = producer.publish("room/alice", Route::default()).unwrap();
7633 let track = first.create_track("video", None).unwrap();
7634 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
7635 let mut subscription = resolved
7636 .track("video")
7637 .unwrap()
7638 .subscribe(None)
7639 .await
7640 .expect("subscribe");
7641 let mut group = track.append_group().unwrap();
7642 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
7643 group.finish().unwrap();
7644 let mut group = subscription.recv_group().await.unwrap().unwrap();
7645 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"before");
7646
7647 let second = producer.publish("room/alice", Route::default()).unwrap();
7649 let _incompatible = second
7650 .create_track("video", track::Info::default().with_timescale(crate::Timescale::MICRO))
7651 .unwrap();
7652 for _ in 0..10 {
7653 tokio::task::yield_now().await;
7654 }
7655
7656 let mut group = track.append_group().unwrap();
7658 group.write_frame(crate::Timestamp::ZERO, b"still".as_ref()).unwrap();
7659 group.finish().unwrap();
7660 let mut group = subscription.recv_group().await.unwrap().unwrap();
7661 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"still");
7662
7663 drop(track);
7665 first.close();
7666 assert!(matches!(subscription.recv_group().await, Err(Error::Unsupported)));
7667 }
7668
7669 #[tokio::test]
7670 async fn multiple_scopes_present_one_broad_prefix() {
7671 let producer = origin(1).produce();
7672 let _a = producer.announce("", Route::default()).unwrap();
7673
7674 let consumer = producer.consume().scope("", &scopes(&["alpha", "beta"])).unwrap();
7675 let mut announced = consumer.announced();
7676 announced.assert_next_active("");
7677 announced.assert_next_wait();
7678 }
7679
7680 #[test]
7681 fn scope_accepts_every_pattern_union() {
7682 let producer = origin(1).produce();
7683
7684 let root = producer.scope("", &Patterns::from(Pattern::all())).unwrap();
7686 assert_eq!(root.allowed(), Patterns::from(Pattern::all()));
7687
7688 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
7690 assert_eq!(scoped.allowed(), scopes(&["room"]));
7691
7692 let multi = producer.scope("", &scopes(&["room", "room/chat", "anon"])).unwrap();
7694 assert_eq!(multi.allowed(), scopes(&["room", "anon"]));
7695
7696 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
7698 assert_eq!(consumer.allowed(), scopes(&["room"]));
7699
7700 for text in ["room", "", "*room", "room/*", "*", "**/room", "room/**/chat", "*.hang"] {
7701 let union = Patterns::from(text.parse::<Pattern>().unwrap());
7702 assert_eq!(producer.scope("", &union).expect(text).allowed(), union, "{text}");
7703 assert_eq!(
7704 producer.consume().scope("", &union).expect(text).allowed(),
7705 union,
7706 "{text}"
7707 );
7708 }
7709
7710 let mixed: Patterns = ["room/**".parse().unwrap(), "other".parse().unwrap()]
7711 .into_iter()
7712 .collect();
7713 assert_eq!(producer.scope("", &mixed).unwrap().allowed(), mixed);
7714 }
7715
7716 #[test]
7717 fn route_table_prunes_to_empty() {
7718 let producer = origin(1).produce();
7719 let consumer = producer.consume();
7720
7721 let cursor = consumer
7724 .scope("", &scopes(&["room/a", "other/deep/head"]))
7725 .unwrap()
7726 .announced();
7727 let route = producer.announce("room/a/b/c", Route::default()).unwrap();
7728 {
7729 let table = producer.shared.lock();
7730 assert!(table.routes.root.find(Path::new("room/a/b/c").parts()).is_some());
7731 assert!(table.routes.root.find(Path::new("other/deep/head").parts()).is_some());
7732 assert_eq!(table.routes.root.cursors_below, 2);
7733 }
7734
7735 drop(route);
7736 drop(cursor);
7737 let table = producer.shared.lock();
7738 assert!(table.routes.root.is_empty());
7739 assert_eq!(table.routes.root.cursors_below, 0);
7740 }
7741
7742 #[test]
7746 fn a_published_broadcast_keeps_the_driver_running() {
7747 let (producer, mut driver) = Producer::new(Config::new(origin(1)));
7748 let waiter = kio::Waiter::noop();
7749 let broadcast = producer.create_broadcast("room/a").unwrap();
7750 drop(producer);
7751 assert!(
7752 driver.poll(Instant::now(), &waiter).is_ok(),
7753 "the broadcast is lifecycle work"
7754 );
7755 drop(broadcast);
7756 assert!(matches!(driver.poll(Instant::now(), &waiter), Err(Error::Closed)));
7757 }
7758
7759 #[tokio::test]
7762 async fn a_dynamic_keeps_the_driver_running() {
7763 let (producer, driver) = Producer::new(Config::new(origin(1)));
7764 let run = tokio::spawn(crate::time::run(driver));
7765 let consumer = producer.consume();
7766 let server = producer.dynamic("room", Route::default()).unwrap();
7767 drop(producer);
7768
7769 let pending = consumer.request_broadcast("room/alice");
7770 let request = queued(&server).await;
7771 let source = broadcast::Info::new().produce();
7772 request.accept(&source);
7773 let resolved = pending.await.expect("the dynamic still serves");
7774
7775 drop(resolved);
7776 drop(source);
7777 drop(server);
7778 tokio::time::timeout(Duration::from_secs(5), run)
7779 .await
7780 .expect("driver must finish once the dynamic is gone")
7781 .unwrap();
7782 drop(consumer);
7783 }
7784
7785 #[test]
7786 fn watch_wakes_only_for_covering_changes() {
7787 let producer = origin(1).produce();
7788 let waiter = kio::Waiter::noop();
7789 let watch = producer.shared.lock().watch(&producer.shared, &Path::new("room/a"));
7790 let seen = watch.seen();
7791
7792 let _other = producer.announce("other", Route::default()).unwrap();
7794 let _below = producer.announce("room/a/b", Route::default()).unwrap();
7795 assert!(watch.poll_changed(&waiter, seen).is_pending());
7796
7797 let above = producer.announce("room", Route::default()).unwrap();
7799 assert!(watch.poll_changed(&waiter, seen).is_ready());
7800 let seen = watch.seen();
7801 drop(above);
7802 assert!(watch.poll_changed(&waiter, seen).is_ready());
7803 let seen = watch.seen();
7804
7805 let _beside = producer.create_broadcast("room/b").unwrap();
7807 assert!(watch.poll_changed(&waiter, seen).is_pending());
7808 let _here = producer.create_broadcast("room/a").unwrap();
7809 assert!(watch.poll_changed(&waiter, seen).is_ready());
7810
7811 drop(watch);
7813 let table = producer.shared.lock();
7814 let node = table
7815 .routes
7816 .root
7817 .find(Path::new("room/a").parts())
7818 .expect("route below keeps the node");
7819 assert!(node.watches.is_empty());
7820 assert_eq!(table.routes.root.watches_below, 0);
7821 }
7822
7823 #[test]
7824 fn a_discarded_front_task_unregisters_its_watch() {
7825 let (producer, _driver) = Producer::new(Config {
7826 hop: origin(1),
7827 ..Default::default()
7828 });
7829 let consumer = producer.consume();
7830 let _served = producer.dynamic("room", Route::default()).unwrap();
7831 drop(producer);
7836 let _pending = consumer.request_broadcast("room/a");
7837 }
7838
7839 #[test]
7840 fn create_broadcast_refuses_a_path_no_pattern_can_spell() {
7841 let producer = origin(1).produce();
7842
7843 assert!(matches!(
7847 producer.create_broadcast("room/*"),
7848 Err(Error::InvalidPath(_))
7849 ));
7850 assert!(matches!(
7851 producer.announce("room/**", Route::default()),
7852 Err(Error::InvalidPath(_))
7853 ));
7854 }
7855
7856 #[test]
7857 fn scope_empty_union_grants_nothing() {
7858 let producer = origin(1).produce();
7859
7860 assert!(matches!(producer.scope("", &Patterns::new()), Err(Error::Unauthorized)));
7862 assert!(matches!(
7863 producer.consume().scope("", &Patterns::new()),
7864 Err(Error::Unauthorized)
7865 ));
7866 }
7867
7868 #[test]
7869 fn scope_nests_and_rebases_roots() {
7870 let producer = origin(1).produce();
7871
7872 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
7874 let nested = scoped.scope("", &scopes(&["room/chat"])).unwrap();
7875 assert_eq!(nested.allowed(), scopes(&["room/chat"]));
7876
7877 assert!(matches!(
7879 scoped.scope("", &scopes(&["other"])),
7880 Err(Error::Unauthorized)
7881 ));
7882
7883 let rooted = nested.scope("room/chat", &Patterns::from(Pattern::all())).unwrap();
7885 assert_eq!(rooted.allowed(), scopes(&[""]));
7886
7887 let broadcast = nested.create_broadcast("room/chat/live").unwrap();
7889 assert!(producer.consume().get_broadcast("room/chat/live").is_some());
7890 broadcast.close();
7891 }
7892
7893 #[test]
7894 fn scope_intersects_and_rebases_arbitrary_grants() {
7895 let producer = origin(1).produce();
7896 let rooms = producer
7897 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
7898 .unwrap();
7899 let chats = rooms
7900 .scope("", &Patterns::from("*/chat".parse::<Pattern>().unwrap()))
7901 .unwrap();
7902 assert_eq!(chats.allowed(), Patterns::from("room/chat".parse::<Pattern>().unwrap()));
7903
7904 let exact = producer
7905 .scope("", &Patterns::from("room/alice".parse::<Pattern>().unwrap()))
7906 .unwrap();
7907 let rooted = exact.scope("room", &Patterns::from(Pattern::all())).unwrap();
7908 assert_eq!(rooted.allowed(), Patterns::from("alice".parse::<Pattern>().unwrap()));
7909 assert!(matches!(
7910 exact.scope("room/bob", &Patterns::from(Pattern::all())),
7911 Err(Error::Unauthorized)
7912 ));
7913
7914 let broadcast = exact.create_broadcast("room/alice").unwrap();
7915 assert!(matches!(
7916 exact.create_broadcast("room/alice/cam"),
7917 Err(Error::Unauthorized)
7918 ));
7919 assert!(producer.consume().get_broadcast("room/alice").is_some());
7920 drop(broadcast);
7921 }
7922
7923 #[tokio::test]
7924 async fn wildcard_scope_filters_announcements_and_reports_captures() {
7925 let producer = origin(1).produce();
7926 let consumer = producer
7927 .consume()
7928 .scope("", &Patterns::from("room/*/chat".parse::<Pattern>().unwrap()))
7929 .unwrap();
7930 let mut announced = consumer.announced();
7931
7932 let alice = producer.create_broadcast("room/alice/chat").unwrap();
7933 alice.announce(Route::default()).unwrap();
7934 let update = announced.try_next().expect("alice's chat");
7935 assert_eq!(update.prefix.as_str(), "room/alice/chat");
7936 assert_eq!(update.captures, Some(vec!["alice".parse::<Pattern>().unwrap()]));
7937
7938 let audio = producer.create_broadcast("room/alice/audio").unwrap();
7939 audio.announce(Route::default()).unwrap();
7940 announced.assert_next_wait();
7941
7942 let broad = producer.announce("room", Route::default()).unwrap();
7943 let update = announced.try_next().expect("overlapping broad route");
7944 assert_eq!(update.prefix.as_str(), "room");
7945 assert_eq!(update.captures, None, "an overlap does not pin the wildcard");
7946
7947 drop(broad);
7948 drop(audio);
7949 drop(alice);
7950 }
7951
7952 #[tokio::test]
7953 async fn local_broadcast_wins_announcement_ties() {
7954 let producer = origin(1).produce();
7955 let remote = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
7956 let local = producer.create_broadcast("room/alice").unwrap();
7957 local.announce(Route::default()).unwrap();
7958
7959 let mut announced = producer.consume().announced();
7960 let update = announced.try_next().expect("one winning route");
7961 assert_eq!(update.prefix.as_str(), "room/alice");
7962 assert_eq!(update.route.cost, Cost::default());
7963 announced.assert_next_wait();
7964
7965 drop(local);
7966 drop(remote);
7967 }
7968
7969 #[test]
7974 fn cost_charge_saturates() {
7975 assert_eq!(Cost { warm: 4, cold: 6 }.charged(5), Cost { warm: 9, cold: 11 });
7976 assert_eq!(Cost::new(u64::MAX).charged(10), Cost::new(MAX_COST));
7977
7978 assert_eq!(Cost::UNKNOWN.charged(3).cold, MAX_COST);
7981 }
7982
7983 fn expiring_origin(expiry: Duration) -> Producer {
7985 let pool = cache::Pool::new(cache::Config::default().with_expiry(expiry));
7986 Config {
7987 pool,
7988 ..Config::default()
7989 }
7990 .produce()
7991 }
7992
7993 #[tokio::test(start_paused = true)]
7997 async fn stalled_publisher_open_group_is_reclaimed() {
7998 let expiry = Duration::from_secs(1);
7999 let origin = expiring_origin(expiry);
8000 let broadcast = origin.create_broadcast("test").unwrap();
8001 let track = broadcast.create_track("video", None).unwrap();
8002
8003 let mut stalled = track.append_group().unwrap();
8004 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
8005 let _successor = track.append_group().unwrap();
8009
8010 let mut reading = stalled.consume();
8011 assert!(reading.read_frame().await.unwrap().is_some());
8012
8013 crate::model::clock::advance(expiry * 2);
8015
8016 let reclaimed = tokio::time::timeout(Duration::from_secs(60), reading.read_frame()).await;
8019 assert!(
8020 matches!(reclaimed, Ok(Err(Error::Old))),
8021 "the sweep must reclaim an idle open group and surface the gap, got {reclaimed:?}"
8022 );
8023 }
8024
8025 #[tokio::test(start_paused = true)]
8028 async fn sweep_respects_a_disabled_expiry() {
8029 let origin = Config {
8030 pool: cache::Pool::unbounded(),
8031 ..Config::default()
8032 }
8033 .produce();
8034 let broadcast = origin.create_broadcast("test").unwrap();
8035 let track = broadcast.create_track("video", None).unwrap();
8036
8037 let mut stalled = track.append_group().unwrap();
8038 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
8039 let _successor = track.append_group().unwrap();
8040
8041 let mut reading = stalled.consume();
8042 assert!(reading.read_frame().await.unwrap().is_some());
8043
8044 crate::model::clock::advance(Duration::from_secs(3600));
8045 tokio::time::advance(Duration::from_secs(3600)).await;
8046
8047 assert!(
8048 reading.read_frame().now_or_never().is_none(),
8049 "a pool without an expiry window never reclaims"
8050 );
8051 }
8052
8053 #[test]
8056 fn drain_cost_is_encodable() {
8057 use crate::coding::Encode;
8058
8059 let mut buf = Vec::new();
8060 Cost::DRAIN
8061 .encode(&mut buf, crate::lite::Version::Lite06)
8062 .expect("a draining route is still forwarded, so its cost must encode");
8063 }
8064}