1use crate::{broadcast, cache, 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 runtime::{Instant, Timers},
23 time::Clock,
24 util::{Keepalive, TaskSet, Tasks, TasksWeak},
25};
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
37pub struct Hop {
38 id: u64,
40}
41
42impl Hop {
43 pub const UNKNOWN: Self = Self { id: 0 };
50
51 pub fn new(id: u64) -> Result<Self, InvalidHop> {
57 if id == 0 || id >= 1u64 << 62 {
58 return Err(InvalidHop::Range);
59 }
60 Ok(Self { id })
61 }
62
63 pub fn random() -> Self {
70 let mut rng = rand::rng();
71 let id = rng.random_range(1..(1u64 << 53));
72 Self { id }
73 }
74
75 pub fn id(self) -> u64 {
77 self.id
78 }
79
80 pub(crate) fn from_wire(id: u64) -> Result<Self, DecodeError> {
82 if id >= 1u64 << 62 {
83 return Err(DecodeError::InvalidValue);
84 }
85 Ok(Self { id })
86 }
87}
88
89#[derive(Clone, Debug)]
95#[non_exhaustive]
96pub struct Config {
97 pub hop: Hop,
100
101 pub pool: cache::Pool,
108
109 pub cache_duration: Duration,
117
118 pub default_max_age: Duration,
128}
129
130impl Default for Config {
131 fn default() -> Self {
133 let pool = cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY));
134 Self {
135 hop: Hop::random(),
136 pool,
137 cache_duration: Duration::MAX,
138 default_max_age: track::DEFAULT_MAX_AGE,
139 }
140 }
141}
142
143impl Config {
144 pub fn new(hop: Hop) -> Self {
146 Self { hop, ..Self::default() }
147 }
148}
149
150impl From<Hop> for Config {
151 fn from(hop: Hop) -> Self {
153 Self::new(hop)
154 }
155}
156
157impl TryFrom<u64> for Hop {
158 type Error = InvalidHop;
159
160 fn try_from(id: u64) -> Result<Self, Self::Error> {
161 Self::new(id)
162 }
163}
164
165impl fmt::Display for Hop {
166 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
167 self.id.fmt(f)
168 }
169}
170
171impl<V: Copy> Encode<V> for Hop
172where
173 u64: Encode<V>,
174{
175 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
176 self.id.encode(w, version)
177 }
178}
179
180impl<V: Copy> Decode<V> for Hop
181where
182 u64: Decode<V>,
183{
184 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
185 Self::from_wire(u64::decode(r, version)?)
186 }
187}
188
189pub(crate) const MAX_HOPS: usize = 32;
195
196#[derive(Debug, Clone, Default, PartialEq, Eq)]
204pub struct Hops(Vec<Hop>);
205
206#[derive(Debug, Clone, Copy, PartialEq, Eq)]
208#[non_exhaustive]
209pub enum InvalidHop {
210 Range,
213
214 TooMany,
217
218 Duplicate,
222}
223
224impl fmt::Display for InvalidHop {
225 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
226 match self {
227 Self::Range => write!(f, "local hop id must be non-zero and below 2^62"),
228 Self::TooMany => write!(f, "too many hops (max {MAX_HOPS})"),
229 Self::Duplicate => write!(f, "hop already in the chain"),
230 }
231 }
232}
233
234impl std::error::Error for InvalidHop {}
235
236impl From<InvalidHop> for DecodeError {
237 fn from(err: InvalidHop) -> Self {
238 match err {
239 InvalidHop::TooMany => DecodeError::BoundsExceeded,
240 InvalidHop::Range | InvalidHop::Duplicate => DecodeError::InvalidValue,
241 }
242 }
243}
244
245impl Hops {
246 pub fn new() -> Self {
248 Self(Vec::new())
249 }
250
251 pub fn push(&mut self, hop: Hop) -> Result<(), InvalidHop> {
257 if self.0.len() >= MAX_HOPS {
258 return Err(InvalidHop::TooMany);
259 }
260 if hop != Hop::UNKNOWN && self.0.contains(&hop) {
261 return Err(InvalidHop::Duplicate);
262 }
263 self.0.push(hop);
264 Ok(())
265 }
266
267 pub fn contains(&self, hop: &Hop) -> bool {
269 self.0.contains(hop)
270 }
271
272 pub fn len(&self) -> usize {
274 self.0.len()
275 }
276
277 pub fn is_empty(&self) -> bool {
279 self.0.is_empty()
280 }
281
282 pub fn iter(&self) -> std::slice::Iter<'_, Hop> {
284 self.0.iter()
285 }
286
287 pub fn as_slice(&self) -> &[Hop] {
289 &self.0
290 }
291}
292
293impl TryFrom<Vec<Hop>> for Hops {
294 type Error = InvalidHop;
295
296 fn try_from(v: Vec<Hop>) -> Result<Self, Self::Error> {
297 if v.len() > MAX_HOPS {
298 return Err(InvalidHop::TooMany);
299 }
300 for (i, hop) in v.iter().enumerate() {
302 if *hop != Hop::UNKNOWN && v[i + 1..].contains(hop) {
303 return Err(InvalidHop::Duplicate);
304 }
305 }
306 Ok(Self(v))
307 }
308}
309
310impl<'a> IntoIterator for &'a Hops {
311 type Item = &'a Hop;
312 type IntoIter = std::slice::Iter<'a, Hop>;
313
314 fn into_iter(self) -> Self::IntoIter {
315 self.iter()
316 }
317}
318
319impl<V: Copy> Encode<V> for Hops
320where
321 u64: Encode<V>,
322 Hop: Encode<V>,
323{
324 fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
325 (self.0.len() as u64).encode(w, version)?;
326 for origin in &self.0 {
327 origin.encode(w, version)?;
328 }
329 Ok(())
330 }
331}
332
333impl<V: Copy> Decode<V> for Hops
334where
335 u64: Decode<V>,
336 Hop: Decode<V>,
337{
338 fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
339 let count = u64::decode(r, version)? as usize;
340 if count > MAX_HOPS {
341 return Err(DecodeError::BoundsExceeded);
342 }
343 let mut list = Self(Vec::with_capacity(count));
346 for _ in 0..count {
347 list.push(Hop::decode(r, version)?)?;
348 }
349 Ok(list)
350 }
351}
352
353const MAX_COST: u64 = (1 << 62) - 1;
360
361#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord)]
370pub struct Cost {
371 pub warm: u64,
380
381 pub cold: u64,
388}
389
390impl Cost {
391 pub const fn new(cost: u64) -> Self {
394 Self { warm: cost, cold: cost }
395 }
396
397 pub const MAX: Self = Self::new(MAX_COST);
404
405 pub const DRAIN: Self = Self::MAX;
408
409 pub(crate) const UNKNOWN: Self = Self {
413 warm: 0,
414 cold: MAX_COST,
415 };
416
417 pub(crate) fn charged(self, link_cost: u64) -> Self {
420 Self {
421 warm: self.warm.saturating_add(link_cost).min(MAX_COST),
422 cold: self.cold.saturating_add(link_cost).min(MAX_COST),
423 }
424 }
425
426 pub(crate) fn clamped(self) -> Self {
429 Self {
430 warm: self.warm.min(MAX_COST),
431 cold: self.cold.min(MAX_COST),
432 }
433 }
434}
435
436impl From<u64> for Cost {
437 fn from(cost: u64) -> Self {
438 Self::new(cost)
439 }
440}
441
442#[derive(Clone, Debug, PartialEq, Eq)]
453#[non_exhaustive]
454pub struct Route {
455 pub hops: Hops,
460
461 pub cost: Cost,
465
466 pub(crate) via: Hop,
472}
473
474impl Default for Route {
475 fn default() -> Self {
476 Self {
477 hops: Hops::new(),
478 cost: Cost::default(),
479 via: Hop::UNKNOWN,
480 }
481 }
482}
483
484impl Route {
485 pub fn with_hops(mut self, hops: Hops) -> Self {
487 self.hops = hops;
488 self
489 }
490
491 pub fn with_cost(mut self, cost: impl Into<Cost>) -> Self {
496 self.cost = cost.into();
497 self
498 }
499
500 pub(crate) fn with_via(mut self, via: Hop) -> Self {
505 self.via = via;
506 self
507 }
508
509 pub fn is_anonymous(&self) -> bool {
517 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
518 }
519}
520
521static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
522
523#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
524struct ConsumerId(u64);
525
526impl ConsumerId {
527 fn new() -> Self {
528 Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
529 }
530}
531
532fn fnv_key(name: &str, origins: impl IntoIterator<Item = Hop>) -> u64 {
541 const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
543
544 let mut hash = SEED;
545 for &byte in name.as_bytes() {
546 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
547 }
548 for origin in origins {
549 for &byte in &origin.id().to_le_bytes() {
550 hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
551 }
552 }
553
554 hash
555}
556
557fn route_order(prefix: &Path, entry: &RouteEntry) -> (bool, Cost, usize, u64, Reverse<u64>) {
564 (
565 entry.is_anonymous(),
566 entry.cost,
567 entry.hops.len(),
568 fnv_key(prefix.as_str(), entry.hops.iter().copied()),
569 Reverse(entry.id),
570 )
571}
572
573type RouteMeta = (Hops, Cost);
575
576type AnnounceMeta = (RouteMeta, Option<Vec<Pattern>>);
583
584enum PendingUpdate {
585 Announce(AnnounceMeta),
586 Unannounce(AnnounceMeta),
587 UnannounceAnnounce { old: AnnounceMeta, new: AnnounceMeta },
588}
589
590#[derive(Default)]
595struct OriginConsumerState {
596 pending: BTreeMap<PathOwned, PendingUpdate>,
597 delivered: BTreeSet<PathOwned>,
602 ended: bool,
605}
606
607impl OriginConsumerState {
608 fn apply_announce(&mut self, prefix: PathOwned, meta: RouteMeta, captures: Option<Vec<Pattern>>) {
609 let meta = (meta, captures);
610 let new = match self.pending.remove(&prefix) {
611 None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(meta),
613 Some(PendingUpdate::Unannounce(old) | PendingUpdate::UnannounceAnnounce { old, .. }) => {
615 PendingUpdate::UnannounceAnnounce { old, new: meta }
616 }
617 };
618 self.pending.insert(prefix, new);
619 }
620
621 fn apply_unannounce(&mut self, prefix: PathOwned, last: RouteMeta, captures: Option<Vec<Pattern>>) {
622 let last = (last, captures);
623 match self.pending.remove(&prefix) {
624 Some(PendingUpdate::Announce(_)) if !self.delivered.contains(&prefix) => {}
627 None | Some(PendingUpdate::Announce(_) | PendingUpdate::Unannounce(_)) => {
630 self.pending.insert(prefix, PendingUpdate::Unannounce(last));
631 }
632 Some(PendingUpdate::UnannounceAnnounce { old, .. }) => {
635 self.pending.insert(prefix, PendingUpdate::Unannounce(old));
636 }
637 }
638 }
639
640 fn take(&mut self) -> Option<AnnounceUpdate> {
642 let prefix = self.pending.keys().next()?.clone();
643 let ((meta, captures), kind) = match self.pending.remove(&prefix).unwrap() {
644 PendingUpdate::Announce(meta) => {
645 let kind = match self.delivered.insert(prefix.clone()) {
647 true => AnnounceKind::Announced,
648 false => AnnounceKind::Updated,
649 };
650 (meta, kind)
651 }
652 PendingUpdate::Unannounce(meta) => {
653 self.delivered.remove(&prefix);
654 (meta, AnnounceKind::Retracted)
655 }
656 PendingUpdate::UnannounceAnnounce { old, new } => {
657 self.delivered.remove(&prefix);
660 self.pending.insert(prefix.clone(), PendingUpdate::Announce(new));
661 (old, AnnounceKind::Retracted)
662 }
663 };
664 Some(AnnounceUpdate {
665 prefix,
666 captures,
667 route: Route {
668 hops: meta.0,
669 cost: meta.1,
670 via: Hop::UNKNOWN,
671 },
672 kind,
673 })
674 }
675}
676
677struct RouteEntry {
679 id: u64,
680 prefix: PathOwned,
681 scope: Patterns,
684 hops: Hops,
685 cost: Cost,
686 via: Hop,
690 local: bool,
693 server: Option<kio::Shared<ServeState>>,
697 source: Option<broadcast::Consumer>,
701 advertised: bool,
704 claim: Pattern,
710}
711
712impl RouteEntry {
713 fn is_anonymous(&self) -> bool {
714 self.hops.iter().any(|hop| *hop == Hop::UNKNOWN)
715 }
716
717 fn serves(&self, path: &Path) -> bool {
721 self.server.is_some() || (self.source.is_some() && self.prefix == *path)
722 }
723
724 fn qualifies(&self, pin: Pin) -> bool {
726 match pin {
727 Pin::Any => true,
728 Pin::Local => self.local,
729 Pin::Publisher(first) => self.hops.iter().next() == Some(&first),
730 Pin::None => false,
731 }
732 }
733
734 fn visible_to(&self, exclude: Option<Hop>) -> bool {
739 match exclude {
740 Some(peer) if peer != Hop::UNKNOWN => self.via != peer && !self.hops.contains(&peer),
741 _ => true,
742 }
743 }
744
745 fn overlaps(&self, allowed: &Patterns) -> bool {
747 self.scope.iter().any(|scope| {
748 scope
749 .intersect(&self.claim)
750 .is_ok_and(|scoped| scoped.iter().any(|restriction| allowed.overlaps(restriction)))
751 })
752 }
753}
754
755fn prefix_claim(prefix: &Path) -> Result<Pattern, InvalidPattern> {
757 if prefix.parts().count() == Path::MAX_PARTS {
758 Pattern::literal(prefix.as_str())
759 } else {
760 Pattern::subtree(prefix.as_str())
761 }
762}
763
764#[derive(Default)]
769struct ServeState {
770 requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
773
774 served: WeakCache<PathOwned, broadcast::WeakConsumer>,
780
781 closed: bool,
784}
785
786type FrontKey = (PathOwned, Option<Hop>);
791
792#[derive(Clone)]
795struct RemoteFront {
796 request: kio::Producer<PendingBroadcast>,
800 broadcast: broadcast::WeakConsumer,
803}
804
805type CursorRoute = (u64, RouteMeta, bool, Option<Vec<Pattern>>);
807
808impl WeakEntry for RemoteFront {
809 fn is_closed(&self) -> bool {
810 self.broadcast.is_closed()
811 }
812
813 fn same_channel(&self, other: &Self) -> bool {
814 self.broadcast.same_channel(&other.broadcast)
815 }
816}
817
818struct TableCursor {
821 root: PathOwned,
823 allowed: Patterns,
825 heads: Vec<PathOwned>,
829 exclude: Option<Hop>,
832 state: kio::Producer<OriginConsumerState>,
834 current: HashMap<PathOwned, CursorRoute>,
839}
840
841impl TableCursor {
842 fn presented(&self, prefix: &Path, claim: &Pattern) -> Option<PathOwned> {
848 if !self.allowed.overlaps(claim) {
849 return None;
850 }
851
852 if let Some(relative) = prefix.strip_prefix(&self.root) {
853 return Some(relative.to_owned());
854 }
855 self.root.has_prefix(prefix).then(PathOwned::default)
856 }
857
858 fn captures(&self, prefix: &Path) -> Option<Vec<Pattern>> {
861 let literal = Pattern::literal(prefix.as_str()).ok()?;
862 self.allowed
863 .iter()
864 .filter_map(|allowed| {
865 allowed
866 .captures(&literal)
867 .map(|captures| (allowed.specificity(), captures))
868 })
869 .max_by_key(|(specificity, _)| *specificity)
870 .map(|(_, captures)| captures)
871 }
872
873 fn visible(&self, entry: &RouteEntry) -> bool {
877 (entry.advertised || (entry.local && self.exclude.is_none()))
878 && entry.visible_to(self.exclude)
879 && entry.overlaps(&self.allowed)
880 }
881}
882
883#[derive(Clone)]
885struct OriginScope {
886 allowed: Patterns,
888}
889
890impl OriginScope {
891 fn empty() -> Self {
893 Self {
894 allowed: Patterns::new(),
895 }
896 }
897
898 fn narrow(&self, patterns: &Patterns) -> Option<Self> {
900 let allowed = self.allowed.intersect(patterns).ok()?;
901 if allowed.is_empty() {
902 None
903 } else {
904 Some(Self { allowed })
905 }
906 }
907
908 fn permits(&self, path: &Path) -> bool {
910 self.allowed.matches(path.as_str())
911 }
912
913 fn relative(&self, root: &Path) -> Patterns {
915 self.allowed.rebase(root.as_str())
916 }
917}
918
919impl Default for OriginScope {
920 fn default() -> Self {
921 Self {
922 allowed: Patterns::from(Pattern::all()),
923 }
924 }
925}
926
927pub(crate) fn interest_prefixes(allowed: &Patterns) -> Vec<PathOwned> {
930 let mut heads: Vec<PathOwned> = allowed
931 .iter()
932 .map(|pattern| Path::new(pattern.head()).to_owned())
933 .collect();
934 heads.sort();
935 heads.dedup();
936 let covered = heads.clone();
937 heads.retain(|head| !covered.iter().any(|other| other != head && head.has_prefix(other)));
938 heads
939}
940
941#[derive(Clone, Copy, Debug, PartialEq, Eq)]
943pub enum AnnounceKind {
944 Announced,
946 Updated,
948 Retracted,
950}
951
952impl AnnounceKind {
953 pub fn is_active(self) -> bool {
955 !matches!(self, Self::Retracted)
956 }
957}
958
959#[derive(Clone, Debug)]
967pub struct AnnounceUpdate {
968 pub prefix: PathOwned,
970 pub captures: Option<Vec<Pattern>>,
973 pub route: Route,
976 pub kind: AnnounceKind,
978}
979
980#[derive(Clone)]
982pub struct Producer {
983 hop: Hop,
986
987 scope: OriginScope,
989
990 root: PathOwned,
992
993 shared: kio::Shared<OriginState>,
996
997 pool: cache::Pool,
1000
1001 cache_duration: Duration,
1004
1005 default_max_age: Duration,
1008
1009 stats: stats::Session,
1013
1014 tasks: Tasks,
1018
1019 timers: Clock,
1021}
1022
1023impl Producer {
1024 pub fn new(config: Config) -> (Self, Driver) {
1031 let (tasks, set) = TaskSet::new();
1032 let scope = OriginScope::default();
1033 let shared = kio::Shared::<OriginState>::default();
1034 let timers = Clock::default();
1035 let pool = config.pool.clone();
1036 let producer = Self {
1037 hop: config.hop,
1038 scope: scope.clone(),
1039 root: PathOwned::default(),
1040 shared: shared.clone(),
1041 pool: config.pool,
1042 cache_duration: config.cache_duration,
1043 default_max_age: config.default_max_age,
1044 stats: stats::Session::default(),
1045 tasks,
1046 timers: timers.clone(),
1047 };
1048 let driver = Driver {
1049 state: DriverState {
1050 set,
1051 shared,
1052 done: false,
1053 },
1054 timers,
1055 pool,
1056 };
1057 (producer, driver)
1058 }
1059
1060 pub fn with_stats(mut self, session: stats::Session) -> Self {
1064 self.stats = session;
1065 self
1066 }
1067
1068 pub fn config(&self) -> Config {
1070 Config {
1071 hop: self.hop,
1072 pool: self.pool.clone(),
1073 cache_duration: self.cache_duration,
1074 default_max_age: self.default_max_age,
1075 }
1076 }
1077
1078 pub fn hop(&self) -> Hop {
1080 self.hop
1081 }
1082
1083 pub(crate) fn default_max_age(&self) -> Duration {
1086 self.default_max_age
1087 }
1088
1089 pub(crate) fn empty(hop: Hop) -> Self {
1094 let (tasks, _) = TaskSet::new();
1097 Self {
1098 hop,
1099 scope: OriginScope::empty(),
1100 root: PathOwned::default(),
1101 shared: kio::Shared::default(),
1102 pool: cache::Pool::default(),
1103 cache_duration: Duration::MAX,
1104 default_max_age: track::DEFAULT_MAX_AGE,
1105 stats: stats::Session::default(),
1106 tasks,
1107 timers: Clock::default(),
1108 }
1109 }
1110
1111 pub fn create_broadcast(&self, path: impl AsPath) -> Result<broadcast::Producer, Error> {
1142 let path = path.as_path();
1143
1144 let full = self.root.join(&path).to_owned();
1145 if !self.scope.permits(&full) {
1146 return Err(Error::Unauthorized);
1147 }
1148 if full.parts().count() > Path::MAX_PARTS {
1152 return Err(BoundsExceeded.into());
1153 }
1154 let claim = prefix_claim(&full)?;
1157
1158 let ingress = self.stats.ingress(&full);
1160
1161 let announcing = Announcing {
1166 hop: self.hop,
1167 shared: self.shared.clone(),
1168 requested: full.clone(),
1169 prefixes: vec![(full.clone(), claim)],
1170 scope: self.scope.allowed.clone(),
1171 local: true,
1172 stats: self.stats.clone(),
1173 };
1174 let info = broadcast::Info {
1175 pool: self.pool.clone(),
1176 cache_duration: self.cache_duration,
1177 path: full,
1178 };
1179 let source = info.produce().with_stats(ingress);
1180 let entry = announcing.announce(
1181 Route::default(),
1182 Serving {
1183 server: None,
1184 source: Some(source.consume()),
1185 advertised: false,
1186 },
1187 )?;
1188 Ok(source.with_announcer(Announcer {
1189 entry,
1190 _keepalive: self.tasks.keepalive(),
1191 }))
1192 }
1193
1194 pub fn publish(&self, path: impl AsPath, route: Route) -> Result<broadcast::Producer, Error> {
1196 let broadcast = self.create_broadcast(path)?;
1197 broadcast.announce(route)?;
1198 Ok(broadcast)
1199 }
1200
1201 pub(crate) fn create_source(&self, path: impl AsPath) -> broadcast::Producer {
1207 let path = path.as_path();
1208 let full = self.root.join(&path).to_owned();
1209 let ingress = self.stats.ingress(&full);
1210 broadcast::Info {
1211 pool: self.pool.clone(),
1212 cache_duration: self.cache_duration,
1213 path: full,
1214 }
1215 .produce()
1216 .with_stats(ingress)
1217 }
1218
1219 #[cfg(test)]
1228 pub(crate) fn announce(&self, prefix: impl AsPath, route: Route) -> Result<AnnounceProducer, Error> {
1229 Announcing::new(self, prefix)?.announce(
1230 route,
1231 Serving {
1232 server: None,
1233 source: None,
1234 advertised: true,
1235 },
1236 )
1237 }
1238
1239 pub fn dynamic(&self, prefix: impl AsPath, route: Route) -> Result<Dynamic, Error> {
1261 let announcing = Announcing::new(self, prefix)?;
1262 let serve = kio::Shared::<ServeState>::default();
1263 serve.lock().requests.add_handler();
1264 let announcement = announcing.announce(
1265 route,
1266 Serving {
1267 server: Some(serve.clone()),
1268 source: None,
1269 advertised: true,
1270 },
1271 )?;
1272 Ok(Dynamic {
1273 announcement,
1274 state: serve,
1275 })
1276 }
1277
1278 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Producer, Error> {
1285 let root = self.root.join(root).to_owned();
1286 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
1287 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
1288 Ok(Producer {
1289 hop: self.hop,
1290 scope,
1291 root,
1292 shared: self.shared.clone(),
1293 pool: self.pool.clone(),
1294 cache_duration: self.cache_duration,
1295 default_max_age: self.default_max_age,
1296 stats: self.stats.clone(),
1297 tasks: self.tasks.clone(),
1298 timers: self.timers.clone(),
1299 })
1300 }
1301
1302 pub fn consume(&self) -> Consumer {
1307 Consumer::from_producer(self, stats::Session::default())
1310 }
1311
1312 pub fn root(&self) -> &Path<'_> {
1314 &self.root
1315 }
1316
1317 pub fn allowed(&self) -> Patterns {
1319 self.scope.relative(&self.root)
1320 }
1321
1322 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
1324 self.root.join(path)
1325 }
1326}
1327
1328struct Announcing {
1332 hop: Hop,
1333 shared: kio::Shared<OriginState>,
1334 requested: PathOwned,
1336 prefixes: Vec<(PathOwned, Pattern)>,
1340 scope: Patterns,
1342 local: bool,
1343 stats: stats::Session,
1344}
1345
1346impl Announcing {
1347 fn new(producer: &Producer, prefix: impl AsPath) -> Result<Self, Error> {
1349 let requested = producer.root.join(prefix.as_path()).to_owned();
1350 if requested.parts().count() > Path::MAX_PARTS {
1351 return Err(BoundsExceeded.into());
1352 }
1353 let claim = prefix_claim(&requested)?;
1354 if !producer.scope.allowed.overlaps(&claim) {
1355 return Err(Error::Unauthorized);
1356 }
1357 Ok(Self {
1358 hop: producer.hop,
1359 shared: producer.shared.clone(),
1360 requested: requested.clone(),
1361 prefixes: vec![(requested, claim)],
1362 scope: producer.scope.allowed.clone(),
1363 local: false,
1364 stats: producer.stats.clone(),
1365 })
1366 }
1367
1368 fn announce(&self, route: Route, serving: Serving) -> Result<AnnounceProducer, Error> {
1369 debug_assert!(
1370 !route.hops.contains(&self.hop),
1371 "announce called with a looping hop chain",
1372 );
1373
1374 let via = route.via;
1375 let meta: RouteMeta = (route.hops, route.cost);
1376
1377 let mut shared = self.shared.lock();
1378 if shared.closed {
1379 return Err(Error::Closed);
1380 }
1381
1382 let mut entries = Vec::with_capacity(self.prefixes.len());
1383 for (prefix, claim) in &self.prefixes {
1384 let id = shared.next_route;
1385 shared.next_route += 1;
1386 shared.routes.insert(RouteEntry {
1387 id,
1388 prefix: prefix.clone(),
1389 scope: self.scope.clone(),
1390 hops: meta.0.clone(),
1391 cost: meta.1,
1392 via,
1393 local: self.local,
1394 server: serving.server.clone(),
1395 source: serving.source.clone(),
1396 advertised: serving.advertised,
1397 claim: claim.clone(),
1398 });
1399 shared.sync_route(prefix, claim);
1400 entries.push((prefix.clone(), id));
1401 }
1402 drop(shared);
1403
1404 let guard = self.stats.ingress(&self.requested).announce();
1406
1407 Ok(AnnounceProducer {
1408 shared: self.shared.clone(),
1409 entries,
1410 _guard: guard,
1411 })
1412 }
1413}
1414
1415struct Serving {
1417 server: Option<kio::Shared<ServeState>>,
1418 source: Option<broadcast::Consumer>,
1419 advertised: bool,
1420}
1421
1422pub(crate) struct Announcer {
1429 entry: AnnounceProducer,
1430 _keepalive: Keepalive,
1434}
1435
1436impl Announcer {
1437 pub(crate) fn announce(&mut self, route: Route) -> Result<(), Error> {
1439 self.entry.update(route)
1440 }
1441
1442 pub(crate) fn withdraw(&mut self) {
1444 self.entry.withdraw();
1445 }
1446}
1447
1448#[must_use = "dropping an announcement retracts the route"]
1455pub(crate) struct AnnounceProducer {
1456 shared: kio::Shared<OriginState>,
1457 entries: Vec<(PathOwned, u64)>,
1460 _guard: stats::Announce,
1462}
1463
1464impl AnnounceProducer {
1465 pub fn update(&self, route: Route) -> Result<(), Error> {
1473 let mut shared = self.shared.lock();
1474 if shared.closed {
1475 return Err(Error::Closed);
1476 }
1477 for (prefix, id) in &self.entries {
1478 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1480 return Err(Error::Closed);
1481 };
1482 entry.hops = route.hops.clone();
1483 entry.cost = route.cost;
1484 entry.via = route.via;
1485 entry.advertised = true;
1486 let claim = entry.claim.clone();
1487 shared.sync_route(prefix, &claim);
1488 }
1489 Ok(())
1490 }
1491
1492 fn withdraw(&self) {
1494 let mut shared = self.shared.lock();
1495 for (prefix, id) in &self.entries {
1496 let Some(entry) = shared.routes.entry_mut(prefix, *id) else {
1497 continue;
1498 };
1499 if !entry.advertised {
1500 continue;
1501 }
1502 entry.advertised = false;
1503 entry.hops = Hops::default();
1504 entry.cost = Cost::default();
1505 entry.via = Hop::UNKNOWN;
1506 let claim = entry.claim.clone();
1507 shared.sync_route(prefix, &claim);
1508 }
1509 }
1510
1511 fn retract(&self) {
1514 let mut shared = self.shared.lock();
1515 for (prefix, id) in &self.entries {
1516 let Some(entry) = shared.routes.remove(prefix, *id) else {
1517 continue;
1518 };
1519 if let Some(server) = &entry.server {
1522 let mut server = server.lock();
1523 server.closed = true;
1524 for producer in server.requests.drain_all() {
1525 if let Ok(mut request) = producer.write() {
1526 request.resolved.get_or_insert(Err(Error::Unroutable));
1527 }
1528 }
1529 }
1530 shared.sync_route(&entry.prefix, &entry.claim);
1531 }
1532 }
1533}
1534
1535impl Drop for AnnounceProducer {
1536 fn drop(&mut self) {
1537 self.retract();
1538 }
1539}
1540
1541#[must_use = "poll the driver or the origin makes no progress"]
1553pub struct Driver {
1554 state: DriverState,
1555 timers: Clock,
1557 pool: cache::Pool,
1560}
1561
1562struct DriverState {
1564 set: TaskSet,
1566 shared: kio::Shared<OriginState>,
1569 done: bool,
1571}
1572
1573impl Driver {
1574 pub fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
1580 self.timers.advance(now);
1581 let result = self.state.poll(waiter);
1582 let gc = self.pool.gc(now);
1583 if result.is_ready() {
1584 return Err(Error::Closed);
1585 }
1586 Ok(self.timers.timeout().into_iter().chain(gc).min())
1587 }
1588}
1589
1590impl crate::time::Driver for Driver {
1591 fn poll(&mut self, now: Instant, waiter: &kio::Waiter) -> Result<Option<Instant>, Error> {
1592 self.poll(now, waiter)
1593 }
1594}
1595
1596impl DriverState {
1597 fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
1598 if !self.done {
1602 ready!(self.set.poll(waiter));
1603 self.done = true;
1604 }
1605 Poll::Ready(())
1606 }
1607
1608 fn teardown(&mut self) {
1612 drop(std::mem::replace(&mut self.set, TaskSet::owned()));
1615
1616 let (servers, cursors, fronts) = {
1621 let mut shared = self.shared.lock();
1622 shared.closed = true;
1623 shared.routes.poke_all();
1625 let servers: Vec<_> = shared
1626 .routes
1627 .entries()
1628 .filter_map(|entry| entry.server.clone())
1629 .collect();
1630 let cursors: Vec<_> = shared.cursors.values().map(|cursor| cursor.state.clone()).collect();
1631 let fronts: Vec<_> = shared.fronts.values().map(|front| front.request.clone()).collect();
1632 (servers, cursors, fronts)
1633 };
1634
1635 for producer in fronts {
1638 if let Ok(mut request) = producer.write() {
1639 request.resolved.get_or_insert(Err(Error::Dropped));
1640 }
1641 }
1642 for server in servers {
1646 let mut server = server.lock();
1647 server.closed = true;
1648 for producer in server.requests.drain_all() {
1649 if let Ok(mut request) = producer.write() {
1650 request.resolved.get_or_insert(Err(Error::Dropped));
1651 }
1652 }
1653 }
1654
1655 for state in cursors {
1658 if let Ok(mut state) = state.write() {
1659 state.ended = true;
1660 }
1661 }
1662 }
1663}
1664
1665impl Drop for DriverState {
1666 fn drop(&mut self) {
1667 self.teardown();
1668 }
1669}
1670
1671const TRACK_IDLE_LINGER: Duration = Duration::from_secs(30);
1684
1685struct WarmCopy {
1690 track: track::Producer,
1691 _dynamic: track::Dynamic,
1692}
1693
1694impl Drop for WarmCopy {
1695 fn drop(&mut self) {
1696 let _ = self.track.finish();
1697 }
1698}
1699
1700fn warm_copy(source: &track::Consumer) -> Option<WarmCopy> {
1702 let info = source.cached_info()?;
1703 let mut track = track::Producer::new(Arc::new(source.broadcast().clone()), source.name(), info);
1704 for (group, visible) in source.cached_groups() {
1705 if group.is_finished() {
1709 let _ = track.adopt_group(group, visible);
1710 }
1711 }
1712 let dynamic = track.dynamic();
1713 Some(WarmCopy {
1714 track,
1715 _dynamic: dynamic,
1716 })
1717}
1718
1719struct FrontTask {
1721 shared: kio::Shared<OriginState>,
1723 broadcast: broadcast::Producer,
1725 path: PathOwned,
1727 exclude: Option<Hop>,
1729 watch: Watch,
1731 request: kio::Producer<PendingBroadcast>,
1733 timers: Clock,
1734}
1735
1736struct TrackIo {
1739 resume: super::resume::Producer,
1740 staged: Option<(u64, track::Consumer)>,
1742 query: Option<(u64, track::Consumer, track::Querying)>,
1744 copy: Option<(u64, track::Consumer)>,
1746 edge: Option<track::Position>,
1751 warm: Option<WarmCopy>,
1754 used: bool,
1756}
1757
1758async fn run_front(task: FrontTask) {
1762 let FrontTask {
1763 shared,
1764 broadcast,
1765 path,
1766 exclude,
1767 watch,
1768 request,
1769 timers,
1770 } = task;
1771
1772 enum Step {
1774 Assigned(Arc<str>, super::resume::Producer),
1775 Resolved(u64, Result<broadcast::Consumer, Error>),
1776 SourceClosed(u64),
1777 Info(Arc<str>, u64, Result<track::Info, Error>),
1778 Ended(Arc<str>, u64, Result<(), Error>),
1779 Demand(Arc<str>),
1780 Deadline,
1781 Table,
1782 }
1783
1784 let mut front = Front::new(TRACK_IDLE_LINGER);
1785 let mut sources: HashMap<u64, broadcast::Consumer> = HashMap::new();
1786 let mut next_source = 0u64;
1787 let mut upstream: Option<(u64, kio::Consumer<PendingBroadcast>)> = None;
1789 let mut tracks: HashMap<Arc<str>, TrackIo> = HashMap::new();
1790 let mut deadline = crate::runtime::Deadline::new(&timers);
1791 let mut seen = 0;
1793 let mut events: VecDeque<Event> = VecDeque::new();
1794
1795 let select = |front: &mut Front, sources: &HashMap<u64, broadcast::Consumer>, seen: &mut u64| -> Event {
1798 let table = shared.read();
1799 if table.closed {
1800 return Event::Closed;
1801 }
1802 *seen = watch.seen();
1804 front.retain_routes(|route| table.routes.covers(&path.as_path(), route));
1805 let best = table
1806 .best_route(&path.as_path(), exclude, front.pin(), front.refused_routes())
1807 .map(|entry| Candidate {
1808 route: entry.id,
1809 first: entry.hops.iter().next().copied(),
1810 local: entry.local,
1811 });
1812 let serving_closing = front
1813 .serving()
1814 .and_then(|id| sources.get(&id))
1815 .is_some_and(|source| source.is_closing());
1816 Event::Selected { best, serving_closing }
1817 };
1818
1819 events.push_back(select(&mut front, &sources, &mut seen));
1820
1821 loop {
1822 while let Some(event) = events.pop_front() {
1823 for action in front.step(event) {
1824 match action {
1825 Action::Reselect => events.push_back(select(&mut front, &sources, &mut seen)),
1826 Action::Request { route } => {
1827 let found = {
1829 let table = shared.read();
1830 table
1831 .routes
1832 .covering(&path.as_path())
1833 .find(|entry| entry.id == route)
1834 .map(|entry| {
1835 (
1836 Candidate {
1837 route,
1838 first: entry.hops.iter().next().copied(),
1839 local: entry.local,
1840 },
1841 entry.source.clone(),
1842 entry.server.clone(),
1843 )
1844 })
1845 };
1846 let Some((candidate, source, server)) = found else {
1847 events.push_back(Event::Resolved {
1848 route,
1849 result: Err(Refusal {
1850 err: Error::Unroutable,
1851 standing: false,
1852 }),
1853 });
1854 continue;
1855 };
1856 front.identify(candidate);
1857 if let Some(source) = source {
1858 let id = next_source;
1859 next_source += 1;
1860 sources.insert(id, source);
1861 events.push_back(Event::Resolved { route, result: Ok(id) });
1862 continue;
1863 }
1864 let Some(server) = server else {
1865 events.push_back(Event::Resolved {
1866 route,
1867 result: Err(Refusal {
1868 err: Error::Unroutable,
1869 standing: true,
1870 }),
1871 });
1872 continue;
1873 };
1874 let mut serve = server.lock();
1875 if serve.closed {
1876 drop(serve);
1879 events.push_back(Event::Resolved {
1880 route,
1881 result: Err(Refusal {
1882 err: Error::Unroutable,
1883 standing: true,
1884 }),
1885 });
1886 continue;
1887 }
1888 if let Some(weak) = serve.served.get(&path) {
1891 drop(serve);
1892 let id = next_source;
1893 next_source += 1;
1894 sources.insert(id, weak.consume());
1895 events.push_back(Event::Resolved { route, result: Ok(id) });
1896 continue;
1897 }
1898 let pending = match serve.requests.join(&path) {
1899 Some(producer) => producer.consume(),
1900 None => {
1901 let producer = kio::Producer::<PendingBroadcast>::default();
1902 let consumer = producer.consume();
1903 match serve.requests.insert(path.clone(), producer) {
1904 Ok(()) => consumer,
1905 Err(_) => {
1908 drop(serve);
1909 events.push_back(Event::Resolved {
1910 route,
1911 result: Err(Refusal {
1912 err: Error::Unroutable,
1913 standing: true,
1914 }),
1915 });
1916 continue;
1917 }
1918 }
1919 }
1920 };
1921 upstream = Some((route, pending));
1922 }
1923 Action::Detach { source } => {
1924 sources.remove(&source);
1925 for io in tracks.values_mut() {
1928 if io.copy.as_ref().is_some_and(|(s, _)| *s == source) {
1929 io.copy = None;
1930 }
1931 if io.query.as_ref().is_some_and(|(s, ..)| *s == source) {
1932 io.query = None;
1933 }
1934 if io.staged.as_ref().is_some_and(|(s, _)| *s == source) {
1935 io.staged = None;
1936 }
1937 }
1938 }
1939 Action::Resolve => {
1940 if let Ok(mut pending) = request.write() {
1941 pending.resolved.get_or_insert(Ok(broadcast.consume()));
1942 }
1943 }
1944 Action::Query { track: name, source } => {
1945 let Some(io) = tracks.get_mut(&name) else { continue };
1946 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
1947 match sources.get(&source).map(|s| s.track(&name)) {
1948 Some(Ok(copy)) => {
1949 let query = copy.query().into_inner();
1952 io.query = Some((source, copy, query));
1953 }
1954 Some(Err(err)) => events.push_back(Event::TrackInfo {
1955 track: name,
1956 source,
1957 closing,
1958 result: Err(err),
1959 }),
1960 None => {}
1961 }
1962 }
1963 Action::Splice { track: name, source } => {
1964 let Some(io) = tracks.get_mut(&name) else { continue };
1965 let Some((staged, copy)) = io.staged.take() else {
1966 continue;
1967 };
1968 if staged != source {
1969 continue;
1970 }
1971 if let Err(err) = io.resume.takeover(©) {
1972 let _ = io.resume.abort(err);
1976 tracks.remove(&name);
1977 continue;
1978 }
1979 io.warm = None;
1980 io.edge = io.resume.resume_position();
1983 io.copy = Some((source, copy));
1984 }
1985 Action::Park { track: name } => {
1986 let Some(io) = tracks.get_mut(&name) else { continue };
1987 let Some((_, copy)) = io.copy.take() else { continue };
1988 let warm = warm_copy(©);
1992 drop(copy);
1993 if io.resume.release().is_err() {
1994 tracks.remove(&name);
1995 continue;
1996 }
1997 if let Some(warm) = warm {
1998 if let Err(err) = io.resume.takeover(&warm.track) {
1999 let _ = io.resume.abort(err);
2000 tracks.remove(&name);
2001 continue;
2002 }
2003 io.warm = Some(warm);
2004 }
2005 }
2006 Action::Release { track: name } => {
2007 let Some(io) = tracks.get_mut(&name) else { continue };
2008 io.warm = None;
2009 if io.resume.release().is_err() {
2010 tracks.remove(&name);
2011 }
2012 }
2013 Action::Finish { track: name } => {
2014 if let Some(mut io) = tracks.remove(&name) {
2015 let _ = io.resume.finish();
2016 }
2017 }
2018 Action::Abort { track: name, err } => {
2019 if let Some(mut io) = tracks.remove(&name) {
2020 tracing::debug!(name = %name, %err, "aborting track");
2021 let _ = io.resume.abort(err);
2022 }
2023 }
2024 Action::Arm { at } => deadline.set(at),
2025 Action::End { err } => {
2026 if let Ok(mut pending) = request.write() {
2027 pending.resolved.get_or_insert(Err(err.clone()));
2028 }
2029 broadcast.finish();
2036 broadcast.release_spliced(err.clone());
2037 for (_, mut io) in tracks.drain() {
2038 if !io.resume.is_used() || !io.resume.is_spliced() || io.warm.is_some() {
2040 let _ = io.resume.abort(err.clone());
2041 }
2042 }
2043 return;
2044 }
2045 }
2046 }
2047 }
2048
2049 let step = kio::wait(|waiter| {
2050 if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
2051 return Poll::Ready(Step::Assigned(name, resume));
2052 }
2053 if let Some((route, pending)) = &upstream
2054 && let Poll::Ready(result) = pending.poll(waiter, |p| match &p.resolved {
2055 Some(result) => Poll::Ready(result.clone()),
2056 None => Poll::Pending,
2057 }) {
2058 return Poll::Ready(Step::Resolved(
2059 *route,
2060 match result {
2061 Ok(resolved) => resolved,
2062 Err(_closed) => Err(Error::Unroutable),
2065 },
2066 ));
2067 }
2068 if let Some(id) = front.serving()
2069 && let Some(source) = sources.get(&id)
2070 && source.poll_closed(waiter).is_ready()
2071 {
2072 return Poll::Ready(Step::SourceClosed(id));
2073 }
2074 for (name, io) in &tracks {
2075 if let Some((source, _, query)) = &io.query
2076 && let Poll::Ready(result) = query.poll(waiter)
2077 {
2078 return Poll::Ready(Step::Info(name.clone(), *source, result));
2079 }
2080 if let Some((source, copy)) = &io.copy
2081 && let Poll::Ready(result) = copy.poll_complete(waiter)
2082 {
2083 return Poll::Ready(Step::Ended(name.clone(), *source, result));
2084 }
2085 let edge = match io.used {
2087 true => io.resume.poll_unused(waiter),
2088 false => io.resume.poll_used(waiter),
2089 };
2090 if edge.is_ready() {
2091 return Poll::Ready(Step::Demand(name.clone()));
2092 }
2093 }
2094 if deadline.poll(waiter).is_ready() {
2095 return Poll::Ready(Step::Deadline);
2096 }
2097 watch.poll_changed(waiter, seen).map(|()| Step::Table)
2098 })
2099 .await;
2100
2101 let event = match step {
2102 Step::Assigned(name, resume) => {
2103 tracks.insert(
2104 name.clone(),
2105 TrackIo {
2106 resume,
2107 staged: None,
2108 query: None,
2109 copy: None,
2110 edge: None,
2111 warm: None,
2112 used: false,
2113 },
2114 );
2115 Event::TrackAssigned { track: name }
2116 }
2117 Step::Resolved(route, result) => {
2118 upstream = None;
2119 match result {
2120 Ok(source) => {
2121 let id = next_source;
2122 next_source += 1;
2123 sources.insert(id, source);
2124 Event::Resolved { route, result: Ok(id) }
2125 }
2126 Err(err) => {
2127 let standing =
2131 !matches!(err, Error::Unroutable) || shared.read().routes.covers(&path.as_path(), route);
2132 Event::Resolved {
2133 route,
2134 result: Err(Refusal { err, standing }),
2135 }
2136 }
2137 }
2138 }
2139 Step::SourceClosed(source) => Event::SourceClosed { source },
2140 Step::Info(name, source, result) => {
2141 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2142 let Some(io) = tracks.get_mut(&name) else { continue };
2143 let Some((_, copy, _)) = io.query.take() else { continue };
2144 let result = match result {
2147 Ok(info) => match copy.poll_complete(&kio::Waiter::noop()) {
2148 Poll::Ready(Err(err)) => Err(err),
2149 _ => Ok(info),
2150 },
2151 Err(err) => Err(err),
2152 };
2153 if result.is_ok() && io.used {
2156 io.staged = Some((source, copy));
2157 }
2158 Event::TrackInfo {
2159 track: name,
2160 source,
2161 closing,
2162 result,
2163 }
2164 }
2165 Step::Ended(name, source, result) => {
2166 let closing = sources.get(&source).is_some_and(|s| s.is_closing());
2167 let Some(io) = tracks.get_mut(&name) else { continue };
2168 io.copy = None;
2169 let delivered = io.resume.resume_position() != io.edge;
2170 Event::TrackEnded {
2171 track: name,
2172 source,
2173 closing,
2174 result,
2175 delivered,
2176 }
2177 }
2178 Step::Demand(name) => {
2179 let Some(io) = tracks.get_mut(&name) else { continue };
2180 io.used = io.resume.is_used();
2181 if !io.used {
2182 io.query = None;
2185 io.staged = None;
2186 }
2187 match io.used {
2188 true => Event::Used { track: name },
2189 false => Event::Unused {
2190 track: name,
2191 now: timers.now(),
2192 },
2193 }
2194 }
2195 Step::Deadline => {
2196 deadline.set(None);
2199 Event::Deadline { now: timers.now() }
2200 }
2201 Step::Table => select(&mut front, &sources, &mut seen),
2202 };
2203 events.push_back(event);
2204 }
2205}
2206
2207#[derive(Default)]
2212struct RouteTable {
2213 root: RouteNode,
2214}
2215
2216#[derive(Default)]
2219struct RouteNode {
2220 entries: Vec<RouteEntry>,
2222 cursors: Vec<ConsumerId>,
2224 cursors_below: usize,
2227 watches: Vec<(u64, kio::Producer<Watched>)>,
2230 watches_below: usize,
2233 children: HashMap<String, RouteNode>,
2234}
2235
2236#[derive(Default)]
2239struct Watched {
2240 generation: u64,
2241}
2242
2243struct Watch {
2249 shared: kio::Shared<OriginState>,
2250 path: PathOwned,
2251 id: u64,
2252 signal: kio::Consumer<Watched>,
2253}
2254
2255impl Watch {
2256 fn seen(&self) -> u64 {
2260 self.signal.read().generation
2261 }
2262
2263 fn poll_changed(&self, waiter: &kio::Waiter, seen: u64) -> Poll<()> {
2265 self.signal
2266 .poll(waiter, |watched| match watched.generation != seen {
2267 true => Poll::Ready(()),
2268 false => Poll::Pending,
2269 })
2270 .map(|_| ())
2271 }
2272}
2273
2274impl Drop for Watch {
2275 fn drop(&mut self) {
2276 self.shared.lock().routes.remove_watch(&self.path, self.id);
2277 }
2278}
2279
2280#[derive(Clone, Copy)]
2282struct Below {
2283 cursors: usize,
2284 watches: usize,
2285}
2286
2287impl Below {
2288 const NONE: Self = Self { cursors: 0, watches: 0 };
2289 const CURSOR: Self = Self { cursors: 1, watches: 0 };
2290 const WATCH: Self = Self { cursors: 0, watches: 1 };
2291}
2292
2293impl RouteNode {
2294 fn is_empty(&self) -> bool {
2296 self.entries.is_empty() && self.cursors.is_empty() && self.watches.is_empty() && self.children.is_empty()
2297 }
2298
2299 fn find<'a>(&self, mut parts: impl Iterator<Item = &'a str>) -> Option<&Self> {
2301 match parts.next() {
2302 None => Some(self),
2303 Some(part) => self.children.get(part)?.find(parts),
2304 }
2305 }
2306
2307 fn reach<'a>(&mut self, mut parts: impl Iterator<Item = &'a str>, below: Below) -> &mut Self {
2310 self.cursors_below += below.cursors;
2311 self.watches_below += below.watches;
2312 match parts.next() {
2313 None => self,
2314 Some(part) => self.children.entry(part.to_string()).or_default().reach(parts, below),
2315 }
2316 }
2317
2318 fn edit<'a, R>(
2322 &mut self,
2323 mut parts: impl Iterator<Item = &'a str>,
2324 below: Below,
2325 f: impl FnOnce(&mut Self) -> R,
2326 ) -> Option<R> {
2327 let result = match parts.next() {
2328 None => f(self),
2329 Some(part) => {
2330 let child = self.children.get_mut(part)?;
2331 let result = child.edit(parts, below, f)?;
2332 if child.is_empty() {
2333 self.children.remove(part);
2334 }
2335 result
2336 }
2337 };
2338 self.cursors_below -= below.cursors;
2339 self.watches_below -= below.watches;
2340 Some(result)
2341 }
2342
2343 fn poke(&self) {
2345 for (_, watch) in &self.watches {
2346 if let Ok(mut watched) = watch.write() {
2347 watched.generation += 1;
2348 }
2349 }
2350 }
2351
2352 fn poke_below(&self) {
2355 if self.watches_below == 0 {
2356 return;
2357 }
2358 self.poke();
2359 for child in self.children.values() {
2360 child.poke_below();
2361 }
2362 }
2363
2364 fn walk<'a>(&'a self, visit: &mut impl FnMut(&'a Self)) {
2366 visit(self);
2367 for child in self.children.values() {
2368 child.walk(visit);
2369 }
2370 }
2371
2372 fn collect_cursors(&self, out: &mut Vec<ConsumerId>) {
2374 if self.cursors_below == 0 {
2375 return;
2376 }
2377 out.extend(&self.cursors);
2378 for child in self.children.values() {
2379 child.collect_cursors(out);
2380 }
2381 }
2382}
2383
2384impl RouteTable {
2385 fn split(&self, path: &Path) -> (Vec<&RouteNode>, Option<&RouteNode>) {
2388 let mut above = Vec::new();
2389 let mut node = &self.root;
2390 for part in path.parts() {
2391 above.push(node);
2392 match node.children.get(part) {
2393 Some(child) => node = child,
2394 None => return (above, None),
2395 }
2396 }
2397 (above, Some(node))
2398 }
2399
2400 fn covering(&self, path: &Path) -> impl Iterator<Item = &RouteEntry> {
2402 let (above, at) = self.split(path);
2403 above.into_iter().chain(at).flat_map(|node| node.entries.iter())
2404 }
2405
2406 fn covers(&self, path: &Path, id: u64) -> bool {
2408 self.covering(path).any(|entry| entry.id == id)
2409 }
2410
2411 fn at(&self, prefix: &Path) -> impl Iterator<Item = &RouteEntry> {
2413 self.root
2414 .find(prefix.parts())
2415 .into_iter()
2416 .flat_map(|node| node.entries.iter())
2417 }
2418
2419 fn entries(&self) -> impl Iterator<Item = &RouteEntry> {
2421 let mut nodes = Vec::new();
2422 self.root.walk(&mut |node| nodes.push(node));
2423 nodes.into_iter().flat_map(|node| node.entries.iter())
2424 }
2425
2426 fn insert(&mut self, entry: RouteEntry) {
2428 let node = self.root.reach(entry.prefix.parts(), Below::NONE);
2429 node.entries.push(entry);
2430 }
2431
2432 fn entry_mut(&mut self, prefix: &Path, id: u64) -> Option<&mut RouteEntry> {
2434 let mut node = &mut self.root;
2435 for part in prefix.parts() {
2436 node = node.children.get_mut(part)?;
2437 }
2438 node.entries.iter_mut().find(|entry| entry.id == id)
2439 }
2440
2441 fn remove(&mut self, prefix: &Path, id: u64) -> Option<RouteEntry> {
2443 self.root
2444 .edit(prefix.parts(), Below::NONE, |node| {
2445 let index = node.entries.iter().position(|entry| entry.id == id)?;
2446 Some(node.entries.swap_remove(index))
2447 })
2448 .flatten()
2449 }
2450
2451 fn add_cursor(&mut self, head: &Path, id: ConsumerId) {
2453 self.root.reach(head.parts(), Below::CURSOR).cursors.push(id);
2454 }
2455
2456 fn remove_cursor(&mut self, head: &Path, id: ConsumerId) {
2459 self.root.edit(head.parts(), Below::CURSOR, |node| {
2460 node.cursors.retain(|cursor| *cursor != id)
2461 });
2462 }
2463
2464 fn add_watch(&mut self, path: &Path, id: u64) -> kio::Consumer<Watched> {
2466 let producer = kio::Producer::<Watched>::default();
2467 let consumer = producer.consume();
2468 self.root.reach(path.parts(), Below::WATCH).watches.push((id, producer));
2469 consumer
2470 }
2471
2472 fn remove_watch(&mut self, path: &Path, id: u64) {
2475 self.root.edit(path.parts(), Below::WATCH, |node| {
2476 node.watches.retain(|(watch, _)| *watch != id)
2477 });
2478 }
2479
2480 fn poke_below(&self, prefix: &Path) {
2482 if let (_, Some(node)) = self.split(prefix) {
2483 node.poke_below();
2484 }
2485 }
2486
2487 fn poke_all(&self) {
2489 self.root.walk(&mut |node| node.poke());
2490 }
2491
2492 fn cursors_touching(&self, prefix: &Path) -> Vec<ConsumerId> {
2496 let (above, at) = self.split(prefix);
2497 let mut cursors: Vec<ConsumerId> = above.iter().flat_map(|node| node.cursors.iter().copied()).collect();
2498 if let Some(node) = at {
2499 node.collect_cursors(&mut cursors);
2500 }
2501 cursors.sort_unstable();
2503 cursors.dedup();
2504 cursors
2505 }
2506}
2507
2508#[derive(Default)]
2515struct OriginState {
2516 routes: RouteTable,
2519 next_route: u64,
2520 next_watch: u64,
2521
2522 cursors: HashMap<ConsumerId, TableCursor>,
2526
2527 fronts: WeakCache<FrontKey, RemoteFront>,
2535
2536 closed: bool,
2539}
2540
2541impl OriginState {
2542 fn sync_route(&mut self, prefix: &Path, claim: &Pattern) {
2547 let routes = &self.routes;
2549 for id in routes.cursors_touching(prefix) {
2550 let Some(cursor) = self.cursors.get_mut(&id) else {
2551 continue;
2552 };
2553 if let Some(presented) = cursor.presented(prefix, claim) {
2554 Self::sync_cursor(routes, cursor, &presented);
2555 }
2556 }
2557 routes.poke_below(prefix);
2559 }
2560
2561 fn watch(&mut self, shared: &kio::Shared<OriginState>, path: &Path) -> Watch {
2563 let id = self.next_watch;
2564 self.next_watch += 1;
2565 let signal = self.routes.add_watch(path, id);
2566 Watch {
2567 shared: shared.clone(),
2568 path: path.to_owned(),
2569 id,
2570 signal,
2571 }
2572 }
2573
2574 fn sync_cursor(routes: &RouteTable, cursor: &mut TableCursor, presented: &PathOwned) {
2577 let candidates: Vec<&RouteEntry> = match presented.is_empty() {
2583 true => routes
2584 .covering(&cursor.root)
2585 .filter(|entry| cursor.visible(entry))
2586 .collect(),
2587 false => {
2588 let absolute = cursor.root.join(presented);
2589 routes.at(&absolute).filter(|entry| cursor.visible(entry)).collect()
2590 }
2591 };
2592 let most = candidates.iter().map(|entry| entry.prefix.len()).max();
2593 let best = most.and_then(|most| {
2594 candidates
2595 .into_iter()
2596 .filter(|entry| entry.prefix.len() == most)
2597 .min_by_key(|entry| (!entry.local, route_order(&entry.prefix, entry)))
2598 });
2599
2600 match best {
2601 Some(entry) => {
2602 let meta = (entry.hops.clone(), entry.cost);
2603 let served = entry.server.is_some();
2604 let captures = cursor.captures(&entry.prefix);
2605 let previous = cursor
2606 .current
2607 .insert(presented.clone(), (entry.id, meta.clone(), served, captures.clone()));
2608 match previous {
2609 Some((_, prev, prev_served, prev_captures))
2616 if prev == meta && prev_served == served && prev_captures == captures => {}
2617 Some((_, prev, _, prev_captures)) if prev_captures != captures => {
2620 if let Ok(mut state) = cursor.state.write() {
2621 state.apply_unannounce(presented.clone(), prev, prev_captures);
2622 state.apply_announce(presented.clone(), meta, captures);
2623 }
2624 }
2625 _ => {
2626 if let Ok(mut state) = cursor.state.write() {
2627 state.apply_announce(presented.clone(), meta, captures);
2628 }
2629 }
2630 }
2631 }
2632 None => {
2633 if let Some((_, last, _, captures)) = cursor.current.remove(presented)
2634 && let Ok(mut state) = cursor.state.write()
2635 {
2636 state.apply_unannounce(presented.clone(), last, captures);
2637 }
2638 }
2639 }
2640 }
2641
2642 fn register_cursor(&mut self, id: ConsumerId, mut cursor: TableCursor) {
2644 let mut presented: BTreeSet<PathOwned> = BTreeSet::new();
2647 for head in &cursor.heads {
2648 let (above, at) = self.routes.split(head);
2649 let mut nodes = above;
2650 if let Some(node) = at {
2651 node.walk(&mut |node| nodes.push(node));
2652 }
2653 for entry in nodes.into_iter().flat_map(|node| node.entries.iter()) {
2654 if let Some(p) = cursor.presented(&entry.prefix, &entry.claim) {
2655 presented.insert(p);
2656 }
2657 }
2658 }
2659 for p in &presented {
2660 Self::sync_cursor(&self.routes, &mut cursor, p);
2661 }
2662 for head in &cursor.heads {
2663 self.routes.add_cursor(head, id);
2664 }
2665 self.cursors.insert(id, cursor);
2666 }
2667
2668 fn best_route(&self, path: &Path, exclude: Option<Hop>, pin: Pin, refused: &HashSet<u64>) -> Option<&RouteEntry> {
2681 let (above, at) = self.routes.split(path);
2685 let mut best = None;
2686 for node in above.into_iter().chain(at) {
2687 let mut candidates = node
2688 .entries
2689 .iter()
2690 .filter(|entry| entry.scope.matches(path.as_str()))
2691 .filter(|entry| entry.visible_to(exclude))
2692 .filter(|entry| entry.qualifies(pin))
2693 .filter(|entry| !refused.contains(&entry.id))
2694 .peekable();
2695 if candidates.peek().is_some() {
2696 best = candidates
2697 .filter(|entry| entry.serves(path))
2698 .min_by_key(|entry| (!entry.local, route_order(&entry.prefix, entry)));
2699 }
2700 }
2701 best
2702 }
2703}
2704
2705#[derive(Default)]
2712struct PendingBroadcast {
2713 resolved: Option<Result<broadcast::Consumer, Error>>,
2714}
2715
2716#[must_use = "dropping an origin::Dynamic retracts the route"]
2728pub struct Dynamic {
2729 announcement: AnnounceProducer,
2731 state: kio::Shared<ServeState>,
2732}
2733
2734impl Dynamic {
2735 pub fn update(&self, route: Route) -> Result<(), Error> {
2743 self.announcement.update(route)
2744 }
2745
2746 pub fn poll_requested_broadcast(&self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
2751 let mut state = ready!(self.state.poll(waiter, |state| {
2752 if state.closed || state.requests.has_queued() {
2753 Poll::Ready(())
2754 } else {
2755 Poll::Pending
2756 }
2757 }));
2758
2759 if state.closed {
2761 return Poll::Ready(Err(Error::Closed));
2762 }
2763
2764 let path = state.requests.pop().expect("predicate guaranteed a request");
2765 let producer = state.requests.get(&path).expect("popped key must be pending").clone();
2771 Poll::Ready(Ok(Request {
2772 path,
2773 producer,
2774 home: self.state.clone(),
2775 }))
2776 }
2777
2778 pub async fn requested_broadcast(&self) -> Result<Request, Error> {
2784 kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
2785 }
2786}
2787
2788impl ServeState {
2789 fn resolve(
2800 shared: &kio::Shared<Self>,
2801 path: &PathOwned,
2802 producer: &kio::Producer<PendingBroadcast>,
2803 result: Result<broadcast::Consumer, Error>,
2804 ) {
2805 let mut state = shared.lock();
2806 if state.closed {
2807 return;
2808 }
2809 let resolved = match result {
2810 Ok(broadcast) => {
2811 let existing = state.served.insert(path.clone(), broadcast.weak());
2815 Ok(existing.map(|weak| weak.consume()).unwrap_or(broadcast))
2816 }
2817 Err(err) => Err(err),
2818 };
2819 state.requests.remove_if(path, |p| p.same_channel(producer));
2820 if let Ok(mut pending) = producer.write() {
2821 pending.resolved.get_or_insert(resolved);
2822 drop(state);
2823 }
2824 }
2825
2826 fn forget(shared: &kio::Shared<Self>, path: &PathOwned, producer: &kio::Producer<PendingBroadcast>) {
2828 shared.lock().requests.remove_if(path, |p| p.same_channel(producer));
2829 }
2830}
2831
2832pub struct Request {
2839 path: PathOwned,
2841
2842 producer: kio::Producer<PendingBroadcast>,
2845
2846 home: kio::Shared<ServeState>,
2849}
2850
2851impl Request {
2852 pub fn path(&self) -> &Path<'_> {
2854 &self.path
2855 }
2856
2857 pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
2863 let broadcast = broadcast.consume();
2864 ServeState::resolve(&self.home, &self.path, &self.producer, Ok(broadcast));
2865 }
2867
2868 pub fn reject(self, err: Error) {
2870 ServeState::resolve(&self.home, &self.path, &self.producer, Err(err));
2871 }
2872}
2873
2874impl Drop for Request {
2875 fn drop(&mut self) {
2876 ServeState::forget(&self.home, &self.path, &self.producer);
2884 }
2885}
2886
2887pub struct Requesting {
2894 inner: RequestState,
2895 path: PathOwned,
2899 stats: stats::Scope,
2902}
2903
2904enum RequestState {
2905 Failed(Error),
2908 Pending(kio::Consumer<PendingBroadcast>),
2910}
2911
2912impl Requesting {
2913 fn failed(error: Error) -> Self {
2914 Self::new(RequestState::Failed(error))
2915 }
2916
2917 fn queued(consumer: kio::Consumer<PendingBroadcast>) -> Self {
2918 Self::new(RequestState::Pending(consumer))
2919 }
2920
2921 pub fn is_queued(&self) -> bool {
2930 matches!(self.inner, RequestState::Pending(_))
2931 }
2932
2933 fn new(inner: RequestState) -> Self {
2934 Self {
2935 inner,
2936 path: PathOwned::default(),
2937 stats: stats::Scope::default(),
2938 }
2939 }
2940
2941 fn with_path(mut self, path: PathOwned) -> Self {
2942 self.path = path;
2943 self
2944 }
2945
2946 fn with_stats(mut self, scope: stats::Scope) -> Self {
2948 self.stats = scope;
2949 self
2950 }
2951
2952 fn hand_out(&self, broadcast: broadcast::Consumer) -> broadcast::Consumer {
2954 broadcast.with_path(self.path.clone()).with_stats(self.stats.clone())
2955 }
2956
2957 pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
2959 match &self.inner {
2960 RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
2961 RequestState::Pending(consumer) => Poll::Ready(
2962 match ready!(consumer.poll(waiter, |state| match &state.resolved {
2963 Some(result) => Poll::Ready(result.clone()),
2964 None => Poll::Pending,
2965 })) {
2966 Ok(result) => result.map(|broadcast| self.hand_out(broadcast)),
2967 Err(_closed) => Err(Error::Unroutable),
2969 },
2970 ),
2971 }
2972 }
2973}
2974
2975impl kio::Pollable for Requesting {
2976 type Output = Result<broadcast::Consumer, Error>;
2977
2978 fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
2979 self.poll_ok(waiter)
2980 }
2981}
2982
2983pub trait Consume<T> {
2991 fn consume(&self) -> T;
2993}
2994
2995impl<T, U: Consume<T>> Consume<T> for &U {
2996 fn consume(&self) -> T {
2997 (**self).consume()
2998 }
2999}
3000
3001impl Consume<Consumer> for Producer {
3002 fn consume(&self) -> Consumer {
3003 Consumer::from_producer(self, stats::Session::default())
3007 }
3008}
3009
3010impl Consume<Consumer> for Consumer {
3011 fn consume(&self) -> Consumer {
3012 self.clone()
3013 }
3014}
3015
3016impl Consume<broadcast::Consumer> for broadcast::Producer {
3017 fn consume(&self) -> broadcast::Consumer {
3018 self.consume()
3020 }
3021}
3022
3023impl Consume<broadcast::Consumer> for broadcast::Consumer {
3024 fn consume(&self) -> broadcast::Consumer {
3025 self.clone()
3026 }
3027}
3028
3029impl Consume<track::Consumer> for track::Producer {
3030 fn consume(&self) -> track::Consumer {
3031 self.consume()
3032 }
3033}
3034
3035impl Consume<track::Consumer> for track::Consumer {
3036 fn consume(&self) -> track::Consumer {
3037 self.clone()
3038 }
3039}
3040
3041#[derive(Clone)]
3047pub struct Consumer {
3048 hop: Hop,
3050 scope: OriginScope,
3051
3052 root: PathOwned,
3054
3055 shared: kio::Shared<OriginState>,
3058
3059 stats: stats::Session,
3063
3064 exclude: Option<Hop>,
3070
3071 pool: cache::Pool,
3074 cache_duration: Duration,
3075
3076 tasks: TasksWeak,
3080
3081 timers: Clock,
3084}
3085
3086impl Consumer {
3087 fn from_producer(producer: &Producer, stats: stats::Session) -> Self {
3088 Self {
3089 hop: producer.hop,
3090 scope: producer.scope.clone(),
3091 root: producer.root.clone(),
3092 shared: producer.shared.clone(),
3093 stats,
3094 exclude: None,
3095 pool: producer.pool.clone(),
3096 cache_duration: producer.cache_duration,
3097 tasks: producer.tasks.downgrade(),
3098 timers: producer.timers.clone(),
3099 }
3100 }
3101
3102 pub fn hop(&self) -> Hop {
3104 self.hop
3105 }
3106
3107 pub(crate) fn excluding(mut self, peer: Hop) -> Self {
3115 self.exclude = Some(peer);
3116 self
3117 }
3118
3119 pub fn with_stats(mut self, session: stats::Session) -> Self {
3123 self.stats = session;
3124 self
3125 }
3126
3127 fn untagged(&self) -> Self {
3131 Self {
3132 stats: stats::Session::default(),
3133 ..self.clone()
3134 }
3135 }
3136
3137 pub(crate) fn empty(&self) -> Self {
3142 Self {
3143 scope: OriginScope::empty(),
3144 ..self.clone()
3145 }
3146 }
3147
3148 pub fn announced(&self) -> AnnounceConsumer {
3155 AnnounceConsumer::new(
3156 self.root.clone(),
3157 self.scope.allowed.clone(),
3158 self.stats.clone(),
3159 self.exclude,
3160 &self.shared,
3161 )
3162 }
3163
3164 pub fn consume(&self) -> Self {
3166 self.clone()
3167 }
3168
3169 #[cfg(test)]
3172 pub(crate) fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
3173 let full = self.root.join(path).to_owned();
3174 if !self.scope.permits(&full) {
3175 return None;
3176 }
3177 let table = self.shared.lock();
3178 table
3179 .routes
3180 .at(&full)
3181 .filter(|entry| entry.local)
3182 .min_by_key(|entry| route_order(&entry.prefix, entry))
3183 .and_then(|entry| entry.source.clone())
3184 }
3185
3186 pub async fn routed(&self, path: impl AsPath) -> Option<Route> {
3197 let path = path.as_path();
3198
3199 let consumer = match Pattern::subtree(path.as_str()) {
3204 Ok(subtree) => self.scope("", &Patterns::from(subtree)).ok()?,
3205 Err(InvalidPattern::TooManySegments) => self.clone(),
3206 Err(_) => return None,
3207 };
3208
3209 if !consumer.allowed().matches(path.as_str()) {
3213 return None;
3214 }
3215
3216 let mut announced = consumer.untagged().announced();
3219 loop {
3220 let update = announced.next().await?;
3221 if update.kind.is_active() && path.has_prefix(&update.prefix) {
3222 return Some(update.route);
3223 }
3224 }
3225 }
3226
3227 pub async fn routed_broadcast(&self, path: impl AsPath) -> Result<broadcast::Consumer, Error> {
3241 let path = path.as_path();
3242
3243 if !self.allowed().matches(path.as_str()) {
3246 return Err(Error::Unauthorized);
3247 }
3248 loop {
3249 let (watch, seen) = {
3258 let mut table = self.shared.lock();
3259 if table.closed {
3260 return Err(Error::Closed);
3261 }
3262 let watch = table.watch(&self.shared, &self.root.join(&path));
3263 let seen = watch.seen();
3264 (watch, seen)
3265 };
3266 match self.request_broadcast(&path).await {
3267 Ok(broadcast) => return Ok(broadcast),
3268 Err(Error::Unroutable) => {
3269 kio::wait(|waiter| watch.poll_changed(waiter, seen)).await;
3270 }
3271 Err(Error::Dropped) if self.shared.lock().closed => return Err(Error::Closed),
3274 Err(err) => return Err(err),
3275 }
3276 }
3277 }
3278
3279 pub fn scope(&self, root: impl AsPath, patterns: &Patterns) -> Result<Consumer, Error> {
3286 let root = self.root.join(root).to_owned();
3287 let rooted = patterns.rooted(root.as_str()).map_err(|_| BoundsExceeded)?;
3288 let scope = self.scope.narrow(&rooted).ok_or(Error::Unauthorized)?;
3289 Ok(Consumer {
3290 scope,
3291 root,
3292 ..self.clone()
3293 })
3294 }
3295
3296 pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
3317 let path = path.as_path();
3318
3319 let absolute = self.root.join(&path).to_owned();
3323 let scope = self.stats.egress(&absolute);
3324 let requested = path.to_owned();
3328
3329 if !self.scope.permits(&absolute) {
3331 return kio::Pending::new(Requesting::failed(Error::Unauthorized));
3332 }
3333
3334 let mut state = self.shared.lock();
3335
3336 if state.closed {
3338 return kio::Pending::new(Requesting::failed(Error::Closed));
3339 }
3340
3341 let key = (absolute.clone(), self.exclude);
3347 if let Some(front) = state.fronts.get(&key) {
3348 let pending = Requesting::queued(front.request.consume())
3349 .with_path(requested)
3350 .with_stats(scope);
3351 return kio::Pending::new(pending);
3352 }
3353
3354 if state
3356 .best_route(&absolute.as_path(), self.exclude, Pin::Any, &HashSet::new())
3357 .is_none()
3358 {
3359 return kio::Pending::new(Requesting::failed(Error::Unroutable));
3360 }
3361
3362 let broadcast = broadcast::Producer::new_spliced(broadcast::Info {
3367 pool: self.pool.clone(),
3368 cache_duration: self.cache_duration,
3369 path: absolute.clone(),
3370 });
3371 let request = kio::Producer::<PendingBroadcast>::default();
3372 let consumer = request.consume();
3373 let watch = state.watch(&self.shared, &absolute);
3374 state.fronts.insert(
3375 key,
3376 RemoteFront {
3377 request: request.clone(),
3378 broadcast: broadcast.consume().weak(),
3379 },
3380 );
3381 drop(state);
3384 self.tasks.push(run_front(FrontTask {
3385 shared: self.shared.clone(),
3386 broadcast,
3387 path: absolute,
3388 exclude: self.exclude,
3389 watch,
3390 request,
3391 timers: self.timers.clone(),
3392 }));
3393 kio::Pending::new(Requesting::queued(consumer).with_path(requested).with_stats(scope))
3394 }
3395
3396 pub fn root(&self) -> &Path<'_> {
3398 &self.root
3399 }
3400
3401 pub fn allowed(&self) -> Patterns {
3403 self.scope.relative(&self.root)
3404 }
3405
3406 pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
3408 self.root.join(path)
3409 }
3410}
3411
3412pub struct AnnounceConsumer {
3417 id: ConsumerId,
3418 shared: kio::Shared<OriginState>,
3419 root: PathOwned,
3420
3421 state: kio::Producer<OriginConsumerState>,
3424
3425 stats: stats::Session,
3428
3429 guards: HashMap<PathOwned, stats::Announce>,
3433
3434 park: kio::Park,
3437}
3438
3439impl AnnounceConsumer {
3440 fn new(
3441 root: PathOwned,
3442 allowed: Patterns,
3443 stats: stats::Session,
3444 exclude: Option<Hop>,
3445 shared: &kio::Shared<OriginState>,
3446 ) -> Self {
3447 let state = kio::Producer::<OriginConsumerState>::default();
3448 let id = ConsumerId::new();
3449
3450 {
3451 let mut table = shared.lock();
3452 if table.closed {
3453 if let Ok(mut state) = state.write() {
3455 state.ended = true;
3456 }
3457 } else {
3458 table.register_cursor(
3459 id,
3460 TableCursor {
3461 root: root.clone(),
3462 heads: interest_prefixes(&allowed),
3463 allowed,
3464 exclude,
3465 state: state.clone(),
3466 current: HashMap::new(),
3467 },
3468 );
3469 }
3470 }
3471
3472 Self {
3473 id,
3474 shared: shared.clone(),
3475 root,
3476 state,
3477 stats,
3478 guards: HashMap::new(),
3479 park: kio::Park::default(),
3480 }
3481 }
3482
3483 fn hand_out(&mut self, update: AnnounceUpdate) -> AnnounceUpdate {
3485 let absolute = self.root.join(&update.prefix).to_owned();
3486 if update.kind.is_active() {
3487 let scope = self.stats.egress(&absolute);
3488 self.guards
3489 .entry(update.prefix.clone())
3490 .or_insert_with(|| scope.announce());
3491 } else {
3492 self.guards.remove(&update.prefix);
3493 }
3494 update
3495 }
3496
3497 pub async fn next(&mut self) -> Option<AnnounceUpdate> {
3505 kio::wait(|waiter| self.poll_next(waiter)).await
3506 }
3507
3508 pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<AnnounceUpdate>> {
3514 let update = {
3515 let mut state = match ready!(self.state.poll(waiter, |state| {
3516 if state.pending.is_empty() && !state.ended {
3517 Poll::Pending
3518 } else {
3519 Poll::Ready(())
3520 }
3521 })) {
3522 Ok(state) => state,
3523 Err(_) => return Poll::Ready(None),
3525 };
3526 match state.take() {
3527 Some(update) => update,
3528 None => {
3529 state.close();
3532 return Poll::Ready(None);
3533 }
3534 }
3535 };
3536 Poll::Ready(Some(self.hand_out(update)))
3537 }
3538
3539 pub fn try_next(&mut self) -> Option<AnnounceUpdate> {
3544 let update = self.state.write().ok()?.take()?;
3545 Some(self.hand_out(update))
3546 }
3547
3548 pub fn is_closed(&self) -> bool {
3550 let state = self.state.read();
3551 state.is_closed() || state.ended
3552 }
3553
3554 pub fn root(&self) -> &Path<'_> {
3556 &self.root
3557 }
3558
3559 pub fn absolute(&self, prefix: impl AsPath) -> Path<'_> {
3561 self.root.join(prefix)
3562 }
3563}
3564
3565impl futures::Stream for AnnounceConsumer {
3566 type Item = AnnounceUpdate;
3567
3568 fn poll_next(self: std::pin::Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> {
3569 let this = self.get_mut();
3570 let waiter = this.park.hold(cx).clone();
3571 this.poll_next(&waiter)
3572 }
3573}
3574
3575impl Drop for AnnounceConsumer {
3576 fn drop(&mut self) {
3577 let mut shared = self.shared.lock();
3578 if let Some(cursor) = shared.cursors.remove(&self.id) {
3579 for head in &cursor.heads {
3580 shared.routes.remove_cursor(head, self.id);
3581 }
3582 }
3583 }
3584}
3585
3586#[cfg(test)]
3587use futures::FutureExt;
3588
3589#[cfg(test)]
3590#[allow(missing_docs)] impl AnnounceConsumer {
3592 pub fn assert_next_active(&mut self, expected: impl AsPath) -> Route {
3594 let expected = expected.as_path();
3595 let update = self.next().now_or_never().expect("next blocked").expect("no next");
3596 assert_eq!(update.prefix, expected, "wrong prefix");
3597 assert!(update.kind.is_active(), "should be an active route");
3598 update.route
3599 }
3600
3601 pub fn assert_try_next_active(&mut self, expected: impl AsPath) -> Route {
3603 let expected = expected.as_path();
3604 let update = self.try_next().expect("no next");
3605 assert_eq!(update.prefix, expected, "wrong prefix");
3606 assert!(update.kind.is_active(), "should be an active route");
3607 update.route
3608 }
3609
3610 pub fn assert_next_ended(&mut self, expected: impl AsPath) {
3612 let expected = expected.as_path();
3613 let update = self.next().now_or_never().expect("next blocked").expect("no next");
3614 assert_eq!(update.prefix, expected, "wrong prefix");
3615 assert_eq!(update.kind, AnnounceKind::Retracted, "should be a retraction");
3616 }
3617
3618 pub fn assert_next_wait(&mut self) {
3619 if let Some(res) = self.next().now_or_never() {
3620 panic!("next should block: got {:?}", res.map(|u| u.prefix));
3621 }
3622 }
3623}
3624
3625#[cfg(test)]
3629pub(crate) trait ProduceTest {
3630 fn produce(self) -> Producer;
3631}
3632
3633#[cfg(test)]
3634impl ProduceTest for Config {
3635 fn produce(self) -> Producer {
3636 let (producer, driver) = Producer::new(self);
3637 if tokio::runtime::Handle::try_current().is_ok() {
3638 tokio::spawn(crate::time::run(driver));
3639 } else {
3640 std::mem::forget(driver);
3643 }
3644 producer
3645 }
3646}
3647
3648#[cfg(test)]
3649impl ProduceTest for Hop {
3650 fn produce(self) -> Producer {
3651 Config::new(self).produce()
3652 }
3653}
3654
3655#[cfg(test)]
3656mod tests {
3657 use super::*;
3658 use futures::FutureExt;
3659
3660 fn origin(id: u64) -> Hop {
3661 Hop::new(id).unwrap()
3662 }
3663
3664 fn hops(ids: &[u64]) -> Hops {
3665 let mut list = Hops::new();
3666 for &id in ids {
3667 list.push(if id == 0 { Hop::UNKNOWN } else { origin(id) }).unwrap();
3668 }
3669 list
3670 }
3671
3672 fn scopes(prefixes: &[&str]) -> Patterns {
3674 prefixes
3675 .iter()
3676 .map(|prefix| Pattern::subtree(prefix).unwrap())
3677 .collect()
3678 }
3679
3680 #[test]
3681 fn default_config_mints_a_real_hop() {
3682 let config = Config::default();
3683 assert_ne!(config.hop, Hop::UNKNOWN);
3684 let (producer, _driver) = Producer::new(config.clone());
3685 assert_eq!(producer.hop(), config.hop);
3686 assert_eq!(producer.consume().hop(), config.hop);
3687 }
3688
3689 #[test]
3690 fn random_hops_fit_legacy_lite_clients() {
3691 for _ in 0..32 {
3692 assert!(Hop::random().id() < 1u64 << 53);
3693 }
3694 }
3695
3696 async fn settle(mut check: impl FnMut() -> bool) {
3699 for _ in 0..100 {
3700 if check() {
3701 return;
3702 }
3703 tokio::task::yield_now().await;
3704 }
3705 panic!("condition never settled");
3706 }
3707
3708 async fn queued(server: &Dynamic) -> Request {
3710 let mut request = None;
3711 settle(|| match server.poll_requested_broadcast(&kio::Waiter::noop()) {
3712 Poll::Ready(Ok(popped)) => {
3713 request = Some(popped);
3714 true
3715 }
3716 _ => false,
3717 })
3718 .await;
3719 request.unwrap()
3720 }
3721
3722 async fn next_group(subscription: &mut crate::track::Subscriber) -> Result<Option<crate::group::Consumer>, Error> {
3724 let mut next = None;
3725 settle(|| match subscription.poll_recv_group(&kio::Waiter::noop()) {
3726 Poll::Ready(result) => {
3727 next = Some(result);
3728 true
3729 }
3730 Poll::Pending => false,
3731 })
3732 .await;
3733 next.unwrap()
3734 }
3735
3736 #[tokio::test]
3737 async fn announce_and_retract() {
3738 let producer = origin(1).produce();
3739 let consumer = producer.consume();
3740 let mut announced = consumer.announced();
3741 announced.assert_next_wait();
3742
3743 let announcement = producer.announce("room/alice", Route::default()).unwrap();
3744 let route = announced.assert_next_active("room/alice");
3745 assert!(route.hops.is_empty());
3746 assert_eq!(route.cost, Cost::default());
3747 announced.assert_next_wait();
3748
3749 drop(announcement);
3750 announced.assert_next_ended("room/alice");
3751 announced.assert_next_wait();
3752 }
3753
3754 #[tokio::test]
3755 async fn broadcast_announces_its_own_path() {
3756 let producer = origin(1).produce();
3757 let consumer = producer.consume();
3758 let mut announced = consumer.announced();
3759
3760 let broadcast = producer.create_broadcast("room/alice").unwrap();
3762 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::default());
3763 let mut peer = consumer.clone().excluding(Hop::UNKNOWN).announced();
3764 peer.assert_next_wait();
3765 let local = consumer.request_broadcast("room/alice").await.expect("resolves");
3766 assert_eq!(local.info().path.as_str(), "room/alice");
3767
3768 broadcast.announce(Route::default().with_cost(3)).unwrap();
3769 let route = announced.assert_next_active("room/alice");
3770 assert_eq!(route.cost, Cost::new(3));
3771 assert_eq!(peer.assert_next_active("room/alice").cost, Cost::new(3));
3772
3773 broadcast.announce(Route::default().with_cost(1)).unwrap();
3775 let route = announced.assert_next_active("room/alice");
3776 assert_eq!(route.cost, Cost::new(1));
3777
3778 broadcast.unannounce();
3781 assert_eq!(announced.assert_next_active("room/alice").cost, Cost::default());
3782 peer.assert_next_ended("room/alice");
3783 broadcast.unannounce();
3784 announced.assert_next_wait();
3785 let local = consumer.request_broadcast("room/alice").await.expect("resolves");
3786 assert_eq!(local.info().path.as_str(), "room/alice");
3787
3788 broadcast.announce(Route::default()).unwrap();
3790 announced.assert_next_wait();
3791 peer.assert_next_active("room/alice");
3792 broadcast.finish();
3793 announced.assert_next_ended("room/alice");
3794 peer.assert_next_ended("room/alice");
3795 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
3796 announced.assert_next_wait();
3797 }
3798
3799 #[tokio::test]
3800 async fn broadcast_announcement_retracts_with_the_last_producer() {
3801 let producer = origin(1).produce();
3802 let consumer = producer.consume();
3803 let mut announced = consumer.announced();
3804
3805 let broadcast = producer.create_broadcast("room/alice").unwrap();
3806 let clone = broadcast.clone();
3807 broadcast.announce(Route::default()).unwrap();
3808 announced.assert_next_active("room/alice");
3809
3810 drop(broadcast);
3812 announced.assert_next_wait();
3813 drop(clone);
3814 announced.assert_next_ended("room/alice");
3815 }
3816
3817 #[tokio::test]
3818 async fn publish_creates_and_announces_together() {
3819 let producer = origin(1).produce();
3820 let mut announced = producer.consume().announced();
3821 let _broadcast = producer.publish("room/alice", Route::default()).unwrap();
3822 announced.assert_next_active("room/alice");
3823 }
3824
3825 #[tokio::test]
3826 async fn standalone_broadcast_cannot_announce() {
3827 let broadcast = broadcast::Info::new().produce();
3828 assert!(matches!(broadcast.announce(Route::default()), Err(Error::Closed)));
3829 broadcast.unannounce();
3831 }
3832
3833 #[tokio::test]
3834 async fn announce_replays_to_late_cursor() {
3835 let producer = origin(1).produce();
3836 let _a = producer.announce("room/alice", Route::default()).unwrap();
3837 let _b = producer.announce("room/bob", Route::default()).unwrap();
3838
3839 let mut announced = producer.consume().announced();
3840 announced.assert_next_active("room/alice");
3842 announced.assert_next_active("room/bob");
3843 announced.assert_next_wait();
3844 }
3845
3846 #[tokio::test]
3847 async fn announce_keeps_its_prefix_under_a_producer_scope() {
3848 let producer = origin(1).produce();
3849 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
3850
3851 let _a = scoped.announce("", Route::default()).unwrap();
3853 let mut announced = producer.consume().announced();
3854 announced.assert_next_active("");
3855
3856 assert!(matches!(
3858 scoped.announce("other", Route::default()),
3859 Err(Error::Unauthorized)
3860 ));
3861 }
3862
3863 #[tokio::test]
3864 async fn cursor_keeps_an_overlapping_prefix_above_its_scope() {
3865 let producer = origin(1).produce();
3866 let _a = producer.announce("", Route::default()).unwrap();
3867
3868 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
3869 let mut announced = consumer.announced();
3870 announced.assert_next_active("");
3871 }
3872
3873 #[tokio::test]
3874 async fn cursor_root_strips_prefix() {
3875 let producer = origin(1).produce();
3876 let _a = producer.announce("room/alice", Route::default()).unwrap();
3877
3878 let consumer = producer
3879 .consume()
3880 .scope("room", &Patterns::from(Pattern::all()))
3881 .unwrap();
3882 let mut announced = consumer.announced();
3883 announced.assert_next_active("alice");
3884 }
3885
3886 #[tokio::test]
3887 async fn best_route_wins_and_fails_over() {
3888 let producer = origin(1).produce();
3889 let mut announced = producer.consume().announced();
3890
3891 let expensive = producer
3892 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
3893 .unwrap();
3894 let route = announced.assert_next_active("room");
3895 assert_eq!(route.cost, Cost::new(5));
3896
3897 let cheap = producer
3899 .announce("room", Route::default().with_hops(hops(&[20])).with_cost(1))
3900 .unwrap();
3901 let route = announced.assert_next_active("room");
3902 assert_eq!(route.cost, Cost::new(1));
3903
3904 drop(cheap);
3906 let route = announced.assert_next_active("room");
3907 assert_eq!(route.cost, Cost::new(5));
3908
3909 drop(expensive);
3911 announced.assert_next_ended("room");
3912 }
3913
3914 #[tokio::test]
3915 async fn identical_reannounce_is_invisible() {
3916 let producer = origin(1).produce();
3917 let mut announced = producer.consume().announced();
3918
3919 let old = producer
3920 .announce("room", Route::default().with_hops(hops(&[10])))
3921 .unwrap();
3922 let first = announced.assert_next_active("room");
3923 assert_eq!(first.hops.as_slice(), hops(&[10]).as_slice());
3924
3925 let _new = producer
3929 .announce("room", Route::default().with_hops(hops(&[10])))
3930 .unwrap();
3931 announced.assert_next_wait();
3932
3933 drop(old);
3935 announced.assert_next_wait();
3936 }
3937
3938 #[tokio::test]
3939 async fn exclude_hides_routes_through_the_peer() {
3940 let producer = origin(1).produce();
3941 let _a = producer
3942 .announce("room", Route::default().with_hops(hops(&[7])))
3943 .unwrap();
3944
3945 let mut hidden = producer.consume().excluding(origin(7)).announced();
3946 hidden.assert_next_wait();
3947
3948 let mut visible = producer.consume().excluding(origin(8)).announced();
3949 visible.assert_next_active("room");
3950 }
3951
3952 #[tokio::test]
3953 async fn exclude_matches_via_when_the_chain_is_anonymous() {
3954 let producer = origin(1).produce();
3955 let assigned = origin(777);
3956 let _echoed = producer
3957 .announce("echoed", Route::default().with_hops(hops(&[0])).with_via(assigned))
3958 .unwrap();
3959 let _local = producer
3960 .announce("local", Route::default().with_hops(hops(&[10])))
3961 .unwrap();
3962
3963 let mut hidden = producer.consume().excluding(assigned).announced();
3964 hidden.assert_next_active("local");
3965 hidden.assert_next_wait();
3966 }
3967
3968 #[tokio::test]
3969 async fn anonymous_route_loses_to_identified_at_any_cost() {
3970 let producer = origin(1).produce();
3971 let mut announced = producer.consume().announced();
3972
3973 let _anonymous = producer
3974 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(1))
3975 .unwrap();
3976 let route = announced.assert_next_active("room");
3977 assert!(route.is_anonymous());
3978 assert_eq!(route.cost, Cost::new(1));
3979
3980 let _identified = producer
3981 .announce("room", Route::default().with_hops(hops(&[10])).with_cost(5))
3982 .unwrap();
3983 let route = announced.assert_next_active("room");
3984 assert!(!route.is_anonymous());
3985 assert_eq!(route.cost, Cost::new(5));
3986 }
3987
3988 #[tokio::test]
3989 async fn anonymous_routes_order_by_cost() {
3990 let producer = origin(1).produce();
3991 let mut announced = producer.consume().announced();
3992
3993 let expensive = producer
3994 .announce("room", Route::default().with_hops(hops(&[0])).with_cost(5))
3995 .unwrap();
3996 let route = announced.assert_next_active("room");
3997 assert_eq!(route.cost, Cost::new(5));
3998
3999 let _cheap = producer
4000 .announce("room", Route::default().with_hops(hops(&[0, 7])).with_cost(1))
4001 .unwrap();
4002 let route = announced.assert_next_active("room");
4003 assert!(route.is_anonymous());
4004 assert_eq!(route.cost, Cost::new(1));
4005
4006 drop(expensive);
4007 announced.assert_next_wait();
4008 }
4009
4010 #[tokio::test]
4011 async fn anonymous_chain_from_identified_peer_still_ranks_last() {
4012 let producer = origin(1).produce();
4013 let mut announced = producer.consume().announced();
4014
4015 let _anonymous = producer
4016 .announce(
4017 "room",
4018 Route::default()
4019 .with_hops(hops(&[0, 7]))
4020 .with_cost(1)
4021 .with_via(origin(7)),
4022 )
4023 .unwrap();
4024 announced.assert_next_active("room");
4025
4026 let _identified = producer
4027 .announce("room", Route::default().with_hops(hops(&[10, 20])).with_cost(5))
4028 .unwrap();
4029 let route = announced.assert_next_active("room");
4030 assert!(!route.is_anonymous());
4031 assert_eq!(route.cost, Cost::new(5));
4032 }
4033
4034 #[tokio::test]
4035 async fn request_prefers_identified_over_cheaper_anonymous() {
4036 let producer = origin(1).produce();
4037 let consumer = producer.consume();
4038
4039 let anonymous = producer
4040 .dynamic("room", Route::default().with_hops(hops(&[0])).with_cost(1))
4041 .unwrap();
4042 let identified = producer
4043 .dynamic("room", Route::default().with_hops(hops(&[10])).with_cost(5))
4044 .unwrap();
4045
4046 let _pending = consumer.request_broadcast("room/alice");
4047 let request = queued(&identified).await;
4048 assert_eq!(request.path().as_str(), "room/alice");
4049 assert!(
4050 anonymous.poll_requested_broadcast(&kio::Waiter::noop()).is_pending(),
4051 "the cheaper anonymous route must not serve"
4052 );
4053 }
4054
4055 #[tokio::test]
4056 async fn update_reprices_in_place() {
4057 let producer = origin(1).produce();
4058 let mut announced = producer.consume().announced();
4059
4060 let announcement = producer.announce("room", Route::default()).unwrap();
4061 announced.assert_next_active("room");
4062
4063 announcement.update(Route::default().with_cost(9)).unwrap();
4064 let route = announced.assert_next_active("room");
4065 assert_eq!(route.cost, Cost::new(9));
4066 }
4067
4068 #[tokio::test]
4069 async fn retract_after_undelivered_reprice_still_delivered() {
4070 let producer = origin(1).produce();
4071 let mut announced = producer.consume().announced();
4072
4073 let announcement = producer.announce("room", Route::default()).unwrap();
4074 announced.assert_next_active("room");
4075
4076 announcement.update(Route::default().with_cost(9)).unwrap();
4080 drop(announcement);
4081 announced.assert_next_ended("room");
4082 announced.assert_next_wait();
4083 }
4084
4085 #[tokio::test]
4086 async fn scoped_cursor_advertises_most_specific_covering_route() {
4087 let producer = origin(1).produce();
4088 let _broad = producer.announce("room", Route::default().with_cost(1)).unwrap();
4091 let _narrow = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
4092
4093 let consumer = producer
4094 .consume()
4095 .scope("room/alice", &Patterns::from(Pattern::all()))
4096 .unwrap();
4097 let mut announced = consumer.announced();
4098 let route = announced.assert_next_active("");
4099 assert_eq!(route.cost, Cost::new(9));
4100 announced.assert_next_wait();
4101 }
4102
4103 #[tokio::test]
4104 async fn capture_change_retracts_before_reannouncing_a_presented_prefix() {
4105 let producer = origin(1).produce();
4106 let _broad = producer.announce("room", Route::default()).unwrap();
4107 let exact = producer.announce("room/alice", Route::default()).unwrap();
4108 let consumer = producer
4109 .consume()
4110 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
4111 .unwrap()
4112 .scope("room/alice", &Patterns::from(Pattern::all()))
4113 .unwrap();
4114 let mut announced = consumer.announced();
4115
4116 let first = announced.next().now_or_never().expect("next").expect("announce");
4117 assert_eq!(first.prefix.as_str(), "");
4118 assert_eq!(first.kind, AnnounceKind::Announced);
4119 assert_eq!(first.captures, Some(Vec::new()));
4120
4121 drop(exact);
4122 let retracted = announced.next().now_or_never().expect("next").expect("retract");
4123 assert_eq!(retracted.prefix.as_str(), "");
4124 assert_eq!(retracted.kind, AnnounceKind::Retracted);
4125 assert_eq!(retracted.captures, Some(Vec::new()));
4126 let replacement = announced.next().now_or_never().expect("next").expect("announce");
4127 assert_eq!(replacement.prefix.as_str(), "");
4128 assert_eq!(replacement.kind, AnnounceKind::Announced);
4129 assert_eq!(replacement.captures, None);
4130 }
4131
4132 #[tokio::test]
4133 async fn routed_broadcast_resolves_once_announced() {
4134 let producer = origin(1).produce();
4135 let consumer = producer.consume();
4136
4137 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
4139 assert!((&mut resolving).now_or_never().is_none());
4140
4141 let broadcast = producer.create_broadcast("room/alice").unwrap();
4142 let resolved = resolving.await.expect("resolves once created");
4143 assert_eq!(resolved.info().path.as_str(), "room/alice");
4144 drop(broadcast);
4145 }
4146
4147 #[tokio::test]
4148 async fn local_broadcast_resolves_by_exact_path() {
4149 let producer = origin(1).produce();
4150 let consumer = producer.consume();
4151
4152 let broadcast = producer.create_broadcast("room/alice").unwrap();
4153 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
4154 assert_eq!(resolved.info().path.as_str(), "room/alice");
4155 drop(broadcast);
4156
4157 let err = consumer
4159 .request_broadcast("room/bob")
4160 .now_or_never()
4161 .expect("unroutable is synchronous")
4162 .err()
4163 .unwrap();
4164 assert!(matches!(err, Error::Unroutable));
4165 }
4166
4167 #[test]
4168 fn create_broadcast_accepts_a_max_depth_path() {
4169 let producer = origin(1).produce();
4170 let path = vec!["a"; Path::MAX_PARTS].join("/");
4171 let _broadcast = producer.create_broadcast(path.as_str()).expect("max depth is allowed");
4172 let deeper = vec!["a"; Path::MAX_PARTS + 1].join("/");
4173 assert!(matches!(
4174 producer.create_broadcast(deeper.as_str()),
4175 Err(Error::BoundsExceeded(_))
4176 ));
4177 }
4178
4179 #[tokio::test]
4180 async fn duplicate_routes_aggregate_until_the_last_leaves() {
4181 let producer = origin(1).produce();
4182 let first = producer.dynamic("live", Route::default().with_cost(3)).unwrap();
4183 let second = producer.dynamic("live", Route::default().with_cost(1)).unwrap();
4184
4185 let mut announced = producer.consume().announced();
4186 let update = announced.next().now_or_never().expect("next").expect("no next");
4187 assert_eq!(update.prefix.as_str(), "live");
4188 assert_eq!(update.kind, AnnounceKind::Announced);
4189 assert_eq!(update.route.cost, Cost::new(1));
4190 announced.assert_next_wait();
4191
4192 drop(second);
4193 let update = announced.next().now_or_never().expect("next").expect("no next");
4194 assert_eq!(update.prefix.as_str(), "live");
4195 assert_eq!(update.kind, AnnounceKind::Updated);
4196 assert_eq!(update.route.cost, Cost::new(3));
4197
4198 drop(first);
4199 announced.assert_next_ended("live");
4200 announced.assert_next_wait();
4201 }
4202
4203 #[test]
4204 fn dynamic_may_cover_a_scope_but_disjoint_prefixes_are_refused() {
4205 let producer = origin(1).produce();
4206 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
4207 let _broad = scoped
4208 .dynamic("", Route::default())
4209 .expect("an overlapping prefix is accepted");
4210
4211 let _ok = scoped
4212 .dynamic("room/alice", Route::default())
4213 .expect("a contained prefix is accepted");
4214 assert!(matches!(
4215 scoped.dynamic("other", Route::default()),
4216 Err(Error::Unauthorized)
4217 ));
4218 }
4219
4220 #[tokio::test]
4221 async fn dynamic_route_keeps_its_producer_scope() {
4222 let producer = origin(1).produce();
4223 let scope = Patterns::from("*/chat".parse::<Pattern>().unwrap());
4224 let scoped = producer.scope("", &scope).unwrap();
4225 let dynamic = scoped.dynamic("", Route::default()).unwrap();
4226
4227 let mut matching = producer
4228 .consume()
4229 .scope("", &scopes(&["room/chat"]))
4230 .unwrap()
4231 .announced();
4232 matching.assert_next_active("");
4233 let mut outside = producer
4234 .consume()
4235 .scope("", &scopes(&["room/video"]))
4236 .unwrap()
4237 .announced();
4238 outside.assert_next_wait();
4239
4240 let refused = producer
4241 .consume()
4242 .request_broadcast("room/video")
4243 .now_or_never()
4244 .expect("an out-of-scope request must be refused synchronously");
4245 assert!(matches!(refused, Err(Error::Unroutable)));
4246 assert!(dynamic.requested_broadcast().now_or_never().is_none());
4247
4248 let _pending = producer.consume().request_broadcast("room/chat");
4249 let request = queued(&dynamic).await;
4250 assert_eq!(request.path().as_str(), "room/chat");
4251 }
4252
4253 #[tokio::test]
4254 async fn dynamic_accepts_a_max_depth_prefix() {
4255 let producer = origin(1).produce();
4256 let path = (0..Path::MAX_PARTS)
4257 .map(|i| format!("s{i}"))
4258 .collect::<Vec<_>>()
4259 .join("/");
4260 let mut announced = producer.consume().announced();
4261
4262 let dynamic = producer.dynamic(&path, Route::default()).expect("max depth is allowed");
4263 announced.assert_next_active(&path);
4264
4265 let _pending = producer.consume().request_broadcast(&path);
4266 let request = queued(&dynamic).await;
4267 assert_eq!(request.path().as_str(), path);
4268 }
4269
4270 #[tokio::test]
4271 async fn dynamic_exclusion_skips_routes_through_the_subscriber() {
4272 let producer = origin(1).produce();
4273 let _server = producer
4274 .dynamic("live", Route::default().with_hops(hops(&[7])))
4275 .unwrap();
4276
4277 let mut excluded = producer.consume().excluding(origin(7)).announced();
4278 excluded.assert_next_wait();
4279
4280 let mut clean = producer.consume().excluding(origin(8)).announced();
4281 clean.assert_next_active("live");
4282 }
4283
4284 #[tokio::test]
4286 async fn announce_consumer_is_a_stream() {
4287 use futures::StreamExt;
4288 let producer = origin(1).produce();
4289 let server = producer.dynamic("live", Route::default()).unwrap();
4290 let mut announced = producer.consume().announced();
4291 let update = StreamExt::next(&mut announced)
4292 .now_or_never()
4293 .expect("next")
4294 .expect("no next");
4295 assert_eq!(update.prefix.as_str(), "live");
4296 assert_eq!(update.kind, AnnounceKind::Announced);
4297 assert!(StreamExt::next(&mut announced).now_or_never().is_none());
4298 drop(server);
4299 let update = StreamExt::next(&mut announced)
4300 .now_or_never()
4301 .expect("next")
4302 .expect("no next");
4303 assert_eq!(update.kind, AnnounceKind::Retracted);
4304 }
4305
4306 #[tokio::test]
4307 async fn dynamic_retracts() {
4308 let producer = origin(1).produce();
4309 let server = producer.dynamic("live", Route::default()).unwrap();
4310 let mut announced = producer.consume().announced();
4311 announced.assert_next_active("live");
4312
4313 drop(server);
4314 announced.assert_next_ended("live");
4315 }
4316
4317 #[test]
4318 fn charged_wildcard_cost_accumulates_across_hops() {
4319 let first = Cost::new(4).charged(1);
4320 let second = first.charged(2);
4321 assert_eq!(second, Cost { warm: 7, cold: 7 });
4322 }
4323
4324 #[tokio::test]
4325 async fn local_broadcast_is_announced_only_locally() {
4326 let producer = origin(1).produce();
4327 let mut local = producer.consume().announced();
4328 let mut peer = producer.consume().excluding(Hop::UNKNOWN).announced();
4329 let broadcast = producer.create_broadcast("room/alice").unwrap();
4330 local.assert_next_active("room/alice");
4331 peer.assert_next_wait();
4332 drop(broadcast);
4333 local.assert_next_ended("room/alice");
4334 peer.assert_next_wait();
4335 }
4336
4337 #[tokio::test]
4338 async fn served_route_materializes_on_demand() {
4339 let producer = origin(1).produce();
4340 let consumer = producer.consume();
4341
4342 let server = producer.dynamic("room", Route::default()).unwrap();
4343
4344 let pending = consumer.request_broadcast("room/alice");
4345 let request = queued(&server).await;
4346 assert_eq!(request.path().as_str(), "room/alice");
4347
4348 let source = broadcast::Info::new().produce();
4349 request.accept(&source);
4350
4351 let resolved = pending.await.expect("resolves");
4352 assert_eq!(resolved.info().path.as_str(), "room/alice");
4354
4355 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
4357 assert!(again.is_clone(&resolved));
4358 }
4359
4360 #[tokio::test]
4361 async fn served_requests_coalesce() {
4362 let producer = origin(1).produce();
4363 let consumer = producer.consume();
4364 let server = producer.dynamic("room", Route::default()).unwrap();
4365
4366 let first = consumer.request_broadcast("room/alice");
4367 let second = consumer.request_broadcast("room/alice");
4368
4369 let request = queued(&server).await;
4370 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
4372
4373 let source = broadcast::Info::new().produce();
4374 request.accept(&source);
4375
4376 let first = first.await.expect("resolves");
4377 let second = second.await.expect("resolves");
4378 assert!(first.is_clone(&second));
4379 }
4380
4381 #[tokio::test]
4382 async fn retract_rejects_pending_requests() {
4383 let producer = origin(1).produce();
4384 let consumer = producer.consume();
4385 let server = producer.dynamic("room", Route::default()).unwrap();
4386
4387 let pending = consumer.request_broadcast("room/alice");
4388 drop(server);
4389
4390 let err = pending.await.err().unwrap();
4391 assert!(matches!(err, Error::Unroutable));
4392
4393 let err = consumer
4395 .request_broadcast("room/alice")
4396 .now_or_never()
4397 .expect("unroutable")
4398 .err()
4399 .unwrap();
4400 assert!(matches!(err, Error::Unroutable));
4401 }
4402
4403 #[tokio::test]
4404 async fn routed_broadcast_survives_serving_route_retraction() {
4405 let producer = origin(1).produce();
4406 let consumer = producer.consume();
4407
4408 let standby_server = producer.dynamic("room", Route::default()).unwrap();
4411 let second_server = producer.dynamic("room", Route::default()).unwrap();
4412 let incumbent_server = producer.dynamic("room", Route::default()).unwrap();
4413
4414 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
4415 assert!((&mut resolving).now_or_never().is_none());
4416
4417 drop(incumbent_server);
4423 assert!((&mut resolving).now_or_never().is_none());
4424 drop(second_server);
4425 assert!((&mut resolving).now_or_never().is_none());
4426
4427 let request = queued(&standby_server).await;
4428 let source = broadcast::Info::new().produce();
4429 request.accept(&source);
4430
4431 let resolved = resolving.await.expect("resolves via the standby");
4432 assert_eq!(resolved.info().path.as_str(), "room/alice");
4433 }
4434
4435 #[tokio::test]
4436 async fn split_horizon_skips_routes_through_the_requester() {
4437 let producer = origin(1).produce();
4438 let _server = producer
4439 .dynamic("room", Route::default().with_hops(hops(&[7])))
4440 .unwrap();
4441
4442 let excluded = producer.consume().excluding(origin(7));
4444 let err = excluded
4445 .request_broadcast("room/alice")
4446 .now_or_never()
4447 .expect("unroutable")
4448 .err()
4449 .unwrap();
4450 assert!(matches!(err, Error::Unroutable));
4451
4452 let clean = producer.consume().excluding(origin(8));
4454 let pending = clean.request_broadcast("room/alice");
4455 assert!(pending.now_or_never().is_none());
4456 }
4457
4458 #[tokio::test]
4462 async fn handler_rejection_is_final() {
4463 let producer = origin(1).produce();
4464 let consumer = producer.consume();
4465 let server = producer.dynamic("room", Route::default()).unwrap();
4466
4467 let pending = consumer.request_broadcast("room/alice");
4468 let request = queued(&server).await;
4469 request.reject(Error::Unroutable);
4470 let err = tokio::time::timeout(Duration::from_secs(5), pending)
4471 .await
4472 .expect("the front must give up, not spin")
4473 .err()
4474 .unwrap();
4475 assert!(matches!(err, Error::Unroutable));
4476
4477 let pending = consumer.request_broadcast("room/bob");
4479 let request = queued(&server).await;
4480 assert_eq!(request.path().as_str(), "room/bob");
4481 let served = broadcast::Info::new().produce();
4482 request.accept(&served);
4483 pending.await.expect("resolves");
4484 }
4485
4486 #[tokio::test]
4489 async fn routed_broadcast_waits_out_a_rejection() {
4490 let producer = origin(1).produce();
4491 let consumer = producer.consume();
4492 let server = producer.dynamic("room", Route::default()).unwrap();
4493
4494 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
4495 assert!((&mut resolving).now_or_never().is_none());
4496 let request = queued(&server).await;
4497 request.reject(Error::Unroutable);
4498
4499 for _ in 0..20 {
4501 tokio::task::yield_now().await;
4502 }
4503 assert!((&mut resolving).now_or_never().is_none());
4504 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
4505
4506 server.update(Route::default().with_cost(2)).unwrap();
4508 assert!((&mut resolving).now_or_never().is_none());
4509 let request = queued(&server).await;
4510 let served = broadcast::Info::new().produce();
4511 request.accept(&served);
4512 resolving.await.expect("resolves");
4513 }
4514
4515 #[tokio::test]
4518 async fn routed_broadcast_resolves_an_unannounced_local_broadcast() {
4519 let producer = origin(1).produce();
4520 let consumer = producer.consume();
4521
4522 let _local = producer.create_broadcast("room/alice").unwrap();
4523 let resolved = tokio::time::timeout(Duration::from_secs(5), consumer.routed_broadcast("room/alice"))
4524 .await
4525 .expect("resolves without an announce")
4526 .expect("resolves locally");
4527 assert_eq!(resolved.info().path.as_str(), "room/alice");
4528 }
4529
4530 #[tokio::test]
4533 async fn routed_broadcast_reports_teardown_as_closed() {
4534 let (producer, driver) = Producer::new(Config::new(origin(1)));
4535 let consumer = producer.consume();
4536 let _server = producer.dynamic("room", Route::default()).unwrap();
4537
4538 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
4540 assert!((&mut resolving).now_or_never().is_none());
4541
4542 drop(driver);
4543
4544 let err = tokio::time::timeout(Duration::from_secs(5), resolving)
4545 .await
4546 .expect("teardown resolves the wait")
4547 .err()
4548 .unwrap();
4549 assert!(matches!(err, Error::Closed), "unexpected end: {err}");
4550 }
4551
4552 #[tokio::test]
4555 async fn routed_broadcast_wakes_for_a_local_broadcast() {
4556 let producer = origin(1).produce();
4557 let consumer = producer.consume();
4558 let server = producer.dynamic("room", Route::default()).unwrap();
4559
4560 let mut resolving = Box::pin(consumer.routed_broadcast("room/alice"));
4561 assert!((&mut resolving).now_or_never().is_none());
4562 queued(&server).await.reject(Error::Unroutable);
4563 for _ in 0..20 {
4564 tokio::task::yield_now().await;
4565 }
4566 assert!((&mut resolving).now_or_never().is_none());
4567
4568 let _local = producer.create_broadcast("room/alice").unwrap();
4570 let resolved = resolving.await.expect("resolves locally");
4571 assert_eq!(resolved.info().path.as_str(), "room/alice");
4572 assert!(server.poll_requested_broadcast(&kio::Waiter::noop()).is_pending());
4573 }
4574
4575 #[tokio::test]
4578 async fn late_track_on_a_served_front_replays() {
4579 let producer = origin(1).produce();
4580 let consumer = producer.consume();
4581 let server = producer.dynamic("room", Route::default()).unwrap();
4582
4583 let source = broadcast::Info::new().produce();
4584 for name in ["a", "b"] {
4585 let track = source.create_track(name, None).unwrap();
4586 let mut group = track.append_group().unwrap();
4587 group.write_frame(crate::Timestamp::ZERO, name.as_bytes()).unwrap();
4588 group.finish().unwrap();
4589 std::mem::forget(track);
4591 }
4592
4593 let pending = consumer.request_broadcast("room/alice");
4594 queued(&server).await.accept(&source);
4595 let resolved = pending.await.expect("resolves");
4596
4597 let budget = track::Subscription::default().with_max_age(Duration::from_secs(3600));
4598 for name in ["a", "b"] {
4599 let mut subscription = resolved
4600 .track(name)
4601 .unwrap()
4602 .subscribe(budget.clone())
4603 .await
4604 .expect("subscribe");
4605 let mut group = tokio::time::timeout(Duration::from_secs(5), subscription.recv_group())
4606 .await
4607 .expect("the late track must replay, not park")
4608 .expect("recv group")
4609 .expect("track ended early");
4610 let frame = group.read_frame().await.expect("read frame").expect("frame");
4611 assert_eq!(&frame.payload[..], name.as_bytes());
4612 }
4613 }
4614
4615 #[tokio::test]
4616 async fn most_specific_prefix_shadows() {
4617 let producer = origin(1).produce();
4618 let consumer = producer.consume();
4619
4620 let broad_server = producer.dynamic("", Route::default()).unwrap();
4621 let _narrow = producer.announce(".dash", Route::default()).unwrap();
4624
4625 let err = consumer
4626 .request_broadcast(".dash/pid")
4627 .now_or_never()
4628 .expect("unroutable")
4629 .err()
4630 .unwrap();
4631 assert!(matches!(err, Error::Unroutable));
4632
4633 let _pending = consumer.request_broadcast("room/alice");
4635 let request = queued(&broad_server).await;
4636 assert_eq!(request.path().as_str(), "room/alice");
4637 }
4638
4639 #[tokio::test]
4640 async fn root_dynamic_serves_any_path() {
4641 let producer = origin(1).produce();
4642 let consumer = producer.consume();
4643 let mut announced = consumer.announced();
4644 let dynamic = producer.dynamic("", Route::default()).unwrap();
4645 announced.assert_next_active("");
4647
4648 let pending = consumer.request_broadcast("anything/at/all");
4649 let request = queued(&dynamic).await;
4650 assert_eq!(request.path().as_str(), "anything/at/all");
4651
4652 let source = broadcast::Info::new().produce();
4653 request.accept(&source);
4654 let resolved = pending.await.expect("resolves");
4655 assert_eq!(resolved.info().path.as_str(), "anything/at/all");
4656
4657 drop(dynamic);
4659 announced.assert_next_ended("");
4660 let err = consumer
4661 .request_broadcast("something/else")
4662 .now_or_never()
4663 .expect("unroutable")
4664 .err()
4665 .unwrap();
4666 assert!(matches!(err, Error::Unroutable));
4667 }
4668
4669 #[tokio::test]
4675 async fn out_of_scope_request_never_reaches_the_dynamic_handler() {
4676 let producer = origin(1).produce();
4677 let dynamic = producer.dynamic("", Route::default()).unwrap();
4678 let scoped = producer.consume().scope("", &scopes(&["tenant-a"])).unwrap();
4679
4680 for path in ["tenant-b/live", "tenant-a-other/live"] {
4683 let refused = scoped
4684 .request_broadcast(path)
4685 .now_or_never()
4686 .expect("an out-of-scope request must be refused synchronously, not queued");
4687 assert!(matches!(refused, Err(Error::Unauthorized)));
4688 assert!(
4689 dynamic.requested_broadcast().now_or_never().is_none(),
4690 "the dynamic handler was asked to create a broadcast the requester may not read"
4691 );
4692 }
4693 }
4694
4695 #[tokio::test]
4696 async fn routed_waits_for_coverage() {
4697 let producer = origin(1).produce();
4698 let consumer = producer.consume();
4699
4700 let mut fut = consumer.routed("room/alice").boxed();
4701 assert!((&mut fut).now_or_never().is_none());
4702
4703 let _a = producer.announce("room", Route::default().with_cost(3)).unwrap();
4705 let route = fut.now_or_never().expect("covered").expect("routed");
4706 assert_eq!(route.cost, Cost::new(3));
4707
4708 consumer
4710 .routed("room/alice/cam")
4711 .now_or_never()
4712 .expect("covered")
4713 .expect("routed");
4714 }
4715
4716 #[tokio::test]
4717 async fn routed_ignores_deeper_routes() {
4718 let producer = origin(1).produce();
4719 let consumer = producer.consume();
4720
4721 let _deep = producer.announce("room/alice/cam", Route::default()).unwrap();
4723 let mut fut = consumer.routed("room/alice").boxed();
4724 assert!((&mut fut).now_or_never().is_none());
4725
4726 let _exact = producer.announce("room/alice", Route::default()).unwrap();
4727 fut.now_or_never().expect("covered").expect("routed");
4728 }
4729
4730 #[tokio::test]
4731 async fn routed_accepts_a_max_depth_path() {
4732 let producer = origin(1).produce();
4733 let consumer = producer.consume();
4734 let path = (0..Path::MAX_PARTS)
4735 .map(|i| format!("s{i}"))
4736 .collect::<Vec<_>>()
4737 .join("/");
4738 assert_eq!(Path::new(&path).parts().count(), Path::MAX_PARTS);
4739
4740 assert!(consumer.allowed().matches(&path));
4741
4742 let mut fut = consumer.routed(&path).boxed();
4743 assert!((&mut fut).now_or_never().is_none());
4744
4745 let _a = producer.announce("", Route::default()).unwrap();
4747 fut.now_or_never().expect("covered").expect("routed");
4748 }
4749
4750 #[tokio::test]
4751 async fn teardown_ends_everything() {
4752 let (producer, driver) = Producer::new(Config::new(origin(1)));
4753 let consumer = producer.consume();
4754 let _announcement = producer.announce("room", Route::default()).unwrap();
4755 let mut announced = consumer.announced();
4756 announced.assert_next_active("room");
4757
4758 let _server = producer.dynamic("served", Route::default()).unwrap();
4759 let pending = consumer.request_broadcast("served/path");
4760
4761 drop(driver);
4762
4763 announced.assert_next_active("served");
4765 assert!(announced.next().now_or_never().expect("ended").is_none());
4766
4767 assert!(pending.now_or_never().expect("rejected").is_err());
4769 assert!(matches!(producer.announce("x", Route::default()), Err(Error::Closed)));
4770 assert!(matches!(producer.create_broadcast("x"), Err(Error::Closed)));
4771 let err = consumer
4772 .request_broadcast("y")
4773 .now_or_never()
4774 .expect("closed")
4775 .err()
4776 .unwrap();
4777 assert!(matches!(err, Error::Closed));
4778
4779 let mut late = consumer.announced();
4781 assert!(late.next().now_or_never().expect("ended").is_none());
4782 }
4783
4784 struct ResumeRig {
4787 producer: Producer,
4788 resolved: broadcast::Consumer,
4789 subscription: track::Subscriber,
4790 incumbent_track: track::Producer,
4793 }
4794
4795 impl ResumeRig {
4796 async fn new(first: &[u64]) -> (Self, Dynamic, broadcast::Producer) {
4799 let producer = origin(1).produce();
4800 let consumer = producer.consume();
4801
4802 let server = producer
4803 .dynamic("room", Route::default().with_hops(hops(first)))
4804 .unwrap();
4805
4806 let pending = consumer.request_broadcast("room/alice");
4807 let request = queued(&server).await;
4808 let source = broadcast::Info::new().produce();
4809 let track = source.create_track("video", None).unwrap();
4810 let mut group = track.append_group().unwrap();
4811 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
4812 group.finish().unwrap();
4813 request.accept(&source);
4814
4815 let resolved = pending.await.expect("resolves");
4816 let mut subscription = resolved
4817 .track("video")
4818 .unwrap()
4819 .subscribe(None)
4820 .await
4821 .expect("subscribe");
4822 let mut group = subscription
4823 .recv_group()
4824 .await
4825 .expect("recv group")
4826 .expect("track ended early");
4827 let frame = group.read_frame().await.expect("read frame").expect("frame");
4828 assert_eq!(&frame.payload[..], b"before");
4829
4830 (
4831 Self {
4832 producer,
4833 resolved,
4834 subscription,
4835 incumbent_track: track,
4836 },
4837 server,
4838 source,
4839 )
4840 }
4841
4842 fn standby(&self, first: &[u64]) -> Dynamic {
4845 self.producer
4846 .dynamic("room", Route::default().with_hops(hops(first)))
4847 .unwrap()
4848 }
4849 }
4850
4851 async fn assert_resumes(rig: &mut ResumeRig, server: &Dynamic) {
4856 let request = queued(server).await;
4857 let replacement = broadcast::Info::new().produce();
4858 let track = replacement.create_track("video", None).unwrap();
4859 let mut group = track.append_group().unwrap();
4862 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
4863 group.finish().unwrap();
4864 request.accept(&replacement);
4865
4866 let mut group = track.append_group().unwrap();
4867 group.write_frame(crate::Timestamp::ZERO, b"resumed".as_ref()).unwrap();
4868 group.finish().unwrap();
4869
4870 let mut group = rig
4871 .subscription
4872 .recv_group()
4873 .await
4874 .expect("subscription survives the failover")
4875 .expect("track ended early");
4876 let frame = group.read_frame().await.expect("read frame").expect("frame");
4877 assert_eq!(&frame.payload[..], b"resumed");
4878 }
4879
4880 #[tokio::test]
4883 async fn driver_resolves_with_live_consumers() {
4884 let (producer, driver) = Producer::new(Config::new(origin(1)));
4885 let consumer = producer.consume();
4886 let run = crate::time::run(driver);
4887 drop(producer);
4888 tokio::time::timeout(Duration::from_secs(5), run)
4889 .await
4890 .expect("driver must finish once the producers are gone");
4891 drop(consumer);
4892 }
4893
4894 #[tokio::test]
4895 async fn remote_source_resumes_through_same_first_hop() {
4896 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
4897 let standby_server = rig.standby(&[10, 20]);
4898
4899 drop(incumbent);
4901 drop(source);
4902
4903 assert_resumes(&mut rig, &standby_server).await;
4905 }
4906
4907 #[tokio::test]
4911 async fn incompatible_successor_is_refused() {
4912 for replacement in [
4913 track::Info::default().with_timescale(crate::Timescale::MICRO),
4914 track::Info::default().with_priority(7),
4915 track::Info::default().with_max_age(Duration::from_secs(7)),
4916 ] {
4917 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
4918 let standby_server = rig.standby(&[10, 20]);
4919 drop(incumbent);
4920 drop(source);
4921
4922 let request = queued(&standby_server).await;
4925 let successor = broadcast::Info::new().produce();
4926 let track = successor.create_track("video", replacement).unwrap();
4927 let mut group = track.append_group().unwrap();
4928 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
4929 group.finish().unwrap();
4930 request.accept(&successor);
4931
4932 assert!(
4933 matches!(rig.subscription.recv_group().await, Err(Error::Unsupported)),
4934 "the subscription must abort rather than resume onto incompatible metadata"
4935 );
4936
4937 let reopened = rig.resolved.track("video").unwrap();
4939 assert!(matches!(reopened.query().await, Err(Error::Unsupported)));
4940 assert!(matches!(reopened.subscribe(None).await, Err(Error::Unsupported)));
4941 }
4942 }
4943
4944 #[tokio::test]
4945 async fn different_first_hop_ends_the_subscription() {
4946 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
4947 let rival_server = rig.standby(&[11]);
4949
4950 drop(incumbent);
4953 drop(source);
4954 rig.incumbent_track.abort(Error::Dropped).unwrap();
4955
4956 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
4958 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
4959
4960 let consumer = rig.producer.consume();
4962 let pending = consumer.request_broadcast("room/alice");
4963 let request = queued(&rival_server).await;
4964 let replacement = broadcast::Info::new().produce();
4965 request.accept(&replacement);
4966 pending.await.expect("re-request resolves through the rival");
4967 }
4968
4969 #[tokio::test]
4970 async fn anonymous_routes_never_resume() {
4971 let (mut rig, incumbent, source) = ResumeRig::new(&[]).await;
4974 let _twin_server = rig.standby(&[]);
4975
4976 drop(incumbent);
4977 drop(source);
4978 rig.incumbent_track.abort(Error::Dropped).unwrap();
4979
4980 let err = rig.subscription.recv_group().await.err().expect("subscription ends");
4981 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
4982 }
4983
4984 #[tokio::test]
4991 async fn anonymous_handoff_serves_the_newcomer_immediately() {
4992 let producer = origin(1).produce();
4993
4994 let server_a = producer
4996 .dynamic("room", Route::default().with_hops(hops(&[10])))
4997 .unwrap();
4998
4999 let consumer = producer.consume().excluding(origin(30));
5002 let pending = consumer.request_broadcast("room/alice");
5003 let request = queued(&server_a).await;
5004 let source_a = broadcast::Info::new().produce();
5005 let track_a = source_a.create_track("video", None).unwrap();
5006 let mut group = track_a.append_group().unwrap();
5007 group.write_frame(crate::Timestamp::ZERO, b"from-a".as_ref()).unwrap();
5008 group.finish().unwrap();
5009 request.accept(&source_a);
5010
5011 let resolved_a = pending.await.expect("resolves");
5012 let mut sub_a = resolved_a
5013 .track("video")
5014 .unwrap()
5015 .subscribe(None)
5016 .await
5017 .expect("subscribe");
5018 let mut group = sub_a
5019 .recv_group()
5020 .await
5021 .expect("recv group")
5022 .expect("track ended early");
5023 assert_eq!(
5024 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
5025 b"from-a"
5026 );
5027
5028 drop(track_a);
5031 drop(source_a);
5032 drop(server_a);
5033
5034 let err = sub_a.recv_group().await.err().expect("front closed");
5036 assert!(matches!(err, Error::Dropped), "unexpected end: {err}");
5037
5038 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
5042 settle(|| {
5043 matches!(
5044 consumer.request_broadcast("room/alice").now_or_never(),
5045 Some(Err(Error::Unroutable))
5046 )
5047 })
5048 .await;
5049
5050 let server_b = producer
5052 .dynamic("room", Route::default().with_hops(hops(&[20])))
5053 .unwrap();
5054 let pending = consumer.request_broadcast("room/alice");
5055 let request = queued(&server_b).await;
5056 let source_b = broadcast::Info::new().produce();
5057 let track_b = source_b.create_track("video", None).unwrap();
5058 let mut group = track_b.append_group().unwrap();
5059 group.write_frame(crate::Timestamp::ZERO, b"from-b".as_ref()).unwrap();
5060 group.finish().unwrap();
5061 request.accept(&source_b);
5062
5063 let resolved_b = pending.await.expect("B's front is served immediately");
5064 assert!(
5065 !resolved_b.is_clone(&resolved_a),
5066 "B must not splice into A's closed front"
5067 );
5068
5069 let mut sub_b = resolved_b
5070 .track("video")
5071 .unwrap()
5072 .subscribe(None)
5073 .await
5074 .expect("subscribe");
5075 let mut group = sub_b
5076 .recv_group()
5077 .await
5078 .expect("recv group")
5079 .expect("track ended early");
5080 assert_eq!(
5081 &group.read_frame().await.expect("read frame").expect("frame").payload[..],
5082 b"from-b"
5083 );
5084 }
5085
5086 #[tokio::test]
5087 async fn reprice_is_invisible_to_the_subscription() {
5088 let (rig, incumbent, source) = ResumeRig::new(&[10]).await;
5089
5090 incumbent
5093 .update(Route::default().with_hops(hops(&[10])).with_cost(9))
5094 .unwrap();
5095
5096 let track = source.create_track("audio", None).unwrap();
5097 let mut group = track.append_group().unwrap();
5098 group.write_frame(crate::Timestamp::ZERO, b"steady".as_ref()).unwrap();
5099 group.finish().unwrap();
5100
5101 let mut audio = rig
5102 .resolved
5103 .track("audio")
5104 .unwrap()
5105 .subscribe(None)
5106 .await
5107 .expect("subscribe survives the reprice");
5108 let mut group = audio
5109 .recv_group()
5110 .await
5111 .expect("recv group")
5112 .expect("track ended early");
5113 let frame = group.read_frame().await.expect("read frame").expect("frame");
5114 assert_eq!(&frame.payload[..], b"steady");
5115 }
5116
5117 #[tokio::test]
5118 async fn drain_reprice_migrates_before_the_session_dies() {
5119 let (mut rig, incumbent, source) = ResumeRig::new(&[10]).await;
5120 let standby_server = rig.standby(&[10, 20]);
5121
5122 incumbent
5126 .update(Route::default().with_hops(hops(&[10])).with_cost(Cost::DRAIN))
5127 .unwrap();
5128
5129 assert_resumes(&mut rig, &standby_server).await;
5130
5131 drop(incumbent);
5133 drop(source);
5134 }
5135
5136 #[tokio::test]
5137 async fn local_sources_splice_newest_first() {
5138 let producer = origin(1).produce();
5139 let consumer = producer.consume();
5140
5141 let first = producer.create_broadcast("room/alice").unwrap();
5142 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5143
5144 let second = producer.create_broadcast("room/alice").unwrap();
5146 let again = consumer.request_broadcast("room/alice").await.expect("resolves");
5147 assert!(again.is_clone(&resolved));
5148
5149 first.finish();
5151 settle(|| consumer.get_broadcast("room/alice").is_some()).await;
5152 second.finish();
5153 settle(|| consumer.get_broadcast("room/alice").is_none()).await;
5154
5155 let _third = producer.create_broadcast("room/alice").unwrap();
5157 assert!(consumer.get_broadcast("room/alice").is_some());
5158 }
5159
5160 #[tokio::test]
5168 async fn a_finished_broadcast_concludes_in_flight_subscriptions() {
5169 let producer = origin(1).produce();
5170 let consumer = producer.consume();
5171
5172 let broadcast = producer.create_broadcast("room/alice").unwrap();
5173 let track = broadcast.create_track("video", None).unwrap();
5174
5175 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5176 let mut subscription = resolved
5177 .track("video")
5178 .unwrap()
5179 .subscribe(None)
5180 .await
5181 .expect("subscribe");
5182 let mut group = track.append_group().unwrap();
5184 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
5185 group.finish().unwrap();
5186 track.finish().unwrap();
5187 drop(track);
5188 broadcast.finish();
5189
5190 let mut group = next_group(&mut subscription)
5191 .await
5192 .expect("a cleanly finished track was served as an error")
5193 .expect("the track ended before its last group");
5194 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
5195 drop(group);
5196
5197 let end = next_group(&mut subscription)
5198 .await
5199 .expect("a cleanly finished track ended as an error");
5200 assert!(end.is_none(), "a group followed the final one");
5201 }
5202
5203 #[tokio::test]
5210 async fn a_retracted_route_concludes_in_flight_subscriptions() {
5211 let producer = origin(1).produce();
5212 let consumer = producer.consume();
5213 let server = producer
5214 .dynamic("room", Route::default().with_hops(hops(&[10])))
5215 .unwrap();
5216
5217 let pending = consumer.request_broadcast("room/alice");
5218 let request = queued(&server).await;
5219 let source = broadcast::Info::new().produce();
5220 let track = source.create_track("video", None).unwrap();
5221 request.accept(&source);
5222
5223 let resolved = pending.await.expect("resolves");
5224 let mut subscription = resolved
5225 .track("video")
5226 .unwrap()
5227 .subscribe(None)
5228 .await
5229 .expect("subscribe");
5230
5231 source.finish();
5234 drop(server);
5235 settle(|| resolved.is_closed()).await;
5236 let mut group = track.append_group().unwrap();
5237 group.write_frame(crate::Timestamp::ZERO, b"tail".as_ref()).unwrap();
5238 group.finish().unwrap();
5239 track.finish().unwrap();
5240 drop(track);
5241
5242 let mut group = next_group(&mut subscription)
5243 .await
5244 .expect("a retracted route's track was served as an error")
5245 .expect("the track ended before its last group");
5246 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"tail");
5247 drop(group);
5248
5249 let end = next_group(&mut subscription)
5250 .await
5251 .expect("a cleanly finished track ended as an error");
5252 assert!(end.is_none(), "a group followed the final one");
5253 }
5254
5255 #[tokio::test]
5260 async fn origin_front_drops_the_source_when_unused() {
5261 let producer = origin(1).produce();
5262 let consumer = producer.consume();
5263
5264 let broadcast = producer.create_broadcast("room/alice").unwrap();
5265 let track = broadcast.create_track("video", None).unwrap();
5266 let mut group = track.append_group().unwrap();
5267 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
5268 group.finish().unwrap();
5269
5270 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5271 let mut subscription = resolved
5272 .track("video")
5273 .unwrap()
5274 .subscribe(None)
5275 .await
5276 .expect("subscribe");
5277 let mut group = subscription.recv_group().await.unwrap().unwrap();
5278 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
5279 drop(group);
5280 drop(subscription);
5281
5282 tokio::time::timeout(Duration::from_secs(1), track.unused())
5283 .await
5284 .expect("source unused should resolve far below TRACK_IDLE_LINGER")
5285 .expect("source closed");
5286
5287 let mut again = resolved
5290 .track("video")
5291 .unwrap()
5292 .subscribe(track::Subscription::default().with_max_age(Duration::from_secs(3600)))
5293 .await
5294 .expect("resubscribe");
5295 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
5296 .await
5297 .expect("cached group is still on the front")
5298 .expect("recv group")
5299 .expect("track ended early");
5300 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
5301
5302 tokio::time::timeout(Duration::from_secs(1), track.used())
5303 .await
5304 .expect("returning reader re-splices the source")
5305 .expect("source closed");
5306
5307 let mut group = track.append_group().unwrap();
5308 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
5309 group.finish().unwrap();
5310 let mut group = tokio::time::timeout(Duration::from_secs(1), again.recv_group())
5311 .await
5312 .expect("groups past the cached edge come from the re-splice")
5313 .expect("recv group")
5314 .expect("track ended early");
5315 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
5316 }
5317
5318 #[tokio::test]
5323 async fn chained_front_drops_the_source_when_unused() {
5324 let leaf = origin(1).produce();
5325 let leaf_consumer = leaf.consume();
5326
5327 let broadcast = leaf.create_broadcast("room/alice").unwrap();
5328 let track = broadcast.create_track("video", None).unwrap();
5329 let mut group = track.append_group().unwrap();
5330 group.write_frame(crate::Timestamp::ZERO, b"cached".as_ref()).unwrap();
5331 group.finish().unwrap();
5332
5333 let leaf_front = leaf_consumer.request_broadcast("room/alice").await.expect("resolves");
5336
5337 let mid = origin(2).produce();
5338 let mid_server = mid.dynamic("room", Route::default().with_hops(hops(&[10]))).unwrap();
5339 let mid_pending = mid.consume().request_broadcast("room/alice");
5340 queued(&mid_server).await.accept(&leaf_front);
5341 let mid_resolved = mid_pending.await.expect("mid resolves");
5342
5343 let edge = origin(3).produce();
5344 let edge_server = edge.dynamic("room", Route::default().with_hops(hops(&[20]))).unwrap();
5345 let edge_pending = edge.consume().request_broadcast("room/alice");
5346 queued(&edge_server).await.accept(&mid_resolved);
5347 let edge_resolved = edge_pending.await.expect("edge resolves");
5348
5349 let mut subscription = edge_resolved
5350 .track("video")
5351 .unwrap()
5352 .subscribe(None)
5353 .await
5354 .expect("subscribe");
5355 let mut group = subscription.recv_group().await.unwrap().unwrap();
5356 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
5357 drop(group);
5358 drop(subscription);
5359
5360 tokio::time::timeout(Duration::from_secs(5), track.unused())
5361 .await
5362 .expect("chained unused should resolve far below TRACK_IDLE_LINGER")
5363 .expect("source closed");
5364
5365 let cached = edge_resolved.track("video").unwrap().cached_groups();
5366 assert_eq!(
5367 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
5368 vec![0],
5369 "every front keeps the delivered groups after releasing its source"
5370 );
5371
5372 let mut subscription = edge_resolved
5373 .track("video")
5374 .unwrap()
5375 .subscribe(None)
5376 .await
5377 .expect("resubscribe");
5378 tokio::time::timeout(Duration::from_secs(5), track.used())
5379 .await
5380 .expect("resubscribe should reach the leaf")
5381 .expect("source open");
5382 let mut group = track.append_group().unwrap();
5383 group.write_frame(crate::Timestamp::ZERO, b"live".as_ref()).unwrap();
5384 group.finish().unwrap();
5385 let mut group = subscription.recv_group().await.unwrap().unwrap();
5386 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"cached");
5387 drop(group);
5388 let mut group = subscription.recv_group().await.unwrap().unwrap();
5389 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"live");
5390 drop(group);
5391 drop(subscription);
5392
5393 tokio::time::timeout(Duration::from_secs(5), track.unused())
5394 .await
5395 .expect("second chained unused should resolve far below TRACK_IDLE_LINGER")
5396 .expect("source closed");
5397
5398 let cached = edge_resolved.track("video").unwrap().cached_groups();
5399 assert_eq!(
5400 cached.iter().map(|(group, _)| group.sequence).collect::<Vec<_>>(),
5401 vec![0, 1],
5402 "repeated demand keeps every complete group while releasing its source"
5403 );
5404
5405 let fetch = edge_resolved.track("video").unwrap().fetch_group(2, None);
5406 let mut fetch = std::pin::pin!(fetch);
5407 assert!(futures::poll!(fetch.as_mut()).is_pending(), "fetch should re-splice");
5408 tokio::time::timeout(Duration::from_secs(5), track.used())
5409 .await
5410 .expect("fetch should reach the leaf")
5411 .expect("source open");
5412 let mut group = track.append_group().unwrap();
5413 group.write_frame(crate::Timestamp::ZERO, b"fetched".as_ref()).unwrap();
5414 group.finish().unwrap();
5415 let mut group = tokio::time::timeout(Duration::from_secs(5), fetch)
5416 .await
5417 .expect("re-spliced source should answer the fetch")
5418 .expect("fetch succeeds");
5419 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"fetched");
5420 }
5421
5422 #[tokio::test]
5426 async fn incompatible_local_source_keeps_the_incumbent() {
5427 let producer = origin(1).produce();
5428 let consumer = producer.consume();
5429
5430 let first = producer.create_broadcast("room/alice").unwrap();
5431 let track = first.create_track("video", None).unwrap();
5432 let resolved = consumer.request_broadcast("room/alice").await.expect("resolves");
5433 let mut subscription = resolved
5434 .track("video")
5435 .unwrap()
5436 .subscribe(None)
5437 .await
5438 .expect("subscribe");
5439 let mut group = track.append_group().unwrap();
5440 group.write_frame(crate::Timestamp::ZERO, b"before".as_ref()).unwrap();
5441 group.finish().unwrap();
5442 let mut group = subscription.recv_group().await.unwrap().unwrap();
5443 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"before");
5444
5445 let second = producer.create_broadcast("room/alice").unwrap();
5447 let _incompatible = second
5448 .create_track("video", track::Info::default().with_timescale(crate::Timescale::MICRO))
5449 .unwrap();
5450 for _ in 0..10 {
5451 tokio::task::yield_now().await;
5452 }
5453
5454 let mut group = track.append_group().unwrap();
5456 group.write_frame(crate::Timestamp::ZERO, b"still".as_ref()).unwrap();
5457 group.finish().unwrap();
5458 let mut group = subscription.recv_group().await.unwrap().unwrap();
5459 assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"still");
5460
5461 drop(track);
5463 first.finish();
5464 assert!(matches!(subscription.recv_group().await, Err(Error::Unsupported)));
5465 }
5466
5467 #[tokio::test]
5468 async fn multiple_scopes_present_one_broad_prefix() {
5469 let producer = origin(1).produce();
5470 let _a = producer.announce("", Route::default()).unwrap();
5471
5472 let consumer = producer.consume().scope("", &scopes(&["alpha", "beta"])).unwrap();
5473 let mut announced = consumer.announced();
5474 announced.assert_next_active("");
5475 announced.assert_next_wait();
5476 }
5477
5478 #[test]
5479 fn scope_accepts_every_pattern_union() {
5480 let producer = origin(1).produce();
5481
5482 let root = producer.scope("", &Patterns::from(Pattern::all())).unwrap();
5484 assert_eq!(root.allowed(), Patterns::from(Pattern::all()));
5485
5486 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
5488 assert_eq!(scoped.allowed(), scopes(&["room"]));
5489
5490 let multi = producer.scope("", &scopes(&["room", "room/chat", "anon"])).unwrap();
5492 assert_eq!(multi.allowed(), scopes(&["room", "anon"]));
5493
5494 let consumer = producer.consume().scope("", &scopes(&["room"])).unwrap();
5496 assert_eq!(consumer.allowed(), scopes(&["room"]));
5497
5498 for text in ["room", "", "*room", "room/*", "*", "**/room", "room/**/chat", "*.hang"] {
5499 let union = Patterns::from(text.parse::<Pattern>().unwrap());
5500 assert_eq!(producer.scope("", &union).expect(text).allowed(), union, "{text}");
5501 assert_eq!(
5502 producer.consume().scope("", &union).expect(text).allowed(),
5503 union,
5504 "{text}"
5505 );
5506 }
5507
5508 let mixed: Patterns = ["room/**".parse().unwrap(), "other".parse().unwrap()]
5509 .into_iter()
5510 .collect();
5511 assert_eq!(producer.scope("", &mixed).unwrap().allowed(), mixed);
5512 }
5513
5514 #[test]
5515 fn route_table_prunes_to_empty() {
5516 let producer = origin(1).produce();
5517 let consumer = producer.consume();
5518
5519 let cursor = consumer
5522 .scope("", &scopes(&["room/a", "other/deep/head"]))
5523 .unwrap()
5524 .announced();
5525 let route = producer.announce("room/a/b/c", Route::default()).unwrap();
5526 {
5527 let table = producer.shared.lock();
5528 assert!(table.routes.root.find(Path::new("room/a/b/c").parts()).is_some());
5529 assert!(table.routes.root.find(Path::new("other/deep/head").parts()).is_some());
5530 assert_eq!(table.routes.root.cursors_below, 2);
5531 }
5532
5533 drop(route);
5534 drop(cursor);
5535 let table = producer.shared.lock();
5536 assert!(table.routes.root.is_empty());
5537 assert_eq!(table.routes.root.cursors_below, 0);
5538 }
5539
5540 #[test]
5544 fn a_published_broadcast_keeps_the_driver_running() {
5545 let (producer, mut driver) = Producer::new(Config::new(origin(1)));
5546 let waiter = kio::Waiter::noop();
5547 let broadcast = producer.create_broadcast("room/a").unwrap();
5548 drop(producer);
5549 assert!(
5550 driver.poll(Instant::now(), &waiter).is_ok(),
5551 "the broadcast is lifecycle work"
5552 );
5553 drop(broadcast);
5554 assert!(matches!(driver.poll(Instant::now(), &waiter), Err(Error::Closed)));
5555 }
5556
5557 #[test]
5558 fn watch_wakes_only_for_covering_changes() {
5559 let producer = origin(1).produce();
5560 let waiter = kio::Waiter::noop();
5561 let watch = producer.shared.lock().watch(&producer.shared, &Path::new("room/a"));
5562 let seen = watch.seen();
5563
5564 let _other = producer.announce("other", Route::default()).unwrap();
5566 let _below = producer.announce("room/a/b", Route::default()).unwrap();
5567 assert!(watch.poll_changed(&waiter, seen).is_pending());
5568
5569 let above = producer.announce("room", Route::default()).unwrap();
5571 assert!(watch.poll_changed(&waiter, seen).is_ready());
5572 let seen = watch.seen();
5573 drop(above);
5574 assert!(watch.poll_changed(&waiter, seen).is_ready());
5575 let seen = watch.seen();
5576
5577 let _beside = producer.create_broadcast("room/b").unwrap();
5579 assert!(watch.poll_changed(&waiter, seen).is_pending());
5580 let _here = producer.create_broadcast("room/a").unwrap();
5581 assert!(watch.poll_changed(&waiter, seen).is_ready());
5582
5583 drop(watch);
5585 let table = producer.shared.lock();
5586 let node = table
5587 .routes
5588 .root
5589 .find(Path::new("room/a").parts())
5590 .expect("route below keeps the node");
5591 assert!(node.watches.is_empty());
5592 assert_eq!(table.routes.root.watches_below, 0);
5593 }
5594
5595 #[test]
5596 fn a_discarded_front_task_unregisters_its_watch() {
5597 let (producer, _driver) = Producer::new(Config {
5598 hop: origin(1),
5599 ..Default::default()
5600 });
5601 let consumer = producer.consume();
5602 let _served = producer.dynamic("room", Route::default()).unwrap();
5603 drop(producer);
5608 let _pending = consumer.request_broadcast("room/a");
5609 }
5610
5611 #[test]
5612 fn create_broadcast_refuses_a_path_no_pattern_can_spell() {
5613 let producer = origin(1).produce();
5614
5615 assert!(matches!(
5619 producer.create_broadcast("room/*"),
5620 Err(Error::InvalidPath(_))
5621 ));
5622 assert!(matches!(
5623 producer.announce("room/**", Route::default()),
5624 Err(Error::InvalidPath(_))
5625 ));
5626 }
5627
5628 #[test]
5629 fn scope_empty_union_grants_nothing() {
5630 let producer = origin(1).produce();
5631
5632 assert!(matches!(producer.scope("", &Patterns::new()), Err(Error::Unauthorized)));
5634 assert!(matches!(
5635 producer.consume().scope("", &Patterns::new()),
5636 Err(Error::Unauthorized)
5637 ));
5638 }
5639
5640 #[test]
5641 fn scope_nests_and_rebases_roots() {
5642 let producer = origin(1).produce();
5643
5644 let scoped = producer.scope("", &scopes(&["room"])).unwrap();
5646 let nested = scoped.scope("", &scopes(&["room/chat"])).unwrap();
5647 assert_eq!(nested.allowed(), scopes(&["room/chat"]));
5648
5649 assert!(matches!(
5651 scoped.scope("", &scopes(&["other"])),
5652 Err(Error::Unauthorized)
5653 ));
5654
5655 let rooted = nested.scope("room/chat", &Patterns::from(Pattern::all())).unwrap();
5657 assert_eq!(rooted.allowed(), scopes(&[""]));
5658
5659 let broadcast = nested.create_broadcast("room/chat/live").unwrap();
5661 assert!(producer.consume().get_broadcast("room/chat/live").is_some());
5662 broadcast.finish();
5663 }
5664
5665 #[test]
5666 fn scope_intersects_and_rebases_arbitrary_grants() {
5667 let producer = origin(1).produce();
5668 let rooms = producer
5669 .scope("", &Patterns::from("room/*".parse::<Pattern>().unwrap()))
5670 .unwrap();
5671 let chats = rooms
5672 .scope("", &Patterns::from("*/chat".parse::<Pattern>().unwrap()))
5673 .unwrap();
5674 assert_eq!(chats.allowed(), Patterns::from("room/chat".parse::<Pattern>().unwrap()));
5675
5676 let exact = producer
5677 .scope("", &Patterns::from("room/alice".parse::<Pattern>().unwrap()))
5678 .unwrap();
5679 let rooted = exact.scope("room", &Patterns::from(Pattern::all())).unwrap();
5680 assert_eq!(rooted.allowed(), Patterns::from("alice".parse::<Pattern>().unwrap()));
5681 assert!(matches!(
5682 exact.scope("room/bob", &Patterns::from(Pattern::all())),
5683 Err(Error::Unauthorized)
5684 ));
5685
5686 let broadcast = exact.create_broadcast("room/alice").unwrap();
5687 assert!(matches!(
5688 exact.create_broadcast("room/alice/cam"),
5689 Err(Error::Unauthorized)
5690 ));
5691 assert!(producer.consume().get_broadcast("room/alice").is_some());
5692 drop(broadcast);
5693 }
5694
5695 #[tokio::test]
5696 async fn wildcard_scope_filters_announcements_and_reports_captures() {
5697 let producer = origin(1).produce();
5698 let consumer = producer
5699 .consume()
5700 .scope("", &Patterns::from("room/*/chat".parse::<Pattern>().unwrap()))
5701 .unwrap();
5702 let mut announced = consumer.announced();
5703
5704 let alice = producer.create_broadcast("room/alice/chat").unwrap();
5705 alice.announce(Route::default()).unwrap();
5706 let update = announced.try_next().expect("alice's chat");
5707 assert_eq!(update.prefix.as_str(), "room/alice/chat");
5708 assert_eq!(update.captures, Some(vec!["alice".parse::<Pattern>().unwrap()]));
5709
5710 let audio = producer.create_broadcast("room/alice/audio").unwrap();
5711 audio.announce(Route::default()).unwrap();
5712 announced.assert_next_wait();
5713
5714 let broad = producer.announce("room", Route::default()).unwrap();
5715 let update = announced.try_next().expect("overlapping broad route");
5716 assert_eq!(update.prefix.as_str(), "room");
5717 assert_eq!(update.captures, None, "an overlap does not pin the wildcard");
5718
5719 drop(broad);
5720 drop(audio);
5721 drop(alice);
5722 }
5723
5724 #[tokio::test]
5725 async fn local_broadcast_wins_announcement_ties() {
5726 let producer = origin(1).produce();
5727 let remote = producer.announce("room/alice", Route::default().with_cost(9)).unwrap();
5728 let local = producer.create_broadcast("room/alice").unwrap();
5729 local.announce(Route::default()).unwrap();
5730
5731 let mut announced = producer.consume().announced();
5732 let update = announced.try_next().expect("one winning route");
5733 assert_eq!(update.prefix.as_str(), "room/alice");
5734 assert_eq!(update.route.cost, Cost::default());
5735 announced.assert_next_wait();
5736
5737 drop(local);
5738 drop(remote);
5739 }
5740
5741 #[test]
5746 fn cost_charge_saturates() {
5747 assert_eq!(Cost { warm: 4, cold: 6 }.charged(5), Cost { warm: 9, cold: 11 });
5748 assert_eq!(Cost::new(u64::MAX).charged(10), Cost::new(MAX_COST));
5749
5750 assert_eq!(Cost::UNKNOWN.charged(3).cold, MAX_COST);
5753 }
5754
5755 fn expiring_origin(expiry: Duration) -> Producer {
5757 let pool = cache::Pool::new(cache::Config::default().with_expiry(expiry));
5758 Config {
5759 pool,
5760 ..Config::default()
5761 }
5762 .produce()
5763 }
5764
5765 #[tokio::test(start_paused = true)]
5769 async fn stalled_publisher_open_group_is_reclaimed() {
5770 let expiry = Duration::from_secs(1);
5771 let origin = expiring_origin(expiry);
5772 let broadcast = origin.create_broadcast("test").unwrap();
5773 let track = broadcast.create_track("video", None).unwrap();
5774
5775 let mut stalled = track.append_group().unwrap();
5776 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
5777 let _successor = track.append_group().unwrap();
5781
5782 let mut reading = stalled.consume();
5783 assert!(reading.read_frame().await.unwrap().is_some());
5784
5785 crate::model::clock::advance(expiry * 2);
5787
5788 let reclaimed = tokio::time::timeout(Duration::from_secs(60), reading.read_frame()).await;
5791 assert!(
5792 matches!(reclaimed, Ok(Err(Error::Old))),
5793 "the sweep must reclaim an idle open group and surface the gap, got {reclaimed:?}"
5794 );
5795 }
5796
5797 #[tokio::test(start_paused = true)]
5800 async fn sweep_respects_a_disabled_expiry() {
5801 let origin = Config {
5802 pool: cache::Pool::unbounded(),
5803 ..Config::default()
5804 }
5805 .produce();
5806 let broadcast = origin.create_broadcast("test").unwrap();
5807 let track = broadcast.create_track("video", None).unwrap();
5808
5809 let mut stalled = track.append_group().unwrap();
5810 stalled.write_frame(crate::Timestamp::ZERO, b"x".as_slice()).unwrap();
5811 let _successor = track.append_group().unwrap();
5812
5813 let mut reading = stalled.consume();
5814 assert!(reading.read_frame().await.unwrap().is_some());
5815
5816 crate::model::clock::advance(Duration::from_secs(3600));
5817 tokio::time::advance(Duration::from_secs(3600)).await;
5818
5819 assert!(
5820 reading.read_frame().now_or_never().is_none(),
5821 "a pool without an expiry window never reclaims"
5822 );
5823 }
5824
5825 #[test]
5828 fn drain_cost_is_encodable() {
5829 use crate::coding::Encode;
5830
5831 let mut buf = Vec::new();
5832 Cost::DRAIN
5833 .encode(&mut buf, crate::lite::Version::Lite06)
5834 .expect("a draining route is still forwarded, so its cost must encode");
5835 }
5836}